This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-12227-38a37104e59ee3317cf0bc299c9bfee22149391c in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit f5b7f33e6d19e9c3ebee14828f7eb40ae641ecca Author: Daniel <[email protected]> AuthorDate: Tue Sep 15 05:41:33 2026 +0000 [Fix][Zeta] Bound retries for metrics snapshot updates (#12227) Co-authored-by: Daniel <[email protected]> --- docs/en/engines/zeta/telemetry.md | 8 + docs/zh/engines/zeta/telemetry.md | 6 + .../HazelcastMetricsSnapshotStateStore.java | 56 +++- .../HazelcastMetricsSnapshotStateStoreTest.java | 302 ++++++++++++++++++++- 4 files changed, 362 insertions(+), 10 deletions(-) diff --git a/docs/en/engines/zeta/telemetry.md b/docs/en/engines/zeta/telemetry.md index e8f8af28cf..71ac33a907 100644 --- a/docs/en/engines/zeta/telemetry.md +++ b/docs/en/engines/zeta/telemetry.md @@ -150,6 +150,14 @@ These metrics are exported by the active master only; scraping a worker node's e ### Report Metrics Operation +Metrics snapshot writes and deletions attempt at most 10 conditional updates per bucket. If contention +persists, the operation fails with `Failed to update metrics partition ... after 10 concurrent +modifications`. A failed worker report is logged and counted in +`report_metrics_operation_total{result="failure"}`; subsequent scheduled reports can retry while +the task context is retained. Pending pipeline cleanup retains its record when metrics deletion +fails. This limit bounds conflict retries, not network latency: individual Hazelcast invocations +still use their configured timeouts. The metrics format and checkpoint/savepoint state are unchanged. + | MetricName | Type | Labels | DESCRIPTION | |---------------------------------------------------|---------|-----------------------------------------------------------------------------------------------------|------------------------------------------------------------------------------------------------------------------| | report_metrics_operation_total | Counter | **address**, worker instance address,for example: "127.0.0.1:5801". **result**, one of "success" "failure" "interrupted" | The total number of `ReportMetricsOperation` invocations sent by a worker | diff --git a/docs/zh/engines/zeta/telemetry.md b/docs/zh/engines/zeta/telemetry.md index af9f83a02f..bd9841c424 100644 --- a/docs/zh/engines/zeta/telemetry.md +++ b/docs/zh/engines/zeta/telemetry.md @@ -149,6 +149,12 @@ engine_state_store_connector_jar_total_references{backend="hazelcast"} ### ReportMetricsOperation 指标 +指标快照写入和删除对同一个分桶最多尝试 10 次竞争更新。持续发生竞争时,操作会抛出 +`Failed to update metrics partition ... after 10 concurrent modifications`。Worker 上报失败会记录日志, +并计入 `report_metrics_operation_total{result="failure"}`;任务上下文仍保留时,后续定时上报可以再次尝试。 +待处理 Pipeline 清理在指标删除失败时会保留清理记录。这个上限限制的是竞争重试次数,单次 Hazelcast +调用的等待时间仍由其超时设置决定。指标格式和 checkpoint/savepoint 状态保持不变。 + | MetricName | Type | Labels | 描述 | |---------------------------------------------------|---------|--------------------------------------------------------------------------------------------|--------------------------------------------------------------------------------------| | report_metrics_operation_total | Counter | **address**,worker 实例地址,例如:"127.0.0.1:5801"。**result**,取值包括:"success" "failure" "interrupted" | worker 发送的 `ReportMetricsOperation` 调用总次数 | diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/common/statestore/metrics/hazelcast/HazelcastMetricsSnapshotStateStore.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/common/statestore/metrics/hazelcast/HazelcastMetricsSnapshotStateStore.java index 5c00c30de8..8955a01fdb 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/common/statestore/metrics/hazelcast/HazelcastMetricsSnapshotStateStore.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/common/statestore/metrics/hazelcast/HazelcastMetricsSnapshotStateStore.java @@ -18,6 +18,7 @@ package org.apache.seatunnel.engine.server.common.statestore.metrics.hazelcast; import org.apache.seatunnel.common.utils.HashUtils; +import org.apache.seatunnel.engine.common.exception.SeaTunnelEngineException; import org.apache.seatunnel.engine.server.common.statestore.metrics.MetricsSnapshotStateStore; import org.apache.seatunnel.engine.server.dag.physical.PipelineLocation; import org.apache.seatunnel.engine.server.execution.TaskLocation; @@ -36,12 +37,17 @@ import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicLong; +import java.util.function.UnaryOperator; /** Implementation backed by a partitioned Hazelcast metrics {@link IMap}. */ public class HazelcastMetricsSnapshotStateStore implements MetricsSnapshotStateStore, AutoCloseable { + // Metrics contention must not keep task completion or pipeline cleanup retrying indefinitely. + private static final int MAX_UPDATE_ATTEMPTS = 10; + private final IMap<Long, Map<TaskLocation, SeaTunnelMetricsContext>> metricsImap; private final int partitionCount; private final AtomicLong activePartitionKeyCount = new AtomicLong(); @@ -78,12 +84,12 @@ public class HazelcastMetricsSnapshotStateStore .parallelStream() .forEach( entry -> { - metricsImap.compute( + updatePartition( entry.getKey(), - (k, oldVal) -> { - if (oldVal == null) oldVal = new HashMap<>(); - oldVal.putAll(entry.getValue()); - return oldVal; + current -> { + if (current == null) current = new HashMap<>(); + current.putAll(entry.getValue()); + return current; }); }); } @@ -100,9 +106,9 @@ public class HazelcastMetricsSnapshotStateStore @Override public void remove(final TaskLocation taskLocation) { - metricsImap.compute( + updatePartition( partition(taskLocation), - (ignored, current) -> { + current -> { if (current == null) { return null; } @@ -115,9 +121,9 @@ public class HazelcastMetricsSnapshotStateStore @Override public void removePipeline(final PipelineLocation pipelineLocation) { for (long partition = 0; partition < partitionCount; partition++) { - metricsImap.compute( + updatePartition( partition, - (ignored, current) -> { + current -> { if (current == null || current.isEmpty()) { return current; } @@ -167,6 +173,38 @@ public class HazelcastMetricsSnapshotStateStore return Math.toIntExact(activePartitionKeyCount.get()); } + /** + * Bounds Hazelcast 5.1's local compute retry loop so contention cannot indefinitely hold up a + * final worker metrics report or pending pipeline cleanup. Each retry invokes this lambda again + * with the latest bucket; exhaustion propagates to the existing report or cleanup failure path. + * + * <p>Keep the non-serializable lambda local: no new execution class needs to exist on an older + * partition owner. Delegating the compare-and-set to Hazelcast also preserves its original + * serialized expected value instead of reserializing a deserialized metrics map. Individual + * Hazelcast invocations still use their configured timeouts. + */ + private void updatePartition( + long partition, UnaryOperator<Map<TaskLocation, SeaTunnelMetricsContext>> mutation) { + AtomicInteger attempts = new AtomicInteger(); + metricsImap.compute( + partition, + (ignored, current) -> { + if (Thread.currentThread().isInterrupted()) { + throw new SeaTunnelEngineException( + "Interrupted while updating metrics partition " + partition); + } + if (attempts.getAndIncrement() >= MAX_UPDATE_ATTEMPTS) { + throw new SeaTunnelEngineException( + "Failed to update metrics partition " + + partition + + " after " + + MAX_UPDATE_ATTEMPTS + + " concurrent modifications"); + } + return mutation.apply(current); + }); + } + private long partition(TaskLocation taskLocation) { return HashUtils.bucketIndex(taskLocation.hashCode(), partitionCount); } diff --git a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/common/statestore/metrics/hazelcast/HazelcastMetricsSnapshotStateStoreTest.java b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/common/statestore/metrics/hazelcast/HazelcastMetricsSnapshotStateStoreTest.java index 7ef1ada64e..b6a76b2223 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/common/statestore/metrics/hazelcast/HazelcastMetricsSnapshotStateStoreTest.java +++ b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/common/statestore/metrics/hazelcast/HazelcastMetricsSnapshotStateStoreTest.java @@ -17,6 +17,7 @@ package org.apache.seatunnel.engine.server.common.statestore.metrics.hazelcast; +import org.apache.seatunnel.engine.common.exception.SeaTunnelEngineException; import org.apache.seatunnel.engine.server.dag.physical.PipelineLocation; import org.apache.seatunnel.engine.server.execution.TaskGroupLocation; import org.apache.seatunnel.engine.server.execution.TaskLocation; @@ -25,6 +26,8 @@ import org.apache.seatunnel.engine.server.metrics.SeaTunnelMetricsContext; import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.EnumSource; import com.hazelcast.config.Config; import com.hazelcast.core.Hazelcast; @@ -35,11 +38,25 @@ import java.util.HashMap; import java.util.LinkedHashMap; import java.util.Map; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.BiFunction; +import java.util.function.LongConsumer; import static org.awaitility.Awaitility.await; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNull; - +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.spy; + +/** + * Covers metrics snapshot persistence and bounded mutation retries against real Hazelcast buckets. + */ class HazelcastMetricsSnapshotStateStoreTest { private static HazelcastInstance hazelcastInstance; @@ -167,6 +184,289 @@ class HazelcastMetricsSnapshotStateStoreTest { awaitSize(store, 2); } + /** + * A concurrent initial report must survive a put-if-absent conflict without losing either + * writer's task snapshot. + */ + @Test + void mergeShouldPreserveConcurrentBucketCreation() { + IMap<Long, Map<TaskLocation, SeaTunnelMetricsContext>> iMap = + hazelcastInstance.getMap("metrics-snapshot-concurrent-create"); + TaskLocation reported = taskLocation(5L, 50, 500L, 0L, 0); + TaskLocation concurrent = taskLocation(6L, 60, 600L, 0L, 0); + AtomicInteger attempts = new AtomicInteger(); + IMap<Long, Map<TaskLocation, SeaTunnelMetricsContext>> intercepted = + withConcurrentWrites( + iMap, + partition -> { + if (attempts.getAndIncrement() == 0) { + iMap.set( + partition, + singletonSnapshot( + concurrent, metricsContextWithCounterValue(2))); + } + }); + + try (HazelcastMetricsSnapshotStateStore store = + new HazelcastMetricsSnapshotStateStore(intercepted, 1)) { + store.merge(singletonSnapshot(reported, metricsContextWithCounterValue(1))); + + assertCounterValue(1, store.get(reported)); + assertCounterValue(2, store.get(concurrent)); + assertEquals(2, attempts.get()); + awaitSize(store, 2); + } + } + + /** + * A retry must reread the bucket and retain another pipeline's snapshots even when the first + * attempt would have removed the entire bucket. + */ + @ParameterizedTest + @EnumSource(Mutation.class) + void mutationShouldPreserveConcurrentPipelineSnapshots(Mutation mutation) { + IMap<Long, Map<TaskLocation, SeaTunnelMetricsContext>> iMap = + hazelcastInstance.getMap("metrics-snapshot-concurrent-update-" + mutation); + TaskLocation target = taskLocation(7L, 70, 700L, 0L, 0); + TaskLocation concurrent = taskLocation(8L, 80, 800L, 0L, 0); + iMap.put(0L, singletonSnapshot(target, metricsContextWithCounterValue(1))); + AtomicInteger attempts = new AtomicInteger(); + IMap<Long, Map<TaskLocation, SeaTunnelMetricsContext>> intercepted = + withConcurrentWrites( + iMap, + partition -> { + if (attempts.getAndIncrement() == 0) { + Map<TaskLocation, SeaTunnelMetricsContext> snapshot = + iMap.get(partition); + snapshot.put(concurrent, metricsContextWithCounterValue(2)); + iMap.set(partition, snapshot); + } + }); + + try (HazelcastMetricsSnapshotStateStore store = + new HazelcastMetricsSnapshotStateStore(intercepted, 1)) { + applyMutation(store, target, mutation); + + assertCounterValue(2, store.get(concurrent)); + assertEquals(2, attempts.get()); + if (mutation == Mutation.MERGE) { + assertCounterValue(100, store.get(target)); + awaitSize(store, 2); + } else { + assertNull(store.get(target)); + awaitSize(store, 1); + } + } + } + + /** + * The last allowed attempt may still succeed; exhausting earlier conflicts must not reject a + * successful update or delete at the retry boundary. + */ + @ParameterizedTest + @EnumSource(Mutation.class) + void mutationShouldSucceedOnLastAllowedAttempt(Mutation mutation) { + IMap<Long, Map<TaskLocation, SeaTunnelMetricsContext>> iMap = + hazelcastInstance.getMap("metrics-snapshot-last-attempt-" + mutation); + TaskLocation target = taskLocation(12L, 120, 1200L, 0L, 0); + iMap.put(0L, singletonSnapshot(target, metricsContextWithCounterValue(0))); + AtomicInteger attempts = new AtomicInteger(); + IMap<Long, Map<TaskLocation, SeaTunnelMetricsContext>> intercepted = + withConcurrentWrites( + iMap, + partition -> { + int attempt = attempts.incrementAndGet(); + if (attempt < 10) { + iMap.set( + partition, + singletonSnapshot( + target, metricsContextWithCounterValue(attempt))); + } + }); + + try (HazelcastMetricsSnapshotStateStore store = + new HazelcastMetricsSnapshotStateStore(intercepted, 1)) { + applyMutation(store, target, mutation); + + assertEquals(10, attempts.get()); + if (mutation == Mutation.MERGE) { + assertCounterValue(100, store.get(target)); + } else { + assertNull(store.get(target)); + } + } + } + + /** + * Sustained conflicts must surface a failure while preserving the winning writer's snapshot. + * The conflict injector also has a cap so the unfixed implementation fails instead of hanging. + */ + @ParameterizedTest + @EnumSource(Mutation.class) + void mutationShouldFailAfterRepeatedConflicts(Mutation mutation) { + IMap<Long, Map<TaskLocation, SeaTunnelMetricsContext>> iMap = + hazelcastInstance.getMap("metrics-snapshot-conflict-limit-" + mutation); + TaskLocation target = taskLocation(9L, 90, 900L, 0L, 0); + iMap.put(0L, singletonSnapshot(target, metricsContextWithCounterValue(0))); + AtomicInteger conflicts = new AtomicInteger(); + IMap<Long, Map<TaskLocation, SeaTunnelMetricsContext>> intercepted = + withConcurrentWrites( + iMap, + partition -> { + int conflict = conflicts.incrementAndGet(); + assertTrue(conflict <= 20, "Metrics mutation kept retrying conflicts"); + iMap.set( + partition, + singletonSnapshot( + target, metricsContextWithCounterValue(conflict))); + }); + + try (HazelcastMetricsSnapshotStateStore store = + new HazelcastMetricsSnapshotStateStore(intercepted, 1)) { + SeaTunnelEngineException failure = + assertThrows( + SeaTunnelEngineException.class, + () -> applyMutation(store, target, mutation)); + + assertTrue(failure.getMessage().contains("metrics partition 0")); + assertEquals(10, conflicts.get()); + assertCounterValue(conflicts.get(), store.get(target)); + awaitSize(store, 1); + } + } + + /** + * Interruption between the bucket read and mutation must stop the retry and retain the flag. + */ + @ParameterizedTest + @EnumSource(Mutation.class) + void mutationShouldPreserveInterruption(Mutation mutation) { + IMap<Long, Map<TaskLocation, SeaTunnelMetricsContext>> iMap = + hazelcastInstance.getMap("metrics-snapshot-interrupted-" + mutation); + TaskLocation target = taskLocation(10L, 100, 1000L, 0L, 0); + iMap.put(0L, singletonSnapshot(target, metricsContextWithCounterValue(1))); + IMap<Long, Map<TaskLocation, SeaTunnelMetricsContext>> intercepted = spy(iMap); + doAnswer( + invocation -> { + Long partition = invocation.getArgument(0); + BiFunction< + Long, + Map<TaskLocation, SeaTunnelMetricsContext>, + Map<TaskLocation, SeaTunnelMetricsContext>> + remapping = invocation.getArgument(1); + return iMap.compute( + partition, + (key, current) -> { + Thread.currentThread().interrupt(); + return remapping.apply(key, current); + }); + }) + .when(intercepted) + .compute(anyLong(), any()); + + try (HazelcastMetricsSnapshotStateStore store = + new HazelcastMetricsSnapshotStateStore(intercepted, 1)) { + try { + assertThrows( + SeaTunnelEngineException.class, + () -> applyMutation(store, target, mutation)); + assertTrue(Thread.currentThread().isInterrupted()); + } finally { + // Do not leak the deliberately injected interrupt into other tests or shutdown. + Thread.interrupted(); + } + assertCounterValue(1, store.get(target)); + } + } + + /** + * Backend failures must reach the existing caller failure paths while leaving the previously + * stored snapshot available for recovery. + */ + @ParameterizedTest + @EnumSource(Mutation.class) + void mutationShouldPropagateBackendFailure(Mutation mutation) { + IMap<Long, Map<TaskLocation, SeaTunnelMetricsContext>> iMap = + hazelcastInstance.getMap("metrics-snapshot-backend-failure-" + mutation); + TaskLocation target = taskLocation(11L, 110, 1100L, 0L, 0); + iMap.put(0L, singletonSnapshot(target, metricsContextWithCounterValue(1))); + IMap<Long, Map<TaskLocation, SeaTunnelMetricsContext>> intercepted = spy(iMap); + IllegalStateException failure = new IllegalStateException("test backend failure"); + doThrow(failure).when(intercepted).compute(anyLong(), any()); + + try (HazelcastMetricsSnapshotStateStore store = + new HazelcastMetricsSnapshotStateStore(intercepted, 1)) { + assertSame( + failure, + assertThrows( + IllegalStateException.class, + () -> applyMutation(store, target, mutation))); + assertCounterValue(1, store.get(target)); + } + } + + /** + * Routes each injected conflict through the public merge, task-removal or pipeline-removal + * entry point, using the same target snapshot. + */ + private static void applyMutation( + HazelcastMetricsSnapshotStateStore store, TaskLocation target, Mutation mutation) { + switch (mutation) { + case MERGE: + store.merge(singletonSnapshot(target, metricsContextWithCounterValue(100))); + break; + case REMOVE_TASK: + store.remove(target); + break; + case REMOVE_PIPELINE: + store.removePipeline(target.getTaskGroupLocation().getPipelineLocation()); + break; + default: + throw new AssertionError("Unexpected mutation: " + mutation); + } + } + + /** + * Injects an actual map write between remapping and Hazelcast's conditional write. The real + * compute loop, binary comparisons and serialization remain responsible for detecting + * conflicts. + */ + private static IMap<Long, Map<TaskLocation, SeaTunnelMetricsContext>> withConcurrentWrites( + IMap<Long, Map<TaskLocation, SeaTunnelMetricsContext>> iMap, + LongConsumer concurrentWrite) { + IMap<Long, Map<TaskLocation, SeaTunnelMetricsContext>> intercepted = spy(iMap); + doAnswer( + invocation -> { + Long partition = invocation.getArgument(0); + BiFunction< + Long, + Map<TaskLocation, SeaTunnelMetricsContext>, + Map<TaskLocation, SeaTunnelMetricsContext>> + remapping = invocation.getArgument(1); + return iMap.compute( + partition, + (key, current) -> { + Map<TaskLocation, SeaTunnelMetricsContext> updated = + remapping.apply(key, current); + concurrentWrite.accept(key); + return updated; + }); + }) + .when(intercepted) + .compute(anyLong(), any()); + return intercepted; + } + + /** + * Mutation entry points that share the retry limit and must preserve the same interruption and + * backend failure behavior. + */ + private enum Mutation { + MERGE, + REMOVE_TASK, + REMOVE_PIPELINE + } + private static void awaitSize(HazelcastMetricsSnapshotStateStore store, int expectedSize) { await().atMost(10, TimeUnit.SECONDS) .untilAsserted(() -> assertEquals(expectedSize, store.size()));
