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 2354206647a Propagate failed IdealState group commits after queue
drains (#19679)
2354206647a is described below
commit 2354206647a1a4e4e7231441a088931bd0e0a249
Author: Xiaotian (Jackie) Jiang <[email protected]>
AuthorDate: Fri Sep 25 21:14:27 2026 -0700
Propagate failed IdealState group commits after queue drains (#19679)
---
.../common/utils/helix/IdealStateGroupCommit.java | 4 +-
.../helix/IdealStateGroupCommitTest.java | 252 ++++++++++++---------
2 files changed, 143 insertions(+), 113 deletions(-)
diff --git
a/pinot-common/src/main/java/org/apache/pinot/common/utils/helix/IdealStateGroupCommit.java
b/pinot-common/src/main/java/org/apache/pinot/common/utils/helix/IdealStateGroupCommit.java
index 4486c23b0cf..c0a4ab52a90 100644
---
a/pinot-common/src/main/java/org/apache/pinot/common/utils/helix/IdealStateGroupCommit.java
+++
b/pinot-common/src/main/java/org/apache/pinot/common/utils/helix/IdealStateGroupCommit.java
@@ -119,8 +119,8 @@ public class IdealStateGroupCommit {
ArrayList<Entry> processed = new ArrayList<>();
try {
if (queue._pending.peek() == null) {
- // All pending entries have been processed, the updatedIdealState
should be set.
- return entry._updatedIdealState;
+ // The entry was already processed; check its outcome below.
+ break;
}
updateIdealState(helixManager, resourceName, idealState -> {
IdealState updatedIdealState = idealState;
diff --git
a/pinot-controller/src/test/java/org/apache/pinot/controller/helix/IdealStateGroupCommitTest.java
b/pinot-controller/src/test/java/org/apache/pinot/controller/helix/IdealStateGroupCommitTest.java
index 257a03c880a..01535124568 100644
---
a/pinot-controller/src/test/java/org/apache/pinot/controller/helix/IdealStateGroupCommitTest.java
+++
b/pinot-controller/src/test/java/org/apache/pinot/controller/helix/IdealStateGroupCommitTest.java
@@ -31,6 +31,7 @@ import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Function;
+import org.apache.helix.HelixAdmin;
import org.apache.helix.HelixManager;
import org.apache.helix.model.IdealState;
import org.apache.helix.store.zk.ZkHelixPropertyStore;
@@ -45,13 +46,14 @@ import
org.apache.pinot.spi.utils.CommonConstants.Segment.Realtime.Status;
import org.apache.pinot.spi.utils.retry.RetryPolicies;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-import org.testng.Assert;
import org.testng.annotations.AfterClass;
import org.testng.annotations.AfterMethod;
import org.testng.annotations.BeforeClass;
import org.testng.annotations.BeforeMethod;
import org.testng.annotations.Test;
+import static org.testng.Assert.*;
+
public class IdealStateGroupCommitTest {
private static final Logger LOGGER =
LoggerFactory.getLogger(IdealStateGroupCommitTest.class);
@@ -64,6 +66,9 @@ public class IdealStateGroupCommitTest {
private static final int NUM_UPDATES = 100 * SYSTEM_MULTIPLIER;
private static final int NUM_TABLES = 20;
+ private String _clusterName;
+ private HelixManager _helixManager;
+ private HelixAdmin _helixAdmin;
private ExecutorService _executorService;
@BeforeClass
@@ -71,6 +76,9 @@ public class IdealStateGroupCommitTest {
throws Exception {
LOGGER.info("Starting IdealStateGroupCommitTest with SYSTEM_MULTIPLIER:
{}", SYSTEM_MULTIPLIER);
TEST_INSTANCE.setupSharedStateAndValidate();
+ _clusterName = TEST_INSTANCE.getHelixClusterName();
+ _helixManager = TEST_INSTANCE.getHelixManager();
+ _helixAdmin = TEST_INSTANCE.getHelixAdmin();
_executorService = Executors.newFixedThreadPool(4);
}
@@ -83,7 +91,7 @@ public class IdealStateGroupCommitTest {
idealState.setRebalanceMode(IdealState.RebalanceMode.CUSTOMIZED);
idealState.setReplicas("1");
idealState.setNumPartitions(0);
-
TEST_INSTANCE.getHelixAdmin().addResource(TEST_INSTANCE.getHelixClusterName(),
tableName, idealState);
+ _helixAdmin.addResource(_clusterName, tableName, idealState);
ControllerMetrics.get().removeTableMeter(tableName,
ControllerMeter.IDEAL_STATE_UPDATE_SUCCESS);
}
}
@@ -92,7 +100,7 @@ public class IdealStateGroupCommitTest {
public void afterTest() {
for (int i = 0; i < NUM_UPDATES; i++) {
String tableName = TABLE_NAME_PREFIX + i + "_OFFLINE";
-
TEST_INSTANCE.getHelixAdmin().dropResource(TEST_INSTANCE.getHelixClusterName(),
tableName);
+ _helixAdmin.dropResource(_clusterName, tableName);
}
}
@@ -134,7 +142,7 @@ public class IdealStateGroupCommitTest {
initialState.setRebalanceMode(IdealState.RebalanceMode.CUSTOMIZED);
initialState.setReplicas("1");
initialState.setNumPartitions(0);
-
TEST_INSTANCE.getHelixAdmin().addResource(TEST_INSTANCE.getHelixClusterName(),
tableName, initialState);
+ _helixAdmin.addResource(_clusterName, tableName, initialState);
try {
IdealStateGroupCommit commit = new IdealStateGroupCommit();
@@ -152,20 +160,17 @@ public class IdealStateGroupCommitTest {
Throwable freshException = null;
try {
- commit.commit(TEST_INSTANCE.getHelixManager(), tableName, freshUpdater,
- RetryPolicies.noDelayRetryPolicy(1), false);
+ commit.commit(_helixManager, tableName, freshUpdater,
RetryPolicies.noDelayRetryPolicy(1), false);
} catch (Throwable e) {
freshException = e;
}
- LOGGER.info("Fresh thread's commit() returned with exception: {}",
- freshException == null ? "(none)"
- : freshException.getClass().getSimpleName() + ": " +
freshException.getMessage());
+ LOGGER.info("Fresh thread's commit() returned with exception: {}",
freshException == null ? "(none)"
+ : freshException.getClass().getSimpleName() + ": " +
freshException.getMessage());
// The fresh thread's commit() throws (the all-or-nothing batch
semantics are preserved).
- Assert.assertNotNull(freshException,
- "Fresh thread's commit() should throw because the batched IdealState
commit aborts when "
- + "the co-batched stuck entry's lambda throws
PermanentUpdaterException.");
+ assertNotNull(freshException, "Fresh thread's commit() should throw
because the batched IdealState commit aborts "
+ + "when the co-batched stuck entry's lambda throws
PermanentUpdaterException.");
// Step 3: A drainer commit runs. If the fresh thread's entry was still
in _pending and
// NOT cancelled (the pre-fix bug), the drainer would iterate it and
write freshPartition
@@ -175,10 +180,9 @@ public class IdealStateGroupCommitTest {
is.setPartitionState("drainerPartition", "instance1", "ONLINE");
return is;
};
- commit.commit(TEST_INSTANCE.getHelixManager(), tableName, drainerUpdater,
- RetryPolicies.noDelayRetryPolicy(1), false);
+ commit.commit(_helixManager, tableName, drainerUpdater,
RetryPolicies.noDelayRetryPolicy(1), false);
- IdealState finalState =
HelixHelper.getTableIdealState(TEST_INSTANCE.getHelixManager(), tableName);
+ IdealState finalState = HelixHelper.getTableIdealState(_helixManager,
tableName);
Map<String, String> freshMap =
finalState.getInstanceStateMap("freshPartition");
Map<String, String> drainerMap =
finalState.getInstanceStateMap("drainerPartition");
@@ -187,21 +191,20 @@ public class IdealStateGroupCommitTest {
LOGGER.info("drainerPartition state map: {}", drainerMap);
// Sanity: the drainer's own change must be present (proves the drainer
commit ran).
- Assert.assertNotNull(drainerMap, "Drainer commit should have written
drainerPartition.");
- Assert.assertEquals(drainerMap.get("instance1"), "ONLINE");
+ assertNotNull(drainerMap, "Drainer commit should have written
drainerPartition.");
+ assertEquals(drainerMap.get("instance1"), "ONLINE");
// The key assertion: the fresh updater's change must NOT be in
IdealState. The cancellation
// flag must have caused the drainer's iteration to skip and remove the
fresh entry. If this
// assertion fails, the cancellation fix has regressed and the
orphan-creation race can
// reoccur (caller threw, ran cleanup, but subsequent leader applied the
update anyway).
- Assert.assertNull(freshMap,
- "Fresh updater's change must NOT be in IdealState. The fresh
thread's commit() threw, "
- + "indicating to its caller that the update did not happen; a
subsequent leader "
- + "must not have applied it. Found: " + freshMap);
+ assertNull(freshMap, "Fresh updater's change must NOT be in IdealState.
The fresh thread's commit() threw, "
+ + "indicating to its caller that the update did not happen; a
subsequent leader must not have applied it. "
+ + "Found: " + freshMap);
LOGGER.info("=== FIX VERIFIED: caller threw, no future leader applied
the cancelled entry ===");
} finally {
-
TEST_INSTANCE.getHelixAdmin().dropResource(TEST_INSTANCE.getHelixClusterName(),
tableName);
+ _helixAdmin.dropResource(_clusterName, tableName);
}
}
@@ -219,7 +222,7 @@ public class IdealStateGroupCommitTest {
initialState.setRebalanceMode(IdealState.RebalanceMode.CUSTOMIZED);
initialState.setReplicas("1");
initialState.setNumPartitions(0);
-
TEST_INSTANCE.getHelixAdmin().addResource(TEST_INSTANCE.getHelixClusterName(),
tableName, initialState);
+ _helixAdmin.addResource(_clusterName, tableName, initialState);
int stuckCount = 3;
int freshCount = 5;
@@ -245,8 +248,7 @@ public class IdealStateGroupCommitTest {
throw new HelixHelper.PermanentUpdaterException("stuck-" + idx);
};
try {
- commit.commit(TEST_INSTANCE.getHelixManager(), tableName,
stuckUpdater,
- RetryPolicies.noDelayRetryPolicy(1), false);
+ commit.commit(_helixManager, tableName, stuckUpdater,
RetryPolicies.noDelayRetryPolicy(1), false);
} catch (Throwable t) {
result.set(t);
}
@@ -269,34 +271,33 @@ public class IdealStateGroupCommitTest {
return is;
};
try {
- commit.commit(TEST_INSTANCE.getHelixManager(), tableName,
freshUpdater,
- RetryPolicies.noDelayRetryPolicy(1), false);
+ commit.commit(_helixManager, tableName, freshUpdater,
RetryPolicies.noDelayRetryPolicy(1), false);
} catch (Throwable t) {
result.set(t);
}
});
}
- Assert.assertTrue(allReady.await(10, TimeUnit.SECONDS), "Worker threads
failed to start");
+ assertTrue(allReady.await(10, TimeUnit.SECONDS), "Worker threads failed
to start");
goSignal.countDown();
pool.shutdown();
- Assert.assertTrue(pool.awaitTermination(30, TimeUnit.SECONDS), "Workers
did not finish");
+ assertTrue(pool.awaitTermination(30, TimeUnit.SECONDS), "Workers did not
finish");
// After the race, a drainer commit ensures any
successful-but-not-yet-applied entries are
// CAS-written, and any cancelled entries are skipped+removed. This
makes the final
// IdealState the ground truth for which fresh updaters actually took
effect.
- commit.commit(TEST_INSTANCE.getHelixManager(), tableName, is -> {
+ commit.commit(_helixManager, tableName, is -> {
is.setPartitionState("drainer", "instance1", "ONLINE");
return is;
}, RetryPolicies.noDelayRetryPolicy(1), false);
- IdealState finalState =
HelixHelper.getTableIdealState(TEST_INSTANCE.getHelixManager(), tableName);
+ IdealState finalState = HelixHelper.getTableIdealState(_helixManager,
tableName);
// Every stuck owner must observe a throw. The stuck lambda always throws
// PermanentUpdaterException, so the batch they participated in always
aborts.
for (int i = 0; i < stuckCount; i++) {
Throwable t = stuckResults.get(i).get();
- Assert.assertNotNull(t, "Stuck owner " + i + " should have thrown");
+ assertNotNull(t, "Stuck owner " + i + " should have thrown");
}
// Consistency check for every fresh owner: observed outcome must match
IdealState.
@@ -307,30 +308,77 @@ public class IdealStateGroupCommitTest {
Map<String, String> inIs = finalState.getInstanceStateMap("freshPart-"
+ i);
if (t == null) {
freshSucceeded++;
- Assert.assertNotNull(inIs,
- "Fresh owner " + i + " observed success but its partition is NOT
in IdealState. "
- + "Owner outcome and IdealState are inconsistent.");
- Assert.assertEquals(inIs.get("instance1"), "ONLINE");
+ assertNotNull(inIs, "Fresh owner " + i + " observed success but its
partition is NOT in IdealState. Owner "
+ + "outcome and IdealState are inconsistent.");
+ assertEquals(inIs.get("instance1"), "ONLINE");
} else {
freshThrew++;
- Assert.assertNull(inIs,
- "Fresh owner " + i + " observed exception [" + t.getMessage()
- + "] but its partition IS in IdealState. This is the
orphan-creating race: "
- + "the cancellation flag must skip a failed leader's entry
so no subsequent "
- + "batch applies it.");
+ assertNull(inIs, "Fresh owner " + i + " observed exception [" +
t.getMessage() + "] but its partition IS in "
+ + "IdealState. This is the orphan-creating race: the
cancellation flag must skip a failed leader's entry "
+ + "so no subsequent batch applies it.");
}
}
- LOGGER.info("=== Consistency verified: {} stuck threw, {} fresh
succeeded, {} fresh threw ===",
- stuckCount, freshSucceeded, freshThrew);
+ LOGGER.info("=== Consistency verified: {} stuck threw, {} fresh
succeeded, {} fresh threw ===", stuckCount,
+ freshSucceeded, freshThrew);
} finally {
if (!pool.isShutdown()) {
pool.shutdownNow();
}
-
TEST_INSTANCE.getHelixAdmin().dropResource(TEST_INSTANCE.getHelixClusterName(),
tableName);
+ _helixAdmin.dropResource(_clusterName, tableName);
}
}
+ /// A follower can acquire the queue after a failed leader removes its entry
but before it signals the owner.
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testProcessedFailurePropagatesWhenQueueIsEmpty()
+ throws Exception {
+ String tableName = TABLE_NAME_PREFIX + "empty_OFFLINE";
+ IdealStateGroupCommit commit = new IdealStateGroupCommit();
+ Field queuesField =
IdealStateGroupCommit.class.getDeclaredField("_queues");
+ queuesField.setAccessible(true);
+ Object[] queues = (Object[]) queuesField.get(commit);
+ Object queue = queues[(tableName.hashCode() & Integer.MAX_VALUE) %
queues.length];
+
+ Field runningField = queue.getClass().getDeclaredField("_running");
+ runningField.setAccessible(true);
+ AtomicReference<Thread> running = (AtomicReference<Thread>)
runningField.get(queue);
+ Field pendingField = queue.getClass().getDeclaredField("_pending");
+ pendingField.setAccessible(true);
+ ConcurrentLinkedQueue<Object> pending = (ConcurrentLinkedQueue<Object>)
pendingField.get(queue);
+
+ AtomicReference<Throwable> observedFailure = new AtomicReference<>();
+ RuntimeException batchFailure = new RuntimeException("simulated failed
batch");
+ running.set(Thread.currentThread());
+ Thread owner = new Thread(() -> {
+ try {
+ commit.commit(_helixManager, tableName, is -> is,
RetryPolicies.noDelayRetryPolicy(1), false);
+ } catch (Throwable t) {
+ observedFailure.set(t);
+ }
+ });
+ owner.start();
+ try {
+ long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
+ while (pending.peek() == null && System.nanoTime() < deadline) {
+ Thread.sleep(1);
+ }
+ Object entry = pending.peek();
+ assertNotNull(entry, "Owner did not enqueue its entry");
+ Field exceptionField = entry.getClass().getDeclaredField("_exception");
+ exceptionField.setAccessible(true);
+ exceptionField.set(entry, batchFailure);
+ assertTrue(pending.remove(entry));
+ } finally {
+ running.set(null);
+ owner.join(TimeUnit.SECONDS.toMillis(10));
+ }
+ assertFalse(owner.isAlive(), "Owner did not finish");
+ assertNotNull(observedFailure.get(), "Processed entry's failure was lost");
+ assertSame(observedFailure.get().getCause(), batchFailure);
+ }
+
/// Integration-style test that walks the same chain the orphan-creating
production bug walked:
///
/// - Step 2 of commitSegmentMetadataInternal: write the new consuming
segment's ZK metadata
@@ -363,7 +411,7 @@ public class IdealStateGroupCommitTest {
initialState.setReplicas("1");
initialState.setNumPartitions(0);
initialState.setPartitionState(oldSegment, instance, "CONSUMING");
-
TEST_INSTANCE.getHelixAdmin().addResource(TEST_INSTANCE.getHelixClusterName(),
tableName, initialState);
+ _helixAdmin.addResource(_clusterName, tableName, initialState);
ZkHelixPropertyStore<ZNRecord> propertyStore =
TEST_INSTANCE.getPropertyStore();
@@ -374,7 +422,7 @@ public class IdealStateGroupCommitTest {
SegmentZKMetadata newSegmentMetadata = new SegmentZKMetadata(newSegment);
newSegmentMetadata.setStatus(Status.IN_PROGRESS);
boolean wrote = ZKMetadataProvider.setSegmentZKMetadata(propertyStore,
tableName, newSegmentMetadata, -1);
- Assert.assertTrue(wrote, "Pre-condition: writing new segment ZK metadata
should succeed");
+ assertTrue(wrote, "Pre-condition: writing new segment ZK metadata should
succeed");
// Co-batched in-flight throwing entry: mimics a pauseless segment that
has timed out
// and whose queued updater will throw PermanentUpdaterException when
iterated.
@@ -389,15 +437,13 @@ public class IdealStateGroupCommitTest {
Throwable callerThrew = null;
try {
- commit.commit(TEST_INSTANCE.getHelixManager(), tableName,
stepThreeUpdater,
- RetryPolicies.noDelayRetryPolicy(1), false);
+ commit.commit(_helixManager, tableName, stepThreeUpdater,
RetryPolicies.noDelayRetryPolicy(1), false);
} catch (Throwable t) {
callerThrew = t;
// Mimic commitSegmentMetadataInternal's catch: delete the new
segment's ZK metadata
// best-effort because Step 3 appeared to fail.
ZKMetadataProvider.removeSegmentZKMetadata(propertyStore, tableName,
newSegment);
- LOGGER.info("Caller observed exception, ran
removeSegmentZKMetadataBestEffort(newSegment): {}",
- t.getMessage());
+ LOGGER.info("Caller observed exception, ran
removeSegmentZKMetadataBestEffort(newSegment): {}", t.getMessage());
}
// Drainer commit: mimics any subsequent IdealStateGroupCommit batch
that processes
@@ -409,14 +455,13 @@ public class IdealStateGroupCommitTest {
is.setPartitionState("drainerPartition", instance, "ONLINE");
return is;
};
- commit.commit(TEST_INSTANCE.getHelixManager(), tableName, drainer,
- RetryPolicies.noDelayRetryPolicy(1), false);
+ commit.commit(_helixManager, tableName, drainer,
RetryPolicies.noDelayRetryPolicy(1), false);
// Inspect ground truth.
- IdealState finalState =
HelixHelper.getTableIdealState(TEST_INSTANCE.getHelixManager(), tableName);
+ IdealState finalState = HelixHelper.getTableIdealState(_helixManager,
tableName);
Map<String, String> newSegInIs =
finalState.getInstanceStateMap(newSegment);
- SegmentZKMetadata newSegMetadataAfter =
ZKMetadataProvider.getSegmentZKMetadata(
- propertyStore, tableName, newSegment);
+ SegmentZKMetadata newSegMetadataAfter =
+ ZKMetadataProvider.getSegmentZKMetadata(propertyStore, tableName,
newSegment);
boolean inIdealState = newSegInIs != null;
boolean hasZkMetadata = newSegMetadataAfter != null;
@@ -426,25 +471,19 @@ public class IdealStateGroupCommitTest {
LOGGER.info("newSegment has ZK metadata: {}", hasZkMetadata);
// The orphan condition is precisely "in IdealState && no ZK metadata".
- // Pre-fix: this assertion FAILS -- caller threw, ran cleanup, but the
queued updater was
- // applied by the drainer.
- // Post-fix (cancellation): the entry is cancelled in commit()'s catch;
the drainer skips
- // it. The newSegment is neither in IdealState nor has ZK metadata -- a
clean failure.
- Assert.assertFalse(inIdealState && !hasZkMetadata,
- "ORPHAN DETECTED: newSegment is in IdealState but has no ZK
metadata. "
- + "callerThrew=" + (callerThrew != null ?
callerThrew.getMessage() : "null")
- + ". The cancellation fix in IdealStateGroupCommit must skip the
caller's "
- + "still-queued entry so no subsequent batch applies it.");
+ // Pre-fix: this assertion FAILS -- caller threw, ran cleanup, but the
queued updater was applied by the drainer.
+ // Post-fix (cancellation): the entry is cancelled in commit()'s catch;
the drainer skips it. The newSegment is
+ // neither in IdealState nor has ZK metadata -- a clean failure.
+ assertFalse(inIdealState && !hasZkMetadata, "ORPHAN DETECTED: newSegment
is in IdealState but has no ZK metadata."
+ + " callerThrew=" + (callerThrew != null ? callerThrew.getMessage()
: "null") + ". The cancellation fix in "
+ + "IdealStateGroupCommit must skip the caller's still-queued entry
so no subsequent batch applies it.");
// Post-fix expected state: caller threw AND newSegment is neither in
IdealState nor has
// ZK metadata. The cancellation prevented the queued entry from being
applied.
- Assert.assertNotNull(callerThrew,
- "Post-fix expectation: caller's commit() throws (all-or-nothing
batch semantics).");
- Assert.assertFalse(inIdealState,
- "newSegment must NOT be in IdealState -- the caller's cancelled
entry must be skipped "
- + "by the drainer's iteration.");
- Assert.assertFalse(hasZkMetadata,
- "newSegment's ZK metadata must have been cleaned up by the caller's
catch.");
+ assertNotNull(callerThrew, "Post-fix expectation: caller's commit()
throws (all-or-nothing batch semantics).");
+ assertFalse(inIdealState, "newSegment must NOT be in IdealState -- the
caller's cancelled entry must be skipped "
+ + "by the drainer's iteration.");
+ assertFalse(hasZkMetadata, "newSegment's ZK metadata must have been
cleaned up by the caller's catch.");
} finally {
// Best-effort cleanup of any stragglers. ControllerTest.cleanup()
asserts /SEGMENTS has
// zero child table directories, so we must remove the table-level node
too.
@@ -456,7 +495,7 @@ public class IdealStateGroupCommitTest {
propertyStore.remove("/SEGMENTS/" + tableName,
org.apache.helix.AccessOption.PERSISTENT);
} catch (Throwable ignored) {
}
-
TEST_INSTANCE.getHelixAdmin().dropResource(TEST_INSTANCE.getHelixClusterName(),
tableName);
+ _helixAdmin.dropResource(_clusterName, tableName);
}
}
@@ -466,10 +505,8 @@ public class IdealStateGroupCommitTest {
@SuppressWarnings("unchecked")
private static void injectStuckEntry(IdealStateGroupCommit commit, String
resourceName)
throws Exception {
- Class<?> queueClass =
-
Class.forName("org.apache.pinot.common.utils.helix.IdealStateGroupCommit$Queue");
- Class<?> entryClass =
-
Class.forName("org.apache.pinot.common.utils.helix.IdealStateGroupCommit$Entry");
+ Class<?> queueClass =
Class.forName("org.apache.pinot.common.utils.helix.IdealStateGroupCommit$Queue");
+ Class<?> entryClass =
Class.forName("org.apache.pinot.common.utils.helix.IdealStateGroupCommit$Entry");
Field queuesField =
IdealStateGroupCommit.class.getDeclaredField("_queues");
queuesField.setAccessible(true);
@@ -483,8 +520,7 @@ public class IdealStateGroupCommitTest {
ConcurrentLinkedQueue<Object> pending = (ConcurrentLinkedQueue<Object>)
pendingField.get(queue);
Function<IdealState, IdealState> stuckUpdater = is -> {
- throw new HelixHelper.PermanentUpdaterException(
- "simulated exceeded max segment completion time");
+ throw new HelixHelper.PermanentUpdaterException("simulated exceeded max
segment completion time");
};
Constructor<?> entryCtor = entryClass.getDeclaredConstructor(String.class,
Function.class);
@@ -492,8 +528,7 @@ public class IdealStateGroupCommitTest {
Object stuckEntry = entryCtor.newInstance(resourceName, stuckUpdater);
pending.add(stuckEntry);
- LOGGER.info("Injected stuck (always-throwing) entry into queue bucket {}
for resource {}",
- bucket, resourceName);
+ LOGGER.info("Injected stuck (always-throwing) entry into queue bucket {}
for resource {}", bucket, resourceName);
}
@Test(invocationCount = 5)
@@ -507,60 +542,55 @@ public class IdealStateGroupCommitTest {
for (int j = 0; j < NUM_TABLES; j++) {
String tableName = TABLE_NAME_PREFIX + j + "_OFFLINE";
IdealStateGroupCommit commit = groupCommitList.get(new
Random().nextInt(NUM_PROCESSORS));
- Runnable runnable = new
IdealStateUpdater(TEST_INSTANCE.getHelixManager(), commit, tableName, i);
+ Runnable runnable = new IdealStateUpdater(commit, tableName, i);
_executorService.submit(runnable);
}
}
for (int i = 0; i < NUM_TABLES; i++) {
String tableName = TABLE_NAME_PREFIX + i + "_OFFLINE";
- IdealState idealState =
HelixHelper.getTableIdealState(TEST_INSTANCE.getHelixManager(), tableName);
+ IdealState idealState = HelixHelper.getTableIdealState(_helixManager,
tableName);
while (idealState.getNumPartitions() < NUM_UPDATES) {
Thread.sleep(500);
- idealState =
HelixHelper.getTableIdealState(TEST_INSTANCE.getHelixManager(), tableName);
+ idealState = HelixHelper.getTableIdealState(_helixManager, tableName);
}
- Assert.assertEquals(idealState.getNumPartitions(), NUM_UPDATES);
+ assertEquals(idealState.getNumPartitions(), NUM_UPDATES);
ControllerMetrics controllerMetrics = ControllerMetrics.get();
long idealStateUpdateSuccessCount =
controllerMetrics.getMeteredTableValue(tableName,
ControllerMeter.IDEAL_STATE_UPDATE_SUCCESS).count();
- Assert.assertTrue(idealStateUpdateSuccessCount <= NUM_UPDATES);
+ assertTrue(idealStateUpdateSuccessCount <= NUM_UPDATES);
LOGGER.info("{} IdealState update are successfully committed with {}
times zk updates.", NUM_UPDATES,
idealStateUpdateSuccessCount);
}
}
-}
-class IdealStateUpdater implements Runnable {
- private static final Logger LOGGER =
LoggerFactory.getLogger(IdealStateGroupCommitTest.class);
+ private class IdealStateUpdater implements Runnable {
+ private static final Logger LOGGER =
LoggerFactory.getLogger(IdealStateGroupCommitTest.class);
- private final HelixManager _helixManager;
- private final IdealStateGroupCommit _commit;
- private final String _tableName;
- private final int _i;
+ private final IdealStateGroupCommit _commit;
+ private final String _tableName;
+ private final int _i;
- public IdealStateUpdater(HelixManager helixManager, IdealStateGroupCommit
commit, String tableName, int i) {
- _helixManager = helixManager;
- _commit = commit;
- _tableName = tableName;
- _i = i;
- }
+ public IdealStateUpdater(IdealStateGroupCommit commit, String tableName,
int i) {
+ _commit = commit;
+ _tableName = tableName;
+ _i = i;
+ }
- @Override
- public void run() {
- Function<IdealState, IdealState> updater = new Function<IdealState,
IdealState>() {
- @Override
- public IdealState apply(IdealState idealState) {
+ @Override
+ public void run() {
+ Function<IdealState, IdealState> updater = idealState -> {
idealState.setPartitionState("test_id" + _i, "test_id" + _i, "ONLINE");
return idealState;
- }
- };
+ };
- while (true) {
- try {
- if (_commit.commit(_helixManager, _tableName, updater,
RetryPolicies.noDelayRetryPolicy(1), false) != null) {
- break;
+ while (true) {
+ try {
+ if (_commit.commit(_helixManager, _tableName, updater,
RetryPolicies.noDelayRetryPolicy(1), false) != null) {
+ break;
+ }
+ } catch (Throwable e) {
+ LOGGER.warn("IdealState updater {} failed to commit.", _i, e);
}
- } catch (Throwable e) {
- LOGGER.warn("IdealState updater {} failed to commit.", _i, e);
}
}
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]