xiangfu0 commented on code in PR #19737:
URL: https://github.com/apache/pinot/pull/19737#discussion_r4212937020


##########
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:
   **[MAJOR] The now-singleton `_reingestionExecutor` has no shutdown path; its 
non-daemon workers outlive server stop.**
   
   Making the resource `@Singleton` gives the fixed pool 
(`reingestion-worker-%d`, `ThreadFactoryBuilder` default daemon=false) a 
well-defined owner, but nothing ever shuts it down: `ReingestionResource` has 
no `@PreDestroy`, `AdminApiApplication.stop()` only calls 
`_httpServer.shutdownNow()`, and `BaseServerStarter.stop()` never touches it. 
Consequences:
   - Once any re-ingestion has run, the idle core threads keep a JVM without 
`System.exit` alive after `stop()` returns (embedded/test usage).
   - In-flight and queued jobs are neither drained nor interrupted on shutdown: 
`_serverInstance.shutDown()` runs at `BaseServerStarter.java:1075`, before 
`_adminApiApplication.stop()` at :1094, so a worker can keep consuming for up 
to the configured timeout (now unbounded via cluster config) and then call 
`uploadReingestedSegment` against a shut-down instance.
   
   Pre-existing in a worse form (a leaked pool per request), but this PR 
creates the single long-lived owner, so it is the natural place to close the 
gap. Suggest a `@PreDestroy` that calls `_reingestionExecutor.shutdownNow()` - 
the wait loop's documented interrupt semantics already turn that into job 
failure with `writer.close()` cleanup - plus `.setDaemon(true)` as a 
belt-and-braces guard. If shutdown is added, wrap the `submit` so a 
`RejectedExecutionException` clears the `_reingestingSegments` flag set just 
before it (today unreachable, reachable once the pool can be shut down).



##########
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:
   **[MINOR] Queued jobs are invisible to `GET /reingestSegment/jobs` while 
duplicate POSTs for them get 409.**
   
   Now that the pool is genuinely shared (`MAX_PARALLEL_REINGESTIONS` can be 1 
on small hosts) with an unbounded queue, "queued" is a real state that can last 
up to a full consumption timeout per job ahead. But `_runningJobs.put(jobId, 
job)` only executes inside the submitted task, while `_reingestingSegments` is 
marked on the request thread before submit. So while a job is queued: POST 
returns 200 with a jobId, a duplicate POST returns 409 "already in progress", 
and `GET /reingestSegment/jobs` returns `[]` - the API simultaneously reports 
"conflict" and "nothing running", and a client treating absence-from-list as 
completion is misled for potentially tens of minutes. (The controller's repair 
loop only relies on the 409, so cluster repair is unaffected - this is 
operator/API observability.)
   
   Secondary: `ReingestionJob._startTimeMs` is set at request time, so a queued 
job's listed runtime overstates its consumption time, while the timeout 
deliberately starts at consumption start.
   
   Suggest moving `_runningJobs.put(jobId, job)` to the request thread just 
before `submit` (the `finally` already removes it), and optionally exposing a 
separate consumption-start time.



##########
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:
   **[MINOR] Config-change log drops the previous value.**
   
   `onChange` logs only the new effective value. When debugging why an 
in-flight re-ingestion used a different timeout than the current cluster config 
(jobs capture the value at consumption start), the previous value is exactly 
what the operator needs, and it is cheap to capture before the assignment. It 
also logs unconditionally even when the effective value did not change (e.g. an 
invalid new value falling back to the same server value).
   
   Suggest `"Updated re-ingestion consumption timeout from: {}ms to: {}ms"`, 
optionally skipping the log when unchanged.



##########
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:
   **[MINOR] The operator-facing wait messages interpolate a lambda whose 
`toString` is meaningless.**
   
   `"Interrupted while waiting for condition: " + condition` (here), 
`LOGGER.info("Condition satisfied: {}", condition)` and `"Timeout waiting for 
condition: " + condition` all render as 
`ReingestionResource$$Lambda$N/0x...@hash`. The timeout message is the one 
surfaced as the re-ingestion job failure, so the error an operator sees carries 
no segment name and not even the timeout that elapsed - and the new test can 
only assert `hasMessageContaining("Timeout")` because the rest of the message 
is unusable.
   
   Suggest passing a human-readable description into `waitForCondition` (e.g. 
`"consumption of segment " + segmentName`) and including `timeoutMs` in the 
timeout message, e.g. `"Timed out after " + timeoutMs + "ms waiting for " + 
description`.



##########
pinot-server/src/test/java/org/apache/pinot/server/api/resources/ReingestionResourceTest.java:
##########
@@ -0,0 +1,67 @@
+/**
+ * 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 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.waitForCondition;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+
+/// Tests how [ReingestionResource] waits for the re-ingestion consumption to 
complete.
+public class ReingestionResourceTest {
+
+  @Test
+  public void testWaitForConditionWithMaxTimeoutDoesNotTimeOutEarly() {
+    // The condition is false on the first check, so the wait must keep 
checking instead of timing out
+    long completionTimeMs = System.currentTimeMillis() + 50;
+    waitForCondition(() -> System.currentTimeMillis() >= completionTimeMs, 10, 
Long.MAX_VALUE);
+  }
+
+  @Test
+  public void 
testWaitForConditionDetectsCompletionWithTimeoutShorterThanCheckInterval() {
+    // The condition becomes true after the first check. It must be checked 
again at the deadline instead of timing out
+    // after sleeping a full check interval.
+    long completionTimeMs = System.currentTimeMillis() + 100;
+    waitForCondition(() -> System.currentTimeMillis() >= completionTimeMs, 
CHECK_INTERVAL_MS, 1_000);

Review Comment:
   **[MINOR] ~900ms stall budget makes this the one new test a single GC/CPU 
pause can flip.**
   
   The test passes only if the loop reaches `Thread.sleep(Math.min(5000, 
remainingMs))` before the 1s deadline expires. If a pause of roughly >900ms 
lands between the first (false) condition check and the `remainingMs` 
computation, `remainingMs <= 0` breaks the loop without re-checking the (by 
then true) condition and the test fails spuriously with "Timeout waiting for 
condition". All other paths are deterministic.
   
   The property under test only needs `timeoutMs < CHECK_INTERVAL_MS`, so 
widening the timeout keeps the regression coverage while raising the tolerated 
pause, e.g. `waitForCondition(() -> System.currentTimeMillis() >= 
completionTimeMs, CHECK_INTERVAL_MS, 4_000)` (~3.9s budget instead of ~0.9s).



##########
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
+      long consumptionTimeoutMs = _consumptionTimeout.getTimeoutMs();

Review Comment:
   **[MINOR] The configured timeout's pass-through into the job is only covered 
compositionally.**
   
   Config resolution (`ReingestionConsumptionTimeoutTest`), the wait loop 
(`ReingestionResourceTest`) and the HK2 binding 
(`testResourceDependenciesAreInjected`) are each tested, but no test observes 
that `doReingestSegment` actually feeds the resolved value into 
`waitForCondition` - a regression swapping this argument back to a constant, or 
reading the timeout at submit time instead of consumption start as the comment 
above promises, would pass every new test. Deferable (the glue is two lines and 
pauseless integration tests exercise the endpoint with the default), but if you 
want to pin it: configure a small timeout via `configureServerConf` in 
`ReingestionResourceApiTest` and assert a job held past it fails with the 
timeout message within the configured bound.



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