This is an automated email from the ASF dual-hosted git repository.
Jackie-Jiang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pinot.git
The following commit(s) were added to refs/heads/master by this push:
new 84becf5957d Force-expire stranded segments in pauseless ingestion
tests instead of shortening the completion time (#19292)
84becf5957d is described below
commit 84becf5957d01d3a1b6d0f13216c21cb6e2ce5cd
Author: Xiaotian (Jackie) Jiang <[email protected]>
AuthorDate: Tue Aug 18 22:31:44 2026 -0700
Force-expire stranded segments in pauseless ingestion tests instead of
shortening the completion time (#19292)
---
.../tests/BasePauselessRealtimeIngestionTest.java | 27 ++++++-------
.../PauselessRealtimeIngestionIntegrationTest.java | 1 -
...sRealtimeIngestionSegmentCommitFailureTest.java | 47 +++++++++++++---------
.../TableRebalancePauselessIntegrationTest.java | 1 -
...ureInjectingPinotLLCRealtimeSegmentManager.java | 30 +++++++-------
.../realtime/utils/PauselessRealtimeTestUtils.java | 26 ++++++++++++
6 files changed, 80 insertions(+), 52 deletions(-)
diff --git
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/BasePauselessRealtimeIngestionTest.java
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/BasePauselessRealtimeIngestionTest.java
index 6fd0f28cb56..be447d786f9 100644
---
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/BasePauselessRealtimeIngestionTest.java
+++
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/BasePauselessRealtimeIngestionTest.java
@@ -56,7 +56,6 @@ public abstract class BasePauselessRealtimeIngestionTest
extends BaseClusterInte
protected static final int NUM_REALTIME_SEGMENTS = 48;
protected static final long DEFAULT_COUNT_STAR_RESULT = 115545L;
protected static final String DEFAULT_TABLE_NAME_2 = DEFAULT_TABLE_NAME +
"_2";
- private static final long MAX_SEGMENT_COMPLETION_TIME_MILLIS = 10_000;
protected List<File> _avroFiles;
protected boolean _failureEnabled = false;
@@ -106,7 +105,6 @@ public abstract class BasePauselessRealtimeIngestionTest
extends BaseClusterInte
startController();
startBroker();
startServer();
- setMaxSegmentCompletionTimeMillis();
setupNonPauselessTable();
injectFailure();
setupPauselessTable();
@@ -156,14 +154,6 @@ public abstract class BasePauselessRealtimeIngestionTest
extends BaseClusterInte
addTableConfig(tableConfig);
}
- protected void setMaxSegmentCompletionTimeMillis() {
- PinotLLCRealtimeSegmentManager realtimeSegmentManager =
_helixResourceManager.getRealtimeSegmentManager();
- if (realtimeSegmentManager instanceof
FailureInjectingPinotLLCRealtimeSegmentManager) {
- ((FailureInjectingPinotLLCRealtimeSegmentManager) realtimeSegmentManager)
-
.setMaxSegmentCompletionTimeoutMs(MAX_SEGMENT_COMPLETION_TIME_MILLIS);
- }
- }
-
protected void injectFailure() {
PinotLLCRealtimeSegmentManager realtimeSegmentManager =
_helixResourceManager.getRealtimeSegmentManager();
if (realtimeSegmentManager instanceof
FailureInjectingPinotLLCRealtimeSegmentManager) {
@@ -210,12 +200,21 @@ public abstract class BasePauselessRealtimeIngestionTest
extends BaseClusterInte
return segmentZKMetadataList.size() ==
getExpectedZKMetadataWithFailure();
}, 1000, 100000, "New Segment ZK Metadata not created");
- Thread.sleep(MAX_SEGMENT_COMPLETION_TIME_MILLIS);
disableFailure();
- Properties periodicTaskProperties = new Properties();
-
periodicTaskProperties.setProperty(ControllerPeriodicTask.RUN_SEGMENT_LEVEL_VALIDATION,
Boolean.TRUE.toString());
-
_controllerStarter.getRealtimeSegmentValidationManager().run(periodicTaskProperties);
+ // Force-expire the segments stranded by the injected failure instead of
waiting out the max segment completion
+ // time: they become immediately eligible for repair, and their in-flight
commit attempts keep getting rejected,
+ // so the single validation run below stays the only recovery path under
test. Segments created afterwards keep
+ // the production completion time and are never artificially rejected.
Clear the expiry right after the run so
+ // that segments re-activated by the repair can commit normally.
+ PauselessRealtimeTestUtils.forceExpireSegments(_helixResourceManager,
tableNameWithType);
+ try {
+ Properties periodicTaskProperties = new Properties();
+
periodicTaskProperties.setProperty(ControllerPeriodicTask.RUN_SEGMENT_LEVEL_VALIDATION,
Boolean.TRUE.toString());
+
_controllerStarter.getRealtimeSegmentValidationManager().run(periodicTaskProperties);
+ } finally {
+
PauselessRealtimeTestUtils.clearForceExpiredSegments(_helixResourceManager);
+ }
waitForAllDocsLoaded(600_000L);
waitForAllDocsLoaded(tableNameWithType2, 600_000L);
diff --git
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/PauselessRealtimeIngestionIntegrationTest.java
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/PauselessRealtimeIngestionIntegrationTest.java
index e09fddc54bd..2133bcff775 100644
---
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/PauselessRealtimeIngestionIntegrationTest.java
+++
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/PauselessRealtimeIngestionIntegrationTest.java
@@ -89,7 +89,6 @@ public class PauselessRealtimeIngestionIntegrationTest
extends BasePauselessReal
startController();
startBroker();
startServer();
- setMaxSegmentCompletionTimeMillis();
// setupNonPauselessTable() also unpacks and publishes the source records.
All scenarios consume the immutable
// topic from the smallest offset and compare their recovered metadata
with this reference table.
_scenario = Scenario.NO_FAILURE;
diff --git
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/PauselessRealtimeIngestionSegmentCommitFailureTest.java
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/PauselessRealtimeIngestionSegmentCommitFailureTest.java
index eb3fae80a3f..66050ddaa6d 100644
---
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/PauselessRealtimeIngestionSegmentCommitFailureTest.java
+++
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/PauselessRealtimeIngestionSegmentCommitFailureTest.java
@@ -31,9 +31,7 @@ import org.apache.pinot.common.utils.LLCSegmentName;
import org.apache.pinot.controller.BaseControllerStarter;
import org.apache.pinot.controller.ControllerConf;
import
org.apache.pinot.controller.helix.core.periodictask.ControllerPeriodicTask;
-import
org.apache.pinot.controller.helix.core.realtime.PinotLLCRealtimeSegmentManager;
import
org.apache.pinot.integration.tests.realtime.utils.FailureInjectingControllerStarter;
-import
org.apache.pinot.integration.tests.realtime.utils.FailureInjectingPinotLLCRealtimeSegmentManager;
import
org.apache.pinot.integration.tests.realtime.utils.FailureInjectingTableConfig;
import
org.apache.pinot.integration.tests.realtime.utils.FailureInjectingTableDataManagerProvider;
import org.apache.pinot.server.starter.helix.HelixInstanceDataManagerConfig;
@@ -60,7 +58,7 @@ import static org.testng.Assert.assertNotNull;
/// Verifies recovery from server-side segment commit and consuming-transition
failures with one shared cluster.
public class PauselessRealtimeIngestionSegmentCommitFailureTest extends
BaseClusterIntegrationTest {
private static final String REFERENCE_TABLE_NAME = DEFAULT_TABLE_NAME +
"_reference";
- private static final long MAX_SEGMENT_COMPLETION_TIME_MILLIS = 10_000L;
+ private static final long VALIDATION_RERUN_INTERVAL_MS = 10_000L;
private FailureScenario _failureScenario;
private File _sampleAvroFile;
@@ -112,7 +110,6 @@ public class
PauselessRealtimeIngestionSegmentCommitFailureTest extends BaseClus
List<File> avroFiles = unpackAvroData(_tempDir);
_sampleAvroFile = avroFiles.get(0);
pushAvroIntoKafka(avroFiles);
- setMaxSegmentCompletionTimeMillis();
setupReferenceTable();
}
@@ -237,28 +234,32 @@ public class
PauselessRealtimeIngestionSegmentCommitFailureTest extends BaseClus
tableConfig.getValidationConfig().setRetentionTimeValue("100000");
}
- private void setMaxSegmentCompletionTimeMillis() {
- PinotLLCRealtimeSegmentManager realtimeSegmentManager =
_helixResourceManager.getRealtimeSegmentManager();
- if (realtimeSegmentManager instanceof
FailureInjectingPinotLLCRealtimeSegmentManager) {
- ((FailureInjectingPinotLLCRealtimeSegmentManager)
realtimeSegmentManager).setMaxSegmentCompletionTimeoutMs(
- MAX_SEGMENT_COMPLETION_TIME_MILLIS);
- }
- }
-
- private void verifyRecovery(FailureScenario failureScenario)
- throws Exception {
+ private void verifyRecovery(FailureScenario failureScenario) {
String pauselessTableName =
TableNameBuilder.REALTIME.tableNameWithType(getPauselessTableName(failureScenario));
List<String> erroredSegments = getSegmentsInEV(pauselessTableName,
SegmentStateModel.ERROR);
assertFalse(erroredSegments.isEmpty(), "No segments found in ERROR state,
expected at least one.");
- Thread.sleep(MAX_SEGMENT_COMPLETION_TIME_MILLIS);
- Properties periodicTaskProperties = new Properties();
-
periodicTaskProperties.setProperty(ControllerPeriodicTask.RUN_SEGMENT_LEVEL_VALIDATION,
Boolean.TRUE.toString());
-
_controllerStarter.getRealtimeSegmentValidationManager().run(periodicTaskProperties);
+ // Segments in ERROR state are repaired by the segment-level validation
(re-ingestion for COMMITTING segments,
+ // reset for IN_PROGRESS segments), which runs periodically in production.
Run it periodically here as well
+ // instead of asserting on a single pass: the failure injection keeps
producing new ERROR states past any single
+ // pass — replicas hit the injected failures independently while the table
is still consuming, and every reset of
+ // a consuming segment creates a new segment data manager that draws from
the remaining failure budget. The
+ // re-runs are spaced out so that a pass does not re-trigger repairs
(re-ingestion, reset) that are still in
+ // flight from the previous pass, while the ERROR-empty check itself polls
at the usual 1s granularity.
+ long[] lastValidationRunMs = new long[1];
+ TestUtils.waitForCondition(aVoid -> {
+ if (getSegmentsInEV(pauselessTableName,
SegmentStateModel.ERROR).isEmpty()) {
+ return true;
+ }
+ long nowMs = System.currentTimeMillis();
+ if (nowMs - lastValidationRunMs[0] >= VALIDATION_RERUN_INTERVAL_MS) {
+ lastValidationRunMs[0] = nowMs;
+ runSegmentLevelValidation();
+ }
+ return false;
+ }, 1000L, 600_000L, "Some segments are still in ERROR state after repeated
validation runs");
- TestUtils.waitForCondition(aVoid -> getSegmentsInEV(pauselessTableName,
SegmentStateModel.ERROR).isEmpty(),
- 600_000L, "Some segments are still in ERROR state after
resetSegments()");
TestUtils.waitForCondition(aVoid -> getSegmentsInEV(pauselessTableName,
SegmentStateModel.OFFLINE).isEmpty(),
30_000L, "Some segments are in OFFLINE state after resetSegments()");
@@ -267,6 +268,12 @@ public class
PauselessRealtimeIngestionSegmentCommitFailureTest extends BaseClus
TableNameBuilder.REALTIME.tableNameWithType(getNonPauselessTableName())));
}
+ private void runSegmentLevelValidation() {
+ Properties periodicTaskProperties = new Properties();
+
periodicTaskProperties.setProperty(ControllerPeriodicTask.RUN_SEGMENT_LEVEL_VALIDATION,
Boolean.TRUE.toString());
+
_controllerStarter.getRealtimeSegmentValidationManager().run(periodicTaskProperties);
+ }
+
private List<String> getSegmentsInEV(String realtimeTableName, String
status) {
ExternalView externalView = _helixResourceManager.getHelixAdmin()
.getResourceExternalView(_helixResourceManager.getHelixClusterName(),
realtimeTableName);
diff --git
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/TableRebalancePauselessIntegrationTest.java
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/TableRebalancePauselessIntegrationTest.java
index 840665a2b18..555efa6c60c 100644
---
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/TableRebalancePauselessIntegrationTest.java
+++
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/TableRebalancePauselessIntegrationTest.java
@@ -82,7 +82,6 @@ public class TableRebalancePauselessIntegrationTest extends
BasePauselessRealtim
startServer();
createServerTenant(getServerTenant(), 0, 1);
createBrokerTenant(getBrokerTenant(), 1);
- setMaxSegmentCompletionTimeMillis();
setupNonPauselessTable();
injectFailure();
setupPauselessTable();
diff --git
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/realtime/utils/FailureInjectingPinotLLCRealtimeSegmentManager.java
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/realtime/utils/FailureInjectingPinotLLCRealtimeSegmentManager.java
index 5e6ad70ac63..42e175f183e 100644
---
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/realtime/utils/FailureInjectingPinotLLCRealtimeSegmentManager.java
+++
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/realtime/utils/FailureInjectingPinotLLCRealtimeSegmentManager.java
@@ -19,23 +19,21 @@
package org.apache.pinot.integration.tests.realtime.utils;
import com.google.common.annotations.VisibleForTesting;
+import java.util.Collection;
import java.util.Map;
+import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import org.apache.pinot.common.metrics.ControllerMetrics;
import org.apache.pinot.controller.ControllerConf;
import org.apache.pinot.controller.helix.core.PinotHelixResourceManager;
import
org.apache.pinot.controller.helix.core.realtime.PinotLLCRealtimeSegmentManager;
import org.apache.pinot.controller.helix.core.util.FailureInjectionUtils;
-import org.apache.zookeeper.data.Stat;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
public class FailureInjectingPinotLLCRealtimeSegmentManager extends
PinotLLCRealtimeSegmentManager {
@VisibleForTesting
private final Map<String, String> _failureConfig;
- private long _maxSegmentCompletionTimeoutMs = 300000L;
- private static final Logger LOGGER =
LoggerFactory.getLogger(FailureInjectingPinotLLCRealtimeSegmentManager.class);
+ private final Set<String> _forceExpiredSegments =
ConcurrentHashMap.newKeySet();
public
FailureInjectingPinotLLCRealtimeSegmentManager(PinotHelixResourceManager
helixResourceManager,
ControllerConf controllerConf, ControllerMetrics controllerMetrics) {
@@ -43,22 +41,22 @@ public class FailureInjectingPinotLLCRealtimeSegmentManager
extends PinotLLCReal
_failureConfig = new ConcurrentHashMap<>();
}
+ /// Treats force-expired segments as exceeding the max segment completion
time regardless of how recently their
+ /// segment ZK metadata was updated. This makes them immediately eligible
for repair while keeping their in-flight
+ /// commit attempts rejected, without shortening the completion time for any
other segment.
@Override
protected boolean isExceededMaxSegmentCompletionTime(String
realtimeTableName, String segmentName,
long currentTimeMs) {
- Stat stat = new Stat();
- getSegmentZKMetadata(realtimeTableName, segmentName, stat);
- if (currentTimeMs > stat.getMtime() + _maxSegmentCompletionTimeoutMs) {
- LOGGER.info("Segment: {} exceeds the max completion time: {}ms, metadata
update time: {}, current time: {}",
- segmentName, _maxSegmentCompletionTimeoutMs, stat.getMtime(),
currentTimeMs);
- return true;
- } else {
- return false;
- }
+ return _forceExpiredSegments.contains(segmentName)
+ || super.isExceededMaxSegmentCompletionTime(realtimeTableName,
segmentName, currentTimeMs);
+ }
+
+ public void forceExpireSegments(Collection<String> segmentNames) {
+ _forceExpiredSegments.addAll(segmentNames);
}
- public void setMaxSegmentCompletionTimeoutMs(long
maxSegmentCompletionTimeoutMs) {
- _maxSegmentCompletionTimeoutMs = maxSegmentCompletionTimeoutMs;
+ public void clearForceExpiredSegments() {
+ _forceExpiredSegments.clear();
}
@VisibleForTesting
diff --git
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/realtime/utils/PauselessRealtimeTestUtils.java
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/realtime/utils/PauselessRealtimeTestUtils.java
index fc7f5264c01..16b3cf5e776 100644
---
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/realtime/utils/PauselessRealtimeTestUtils.java
+++
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/realtime/utils/PauselessRealtimeTestUtils.java
@@ -18,6 +18,7 @@
*/
package org.apache.pinot.integration.tests.realtime.utils;
+import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@@ -26,6 +27,8 @@ import org.apache.helix.model.IdealState;
import org.apache.pinot.common.metadata.segment.SegmentZKMetadata;
import org.apache.pinot.common.utils.LLCSegmentName;
import org.apache.pinot.common.utils.helix.HelixHelper;
+import org.apache.pinot.controller.helix.core.PinotHelixResourceManager;
+import
org.apache.pinot.controller.helix.core.realtime.PinotLLCRealtimeSegmentManager;
import org.apache.pinot.spi.utils.CommonConstants;
import static org.testng.Assert.assertEquals;
@@ -42,6 +45,29 @@ public class PauselessRealtimeTestUtils {
assertEquals(segmentAssignment.size(), numSegmentsExpected);
}
+ /// Marks all current segments of the given table as exceeding the max
segment completion time, making them
+ /// immediately eligible for repair by the next validation run while keeping
their in-flight commit attempts
+ /// rejected. Segments created after this call are not affected. No-op when
the cluster is not started with a
+ /// [FailureInjectingPinotLLCRealtimeSegmentManager].
+ public static void forceExpireSegments(PinotHelixResourceManager
helixResourceManager, String tableNameWithType) {
+ PinotLLCRealtimeSegmentManager realtimeSegmentManager =
helixResourceManager.getRealtimeSegmentManager();
+ if (realtimeSegmentManager instanceof
FailureInjectingPinotLLCRealtimeSegmentManager) {
+ List<String> segmentNames = new ArrayList<>();
+ for (SegmentZKMetadata segmentZKMetadata :
helixResourceManager.getSegmentsZKMetadata(tableNameWithType)) {
+ segmentNames.add(segmentZKMetadata.getSegmentName());
+ }
+ ((FailureInjectingPinotLLCRealtimeSegmentManager)
realtimeSegmentManager).forceExpireSegments(segmentNames);
+ }
+ }
+
+ /// Clears the force-expired segments so that segments re-activated by the
repair can complete normally.
+ public static void clearForceExpiredSegments(PinotHelixResourceManager
helixResourceManager) {
+ PinotLLCRealtimeSegmentManager realtimeSegmentManager =
helixResourceManager.getRealtimeSegmentManager();
+ if (realtimeSegmentManager instanceof
FailureInjectingPinotLLCRealtimeSegmentManager) {
+ ((FailureInjectingPinotLLCRealtimeSegmentManager)
realtimeSegmentManager).clearForceExpiredSegments();
+ }
+ }
+
public static boolean assertUrlPresent(List<SegmentZKMetadata>
segmentZKMetadataList) {
for (SegmentZKMetadata segmentZKMetadata : segmentZKMetadataList) {
if (segmentZKMetadata.getStatus() ==
CommonConstants.Segment.Realtime.Status.COMMITTING
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]