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]