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


##########
pinot-server/src/main/java/org/apache/pinot/server/api/resources/ReingestionResource.java:
##########
@@ -79,12 +82,12 @@
         description = "Database context passed through http header. If no 
context is provided 'default' database "
             + "context will be considered.")}))
 @Path("/")
+@Singleton

Review Comment:
   Done in 13ba976e3a. `ReingestionResource` now has a `@PreDestroy` 
`shutDown()` that calls `shutdownNow()` and logs how many queued jobs were 
dropped. I checked that Jersey does call it for the `@Singleton` resource when 
the admin API stops: 
`ReingestionResourceApiTest.testRunningReingestionIsStoppedOnShutdown` holds a 
job, stops the admin API, and asserts the worker thread exits. It fails without 
the `@PreDestroy`. The worker threads are now daemon threads, and a 
`RejectedExecutionException` from `submit` clears the in-progress tracking and 
returns 503.
   
   While at it, I made `StatelessRealtimeSegmentWriter.stopConsumption()` use 
`Uninterruptibles.joinUninterruptibly`. Before, an interrupt arriving while 
`close()` was joining the consumer made it stop waiting, and the mutable 
segment could then be destroyed while the consumer thread was still indexing. 
The shutdown interrupt makes that reachable.



##########
pinot-server/src/main/java/org/apache/pinot/server/api/resources/ReingestionResource.java:
##########
@@ -229,8 +235,12 @@ private void doReingestSegment(String realtimeTableName, 
SegmentZKMetadata segme
     String segmentName = segmentZKMetadata.getSegmentName();
     try (StatelessRealtimeSegmentWriter writer = new 
StatelessRealtimeSegmentWriter(segmentZKMetadata,
         indexLoadingConfig, segmentBuildSemaphore)) {
+      // Read when the consumption starts, so that a job waiting for a 
re-ingestion thread uses the latest timeout

Review Comment:
   Done in 13ba976e3a. The job is added to `_runningJobs` on the request thread 
before `submit`, and is removed in the task's `finally` or if `submit` is 
rejected. `testQueuedReingestionIsListed` fills every re-ingestion thread and 
asserts the queued job is listed; it fails without the move. The endpoint's 
Swagger text now mentions queued jobs. I left `startTimeMs` as the time the job 
was accepted, which is the existing API field, rather than add a 
consumption-start field in this PR.



##########
pinot-server/src/main/java/org/apache/pinot/server/api/resources/ReingestionConsumptionTimeout.java:
##########
@@ -0,0 +1,84 @@
+/**
+ * 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 com.google.common.primitives.Longs;
+import java.util.Map;
+import java.util.Set;
+import javax.annotation.Nullable;
+import javax.annotation.concurrent.ThreadSafe;
+import org.apache.pinot.spi.config.provider.PinotClusterConfigChangeListener;
+import org.apache.pinot.spi.env.PinotConfiguration;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+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;
+
+
+/// Holds the consumption timeout of the pauseless segment re-ingestion jobs 
run by [ReingestionResource]. The cluster
+/// config takes precedence over the server config, which takes precedence 
over the default. Cluster config changes are
+/// applied as they arrive. A non-numeric or non-positive value is ignored 
with a warning when the config is loaded or
+/// changed.
+///
+/// Thread-safe: cluster config changes are delivered one at a time, and the 
latest value is visible to all readers.
+@ThreadSafe
+public class ReingestionConsumptionTimeout implements 
PinotClusterConfigChangeListener {
+  private static final Logger LOGGER = 
LoggerFactory.getLogger(ReingestionConsumptionTimeout.class);
+
+  // Timeout from the server config or the default, used when the cluster 
config is absent or invalid
+  private final long _serverTimeoutMs;
+  private volatile long _timeoutMs;
+
+  public ReingestionConsumptionTimeout(PinotConfiguration serverConf) {
+    Long serverTimeoutMs =
+        
parseTimeoutMs(serverConf.getProperty(CONFIG_OF_REINGESTION_CONSUMPTION_TIMEOUT_MS),
 "server");
+    _serverTimeoutMs = serverTimeoutMs != null ? serverTimeoutMs : 
DEFAULT_REINGESTION_CONSUMPTION_TIMEOUT_MS;
+    _timeoutMs = _serverTimeoutMs;
+    LOGGER.info("Initialized re-ingestion consumption timeout to: {}ms", 
_timeoutMs);
+  }
+
+  public long getTimeoutMs() {
+    return _timeoutMs;
+  }
+
+  @Override
+  public void onChange(Set<String> changedConfigs, Map<String, String> 
clusterConfigs) {
+    if 
(!changedConfigs.contains(CONFIG_OF_REINGESTION_CONSUMPTION_TIMEOUT_MS)) {
+      return;
+    }
+    Long clusterTimeoutMs = 
parseTimeoutMs(clusterConfigs.get(CONFIG_OF_REINGESTION_CONSUMPTION_TIMEOUT_MS),
 "cluster");
+    _timeoutMs = clusterTimeoutMs != null ? clusterTimeoutMs : 
_serverTimeoutMs;
+    LOGGER.info("Updated re-ingestion consumption timeout to: {}ms", 
_timeoutMs);

Review Comment:
   Done in 13ba976e3a: it now logs `Updated re-ingestion consumption timeout 
from: {}ms to: {}ms`, and only when the effective value changes.



##########
pinot-server/src/main/java/org/apache/pinot/server/api/resources/ReingestionResource.java:
##########
@@ -247,27 +257,27 @@ private void doReingestSegment(String realtimeTableName, 
SegmentZKMetadata segme
     }
   }
 
-  private void waitForCondition(
-      Function<Void, Boolean> condition, long checkIntervalMs, long timeoutMs, 
long gracePeriodMs) {
-    long endTime = System.currentTimeMillis() + timeoutMs;
+  @VisibleForTesting
+  static void waitForCondition(BooleanSupplier condition, long 
checkIntervalMs, long timeoutMs) {
+    // Saturate to avoid overflow, since the timeout is configurable
+    long endTime = LongMath.saturatedAdd(System.currentTimeMillis(), 
timeoutMs);
 
-    // Adding grace period before starting the condition checks
-    if (gracePeriodMs > 0) {
-      LOGGER.info("Waiting for a grace period of {} ms before starting 
condition checks", gracePeriodMs);
+    while (true) {
       try {
-        Thread.sleep(gracePeriodMs);
-      } catch (InterruptedException e) {
-        throw new RuntimeException("Interrupted during grace period wait", e);
-      }
-    }
-
-    while (System.currentTimeMillis() < endTime) {
-      try {
-        if (Boolean.TRUE.equals(condition.apply(null))) {
+        if (condition.getAsBoolean()) {
           LOGGER.info("Condition satisfied: {}", condition);
           return;
         }
-        Thread.sleep(checkIntervalMs);
+        long remainingMs = endTime - System.currentTimeMillis();
+        if (remainingMs <= 0) {
+          break;
+        }
+        // Do not sleep past the deadline, so that the condition is checked 
once more at the deadline before timing out
+        Thread.sleep(Math.min(checkIntervalMs, remainingMs));
+      } catch (InterruptedException e) {
+        // Do not restore the interrupt flag. The caller closes the segment 
writer next, which must wait for the
+        // consumer thread to stop before destroying the segment, and the 
interrupt is preserved as the cause.
+        throw new RuntimeException("Interrupted while waiting for condition: " 
+ condition, e);

Review Comment:
   Done in 13ba976e3a. `waitForCondition` takes a description (`consumption of 
segment: <name>`), and the timeout message is now `Timed out after 
<timeoutMs>ms waiting for consumption of segment: <name>`. The tests assert the 
full messages.



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