suvodeep-pyne commented on code in PR #19737:
URL: https://github.com/apache/pinot/pull/19737#discussion_r4201362176


##########
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");

Review Comment:
   Done in 7c753f3b71. The javadoc now says that removing *or invalidating* the 
cluster value falls back to the server config, which includes the copy taken at 
startup. With the listener, the WARN names the cluster config as its source and 
is followed by an INFO with the effective timeout (`Updated re-ingestion 
consumption timeout to: ...ms`), so the value it fell back to shows in the log.



##########
pinot-spi/src/main/java/org/apache/pinot/spi/utils/CommonConstants.java:
##########
@@ -1754,6 +1754,15 @@ public static class Server {
         "ingestion.oom.protection.gcIntervalMs";
     public static final long 
DEFAULT_SERVER_INGESTION_OOM_PROTECTION_GC_INTERVAL_MS = 30_000L;
 
+    /// Max time a pauseless segment re-ingestion job (server API `POST 
/reingestSegment/{segmentName}`) waits for
+    /// consumption to reach the segment end offset before failing. Can also 
be set via cluster config, which takes
+    /// precedence over the server config and is picked up without a restart. 
Read when each re-ingestion request is
+    /// accepted. Note that cluster configs present at server startup are also 
copied into the server config, so to

Review Comment:
   Added to the javadoc in 7c753f3b71: only the exact 
`pinot.server.reingestion.consumption.timeoutMs` key is applied without a 
restart, not a `pinot.all.` prefixed or differently cased one. This matches the 
other `PinotClusterConfigChangeListener`s, which also look up exact keys.



##########
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,

Review Comment:
   Intentional. The cluster config is the live operator knob, e.g. to raise the 
timeout on every server during a repair without a restart. It's also how the 
other dynamic server configs behave at runtime: 
`KeepPipelineBreakerStatsPredicate`, for example, starts from the server config 
and is overwritten by the cluster value on registration and on every change. I 
added that rationale to the javadoc in 7c753f3b71.



##########
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:
   Removed in 7c753f3b71: the parameter and its branch are gone, since the only 
caller passed 0.



##########
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:
   Fixed in 7c753f3b71 by moving to a listener. The value is parsed and 
validated once when the config loads or changes, so an invalid value logs one 
WARN per change instead of one per request.



-- 
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]

Reply via email to