xiangfu0 commented on code in PR #19737:
URL: https://github.com/apache/pinot/pull/19737#discussion_r4194103050
##########
pinot-server/src/main/java/org/apache/pinot/server/api/resources/ReingestionResource.java:
##########
@@ -247,9 +266,38 @@ private void doReingestSegment(String realtimeTableName,
SegmentZKMetadata segme
}
}
- private void waitForCondition(
+ /// Resolves the consumption timeout for a re-ingestion job: the cluster
config takes precedence over the server
+ /// config, which takes precedence over the default. A non-numeric or
non-positive value is ignored with a warning.
+ @VisibleForTesting
+ static long getConsumptionTimeoutMs(PinotClusterConfigProvider
clusterConfigProvider,
+ PinotConfiguration serverConf) {
+ Long timeoutMs = parseTimeoutMs(
+
clusterConfigProvider.getClusterConfigs().get(CONFIG_OF_REINGESTION_CONSUMPTION_TIMEOUT_MS),
"cluster");
+ if (timeoutMs == null) {
+ timeoutMs =
parseTimeoutMs(serverConf.getProperty(CONFIG_OF_REINGESTION_CONSUMPTION_TIMEOUT_MS),
"server");
+ }
+ return timeoutMs != null ? timeoutMs :
DEFAULT_REINGESTION_CONSUMPTION_TIMEOUT_MS;
+ }
+
+ @Nullable
+ private static Long parseTimeoutMs(@Nullable String value, String
configSource) {
+ if (value == null) {
+ return null;
+ }
+ Long timeoutMs = Longs.tryParse(value.trim());
+ if (timeoutMs == null || timeoutMs <= 0) {
+ LOGGER.warn("Ignoring invalid {} config: {}={}, expecting a positive
number of milliseconds", configSource,
+ CONFIG_OF_REINGESTION_CONSUMPTION_TIMEOUT_MS, value);
+ return null;
+ }
+ return timeoutMs;
+ }
+
+ @VisibleForTesting
+ static void waitForCondition(
Function<Void, Boolean> condition, long checkIntervalMs, long timeoutMs,
long gracePeriodMs) {
- long endTime = System.currentTimeMillis() + timeoutMs;
+ // Saturate to avoid overflow, since the timeout is configurable
+ long endTime = LongMath.saturatedAdd(System.currentTimeMillis(),
timeoutMs);
Review Comment:
**`gracePeriodMs` silently consumes the timeout budget — and is dead
generality.**
`endTime` is computed before the grace sleep, so the effective wait is
`timeoutMs - gracePeriodMs`; with `gracePeriodMs >= timeoutMs` the condition
gets exactly one post-grace check before `Timeout` (latent — the single
production caller and all 7 test invocations pass 0). Since this PR promotes
the method to a tested `@VisibleForTesting` static API, consider dropping the
always-zero parameter and its untested branch instead of carrying the trap
forward.
##########
pinot-server/src/main/java/org/apache/pinot/server/api/resources/ReingestionResource.java:
##########
@@ -247,9 +266,38 @@ private void doReingestSegment(String realtimeTableName,
SegmentZKMetadata segme
}
}
- private void waitForCondition(
+ /// Resolves the consumption timeout for a re-ingestion job: the cluster
config takes precedence over the server
+ /// config, which takes precedence over the default. A non-numeric or
non-positive value is ignored with a warning.
+ @VisibleForTesting
+ static long getConsumptionTimeoutMs(PinotClusterConfigProvider
clusterConfigProvider,
+ PinotConfiguration serverConf) {
+ Long timeoutMs = parseTimeoutMs(
+
clusterConfigProvider.getClusterConfigs().get(CONFIG_OF_REINGESTION_CONSUMPTION_TIMEOUT_MS),
"cluster");
+ if (timeoutMs == null) {
+ timeoutMs =
parseTimeoutMs(serverConf.getProperty(CONFIG_OF_REINGESTION_CONSUMPTION_TIMEOUT_MS),
"server");
+ }
+ return timeoutMs != null ? timeoutMs :
DEFAULT_REINGESTION_CONSUMPTION_TIMEOUT_MS;
+ }
+
+ @Nullable
+ private static Long parseTimeoutMs(@Nullable String value, String
configSource) {
+ if (value == null) {
+ return null;
+ }
+ Long timeoutMs = Longs.tryParse(value.trim());
+ if (timeoutMs == null || timeoutMs <= 0) {
+ LOGGER.warn("Ignoring invalid {} config: {}={}, expecting a positive
number of milliseconds", configSource,
Review Comment:
**An invalid configured value is re-parsed and re-WARNed on every request —
including requests rejected with 409.**
The timeout is resolved before the duplicate-segment conflict check, so
while a bad value sits in ZK (e.g. `"30m"`), every `POST /reingestSegment/...`
— even duplicates that return CONFLICT — emits the "Ignoring invalid cluster
config" warning; during a repair storm over hundreds of segments this floods
the logs. Resolving after the `compareAndSet`, and/or logging once per distinct
invalid value, would avoid the spam.
##########
pinot-server/src/main/java/org/apache/pinot/server/api/AdminApiApplication.java:
##########
@@ -84,6 +85,7 @@ protected void configure() {
bind(_serverInstance.getServerMetrics()).to(ServerMetrics.class);
bind(_accessControlFactory).to(AccessControlFactory.class);
bind(reloadJobStatusCache).to(ServerReloadJobStatusCache.class);
+ bind(clusterConfigProvider).to(PinotClusterConfigProvider.class);
Review Comment:
**Design note: this bypasses the established
`PinotClusterConfigChangeListener` pattern.**
Every other dynamically-tunable server config registers a listener on
`DefaultClusterConfigChangeHandler` in `BaseServerStarter` (rate limits, query
killing, JFR, throttlers, ...). Binding the raw provider into HK2 for ad-hoc
per-request reads creates a second config-distribution pattern: any resource
can now inject the raw map and invent its own precedence/parsing. A listener
maintaining a validated volatile timeout would also validate once per change
(fixing the per-request WARN noted above) and keep Helix concerns out of JAX-RS
resources. Not a blocker, but worth considering before the pattern spreads.
##########
pinot-server/src/test/java/org/apache/pinot/server/api/resources/ReingestionResourceTest.java:
##########
@@ -0,0 +1,156 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.server.api.resources;
+
+import java.util.HashMap;
+import java.util.Map;
+import javax.ws.rs.core.Response;
+import org.apache.helix.model.ClusterConfig;
+import org.apache.pinot.common.config.DefaultClusterConfigChangeHandler;
+import org.apache.pinot.server.api.BaseResourceTest;
+import org.apache.pinot.spi.env.PinotConfiguration;
+import org.testng.annotations.DataProvider;
+import org.testng.annotations.Test;
+
+import static
org.apache.pinot.server.api.resources.ReingestionResource.CHECK_INTERVAL_MS;
+import static
org.apache.pinot.server.api.resources.ReingestionResource.getConsumptionTimeoutMs;
+import static
org.apache.pinot.server.api.resources.ReingestionResource.waitForCondition;
+import static
org.apache.pinot.spi.utils.CommonConstants.Server.CONFIG_OF_REINGESTION_CONSUMPTION_TIMEOUT_MS;
+import static
org.apache.pinot.spi.utils.CommonConstants.Server.DEFAULT_REINGESTION_CONSUMPTION_TIMEOUT_MS;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+
+/// Tests the consumption timeout resolution of [ReingestionResource] and the
wiring of its injected dependencies.
+public class ReingestionResourceTest extends BaseResourceTest {
Review Comment:
**Test-structure nit: 11 of 12 tests here are pure static-method unit tests,
but pay BaseResourceTest's full fixture.**
`BaseResourceTest.setUp` builds real segments and starts a Grizzly HTTP
server per class; only `testResourceDependenciesAreInjected` needs any of it. A
plain TestNG class for the `getConsumptionTimeoutMs`/`waitForCondition` tests
(keeping the single injection smoke test with the fixture) would cut CI time
and stop segment/server bootstrap changes from failing timeout-parsing tests.
##########
pinot-server/src/main/java/org/apache/pinot/server/api/resources/ReingestionResource.java:
##########
@@ -247,9 +266,38 @@ private void doReingestSegment(String realtimeTableName,
SegmentZKMetadata segme
}
}
- private void waitForCondition(
+ /// Resolves the consumption timeout for a re-ingestion job: the cluster
config takes precedence over the server
+ /// config, which takes precedence over the default. A non-numeric or
non-positive value is ignored with a warning.
+ @VisibleForTesting
+ static long getConsumptionTimeoutMs(PinotClusterConfigProvider
clusterConfigProvider,
+ PinotConfiguration serverConf) {
+ Long timeoutMs = parseTimeoutMs(
+
clusterConfigProvider.getClusterConfigs().get(CONFIG_OF_REINGESTION_CONSUMPTION_TIMEOUT_MS),
"cluster");
+ if (timeoutMs == null) {
+ timeoutMs =
parseTimeoutMs(serverConf.getProperty(CONFIG_OF_REINGESTION_CONSUMPTION_TIMEOUT_MS),
"server");
+ }
+ return timeoutMs != null ? timeoutMs :
DEFAULT_REINGESTION_CONSUMPTION_TIMEOUT_MS;
+ }
+
+ @Nullable
+ private static Long parseTimeoutMs(@Nullable String value, String
configSource) {
+ if (value == null) {
+ return null;
+ }
+ Long timeoutMs = Longs.tryParse(value.trim());
+ if (timeoutMs == null || timeoutMs <= 0) {
+ LOGGER.warn("Ignoring invalid {} config: {}={}, expecting a positive
number of milliseconds", configSource,
+ CONFIG_OF_REINGESTION_CONSUMPTION_TIMEOUT_MS, value);
+ return null;
+ }
+ return timeoutMs;
+ }
+
+ @VisibleForTesting
+ static void waitForCondition(
Function<Void, Boolean> condition, long checkIntervalMs, long timeoutMs,
long gracePeriodMs) {
Review Comment:
**Nit: now that `waitForCondition` is a tested quasi-public static, consider
`BooleanSupplier`.**
The Guava `Function<Void, Boolean>` signature forces
`Boolean.TRUE.equals(condition.apply(null))` and `(Void) -> ...` lambdas at
every call site (four new ones in the test).
`java.util.function.BooleanSupplier` with `condition.getAsBoolean()` expresses
the same thing with no null argument, no boxing, and one less Guava import.
##########
pinot-server/src/main/java/org/apache/pinot/server/starter/helix/BaseServerStarter.java:
##########
@@ -1280,7 +1280,8 @@ private void initSegmentFetcher(PinotConfiguration config)
}
protected AdminApiApplication createServerAdminApp() {
- return new AdminApiApplication(_serverInstance, _accessControlFactory,
_reloadJobStatusCache, _serverConf);
+ return new AdminApiApplication(_serverInstance, _accessControlFactory,
_reloadJobStatusCache,
+ _clusterConfigChangeHandler, _serverConf);
Review Comment:
**[pre-existing] The restart-free contract silently degrades if the Helix
listener registration failed.**
`addClusterfigChangeListener(_clusterConfigChangeHandler)` earlier in
`start()` catches and only logs failures; in that case `getClusterConfigs()`
returns `Map.of()` forever, so the cluster layer of the new timeout resolution
(and the javadoc's "picked up without a restart" promise) is silently dead
while the server serves re-ingestion requests from the startup server-conf copy
/ default. Pre-existing pattern, but the new documented contract now depends on
it — worth at least a WARN at resolution time or a note in the javadoc.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]