xiangfu0 commented on code in PR #19737:
URL: https://github.com/apache/pinot/pull/19737#discussion_r4194100811
##########
pinot-server/src/main/java/org/apache/pinot/server/api/resources/ReingestionResource.java:
##########
@@ -261,13 +309,18 @@ private void waitForCondition(
}
}
- while (System.currentTimeMillis() < endTime) {
+ while (true) {
try {
if (Boolean.TRUE.equals(condition.apply(null))) {
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));
Review Comment:
**Interrupt is swallowed and misreported in the rewritten wait loop.**
`catch (Exception e)` wraps `InterruptedException` from `Thread.sleep` into
`RuntimeException("Caught exception while checking the condition", e)` without
restoring the interrupt flag (the grace-period block above has the same issue).
Pre-existing pattern, but this PR rewrote these lines and makes multi-hour
sleeps configurable: an executor shutdown/cancellation interrupt during the
wait is cleared, subsequent blocking calls in the unwind path (writer close,
segment build semaphore) no longer observe it, and the log falsely claims the
condition check threw. Suggest catching `InterruptedException` separately and
calling `Thread.currentThread().interrupt()` before rethrowing.
##########
pinot-server/src/main/java/org/apache/pinot/server/api/resources/ReingestionResource.java:
##########
@@ -189,6 +204,8 @@ public Response reingestSegment(@PathParam("segmentName")
@Encoded String encode
}
IndexLoadingConfig indexLoadingConfig =
tableDataManager.fetchIndexLoadingConfig();
+ long consumptionTimeoutMs = getConsumptionTimeoutMs(_clusterConfigProvider,
+ (PinotConfiguration)
_application.getProperties().get(AdminApiApplication.PINOT_CONFIGURATION));
// Check if this segment is already being re-ingested
AtomicBoolean isIngesting =
_reingestingSegments.computeIfAbsent(segmentName, k -> new
AtomicBoolean(false));
Review Comment:
**[pre-existing, surfaced by this PR] Resource is created per request, so
all of its stateful fields are ineffective and each POST leaks a non-daemon
thread.**
`ReingestionResource` has no `@Singleton` and is served via package
scanning, so Jersey instantiates it per request — the new test's own comment
says "The resource is created per request". That means:
- `_reingestingSegments` can never reject a duplicate across requests: two
concurrent `POST /reingestSegment/{seg}` each get a fresh map, both CAS
false→true, and the same segment is consumed and uploaded twice in parallel
(the 409 path is dead code).
- `GET /reingestSegment/jobs` always returns `[]` — it reads a brand-new
instance's empty `_runningJobs`.
- every accepted POST creates a new fixed thread pool whose non-daemon
`reingestion-worker` thread is never shut down — a permanent thread leak per
request.
This PR makes it worse in practice: with multi-hour timeouts now
configurable, the controller's
`RealtimeSegmentValidationManager`/`repairSegmentsInErrorState` loop will
re-POST the same stuck segment on every validation cycle while a prior attempt
is still consuming, and none of those duplicates are rejected. Consider adding
`@Singleton` (cf. `ControllerJobStatusResource`) as part of this change.
##########
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 {
+ private static final PinotConfiguration EMPTY_SERVER_CONF = new
PinotConfiguration();
+ private static final PinotConfiguration SERVER_CONF =
+ new
PinotConfiguration(Map.of(CONFIG_OF_REINGESTION_CONSUMPTION_TIMEOUT_MS,
"3600000"));
+
+ @Test
+ public void testResourceDependenciesAreInjected() {
+ // The resource is created per request, so any request fails if one of its
injected dependencies is not bound
+ Response response =
_webTarget.path("/reingestSegment/jobs").request().get(Response.class);
+
assertThat(response.getStatus()).isEqualTo(Response.Status.OK.getStatusCode());
+ }
+
+ @Test
+ public void testConsumptionTimeoutDefaultsTo30Minutes() {
+ assertThat(DEFAULT_REINGESTION_CONSUMPTION_TIMEOUT_MS).isEqualTo(30 * 60 *
1000L);
+ assertThat(getConsumptionTimeoutMs(clusterConfigProvider(Map.of()),
EMPTY_SERVER_CONF)).isEqualTo(
+ DEFAULT_REINGESTION_CONSUMPTION_TIMEOUT_MS);
+ }
+
+ @Test
+ public void testConsumptionTimeoutFromServerConfig() {
+ assertThat(getConsumptionTimeoutMs(clusterConfigProvider(Map.of()),
SERVER_CONF)).isEqualTo(3_600_000L);
+ }
+
+ @Test
+ public void testClusterConfigOverridesServerConfig() {
+ DefaultClusterConfigChangeHandler clusterConfigProvider =
+
clusterConfigProvider(Map.of(CONFIG_OF_REINGESTION_CONSUMPTION_TIMEOUT_MS,
"7200000"));
+ assertThat(getConsumptionTimeoutMs(clusterConfigProvider,
SERVER_CONF)).isEqualTo(7_200_000L);
+ assertThat(getConsumptionTimeoutMs(clusterConfigProvider,
EMPTY_SERVER_CONF)).isEqualTo(7_200_000L);
+ }
+
+ @DataProvider
+ public Object[][] validTimeouts() {
+ return new Object[][]{{"1", 1L}, {" 7200000 ", 7_200_000L},
{String.valueOf(Long.MAX_VALUE), Long.MAX_VALUE}};
+ }
+
+ @Test(dataProvider = "validTimeouts")
+ public void testValidTimeouts(String timeout, long expectedTimeoutMs) {
+ PinotConfiguration serverConf =
+ new
PinotConfiguration(Map.of(CONFIG_OF_REINGESTION_CONSUMPTION_TIMEOUT_MS,
timeout));
+ assertThat(getConsumptionTimeoutMs(clusterConfigProvider(Map.of()),
serverConf)).isEqualTo(expectedTimeoutMs);
+ assertThat(getConsumptionTimeoutMs(
+
clusterConfigProvider(Map.of(CONFIG_OF_REINGESTION_CONSUMPTION_TIMEOUT_MS,
timeout)), SERVER_CONF)).isEqualTo(
+ expectedTimeoutMs);
+ }
+
+ @DataProvider
+ public Object[][] invalidTimeouts() {
+ return new Object[][]{{"abc"}, {""}, {"0"}, {"-1"}, {"1.5"}, {"30m"},
{"9223372036854775808"}};
+ }
+
+ @Test(dataProvider = "invalidTimeouts")
+ public void testInvalidClusterConfigFallsBackToServerConfig(String
invalidTimeout) {
+ DefaultClusterConfigChangeHandler clusterConfigProvider =
+
clusterConfigProvider(Map.of(CONFIG_OF_REINGESTION_CONSUMPTION_TIMEOUT_MS,
invalidTimeout));
+ assertThat(getConsumptionTimeoutMs(clusterConfigProvider,
SERVER_CONF)).isEqualTo(3_600_000L);
+ }
+
+ @Test(dataProvider = "invalidTimeouts")
+ public void testInvalidServerConfigFallsBackToDefault(String invalidTimeout)
{
+ PinotConfiguration serverConf =
+ new
PinotConfiguration(Map.of(CONFIG_OF_REINGESTION_CONSUMPTION_TIMEOUT_MS,
invalidTimeout));
+ assertThat(getConsumptionTimeoutMs(clusterConfigProvider(Map.of()),
serverConf)).isEqualTo(
+ DEFAULT_REINGESTION_CONSUMPTION_TIMEOUT_MS);
+ }
+
+ @Test
+ public void testClusterConfigChangeIsPickedUpWithoutRestart() {
+ DefaultClusterConfigChangeHandler clusterConfigProvider =
clusterConfigProvider(Map.of());
+ assertThat(getConsumptionTimeoutMs(clusterConfigProvider,
SERVER_CONF)).isEqualTo(3_600_000L);
+
+ // Set the cluster config on the same provider
+ setClusterConfigs(clusterConfigProvider,
Map.of(CONFIG_OF_REINGESTION_CONSUMPTION_TIMEOUT_MS, "10800000"));
+ assertThat(getConsumptionTimeoutMs(clusterConfigProvider,
SERVER_CONF)).isEqualTo(10_800_000L);
+
+ // Remove the cluster config, which falls back to the server config
+ setClusterConfigs(clusterConfigProvider, Map.of());
+ assertThat(getConsumptionTimeoutMs(clusterConfigProvider,
SERVER_CONF)).isEqualTo(3_600_000L);
+ }
+
+ @Test
+ public void testWaitForConditionWithMaxTimeoutDoesNotOverflow() {
+ // Without saturation, the deadline overflows to the past and the wait
times out before checking the condition
+ waitForCondition(v -> true, 1, Long.MAX_VALUE, 0);
Review Comment:
**This test cannot detect removal of the overflow guard.**
With `v -> true`, the rewritten loop returns on the first condition check
before the deadline is ever consulted, so this passes even if
`LongMath.saturatedAdd` is replaced with a plain overflowing `+`. The comment
("the wait times out before checking the condition") describes the old
deadline-first loop, not the shipped one. To pin the saturation behavior, use a
condition that is false on the first check (e.g. flips true after ~50ms) with
`timeoutMs = Long.MAX_VALUE` — an overflowed negative `endTime` would then
throw `Timeout` and fail the test.
##########
pinot-server/src/main/java/org/apache/pinot/server/api/resources/ReingestionResource.java:
##########
@@ -189,6 +204,8 @@ public Response reingestSegment(@PathParam("segmentName")
@Encoded String encode
}
IndexLoadingConfig indexLoadingConfig =
tableDataManager.fetchIndexLoadingConfig();
+ long consumptionTimeoutMs = getConsumptionTimeoutMs(_clusterConfigProvider,
+ (PinotConfiguration)
_application.getProperties().get(AdminApiApplication.PINOT_CONFIGURATION));
Review Comment:
**Fragile, untested wiring: string-key properties lookup + unchecked cast,
dereferenced with no null check.**
`getConsumptionTimeoutMs` calls `serverConf.getProperty(...)` on whatever
the cast yields; if `AdminApiApplication`'s `property(PINOT_CONFIGURATION,
serverConf)` is ever renamed/removed or the resource is hosted without it,
every POST becomes a request-time NPE/500 — and no test would notice:
`testResourceDependenciesAreInjected` only GETs `/reingestSegment/jobs`
(validating the HK2 bindings, not this property), and the unit tests construct
configs directly. Since this PR already edits the binder in
`AdminApiApplication.configure()`, a one-line
`bind(serverConf).to(PinotConfiguration.class)` + `@Inject PinotConfiguration`
would fail at startup instead of request time and drop the `_application`
field, the cast, and the `AdminApiApplication` import.
##########
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:
**An invalid live cluster value silently falls back to the stale startup
cluster snapshot, not the default.**
If the key was present in cluster config when the server started,
`ServiceStartableUtils.applyClusterConfig` copied that value into the server
conf. When the live cluster value later becomes invalid (fat-fingered edit),
resolution falls through to the "server" layer — which is actually the old
cluster value frozen at startup — and any WARN for it is attributed to the
server config, misleading diagnosis. The javadoc on the constant covers the
key-deletion case but not this invalid-value case; consider extending the note
("update the cluster config value instead of removing it **or leaving it
invalid**").
##########
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:
**Only this exact key spelling is picked up without a restart — worth
stating in the javadoc.**
The live lookup is an exact-key `Map.get` on the raw Helix cluster-config
map, while the startup path is lenient: `applyClusterConfig` rewrites
`pinot.all.*` to `pinot.server.*` and `PinotConfiguration` applies relaxed key
matching (case/`-`/`_`). So a cluster config set as
`pinot.all.reingestion.consumption.timeoutMs` (or with different casing) takes
effect at startup via the server-conf copy, then silently stops being dynamic —
later edits to it are invisible to the live read. A one-line note that
restart-free pickup requires the exact `pinot.server.` spelling would save
operators a confusing debugging session.
##########
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:
**Cluster-over-server precedence inverts the platform's usual semantics —
intentional?**
Everywhere else, the server's own config wins:
`ServiceStartableUtils.applyClusterConfig` only copies a cluster value into the
instance config when the key is *not* already set. Here the cluster value
always beats the server conf at runtime, so a per-instance override (e.g. a
beefier server given a larger timeout in its server conf) is silently ignored
once any cluster-wide value exists — this key alone resolves with unique
precedence. If intentional (cluster config as the live operator knob), a short
rationale in the javadoc would help; otherwise consider server-conf-wins for
consistency.
--
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]