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]

Reply via email to