diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/consumer/runtime/CommitProgressSyncProcedure.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/consumer/runtime/CommitProgressSyncProcedure.java index efb89b4473d1..39ace4419b56 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/consumer/runtime/CommitProgressSyncProcedure.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/consumer/runtime/CommitProgressSyncProcedure.java @@ -147,10 +147,13 @@ private static void syncCommitProgress( for (final Map.Entry entry : respMap.entrySet()) { final TPullCommitProgressResp resp = entry.getValue(); if (!isSuccessfulResponse(resp)) { - LOGGER.warn( - ProcedureMessages.LOG_FAILED_PULL_COMMIT_PROGRESS_DATANODE_ARG_STATUS_ARG_33037B29, - entry.getKey(), - Objects.isNull(resp) ? null : resp.getStatus()); + // DataNodes with subscription disabled are expected to reject best-effort pulls. + if (!isUnsupportedOperationResponse(resp)) { + LOGGER.warn( + ProcedureMessages.LOG_FAILED_PULL_COMMIT_PROGRESS_DATANODE_ARG_STATUS_ARG_33037B29, + entry.getKey(), + Objects.isNull(resp) ? null : resp.getStatus()); + } continue; } if (resp.isSetCommitRegionProgress()) { @@ -192,6 +195,12 @@ private static boolean isSuccessfulResponse(final TPullCommitProgressResp respon && response.getStatus().getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode(); } + private static boolean isUnsupportedOperationResponse(final TPullCommitProgressResp response) { + return Objects.nonNull(response) + && response.isSetStatus() + && response.getStatus().getCode() == TSStatusCode.UNSUPPORTED_OPERATION.getStatusCode(); + } + @Override public void executeFromOperateOnDataNodes(ConfigNodeProcedureEnv env) { LOGGER.info( diff --git a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/subscription/consumer/runtime/CommitProgressSyncProcedureTest.java b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/subscription/consumer/runtime/CommitProgressSyncProcedureTest.java index 3797563c8506..2314a952086c 100644 --- a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/subscription/consumer/runtime/CommitProgressSyncProcedureTest.java +++ b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/subscription/consumer/runtime/CommitProgressSyncProcedureTest.java @@ -33,6 +33,7 @@ import org.apache.iotdb.rpc.subscription.payload.poll.WriterProgress; import org.junit.Test; +import org.mockito.ArgumentCaptor; import org.mockito.Mockito; import java.io.ByteArrayOutputStream; @@ -51,12 +52,15 @@ public class CommitProgressSyncProcedureTest { @Test public void requiredSyncShouldRejectFailedResponseBeforeConsensusWrite() throws Exception { + assertRequiredSyncRejectsResponse(TSStatusCode.EXECUTE_STATEMENT_ERROR); + assertRequiredSyncRejectsResponse(TSStatusCode.UNSUPPORTED_OPERATION); + } + + private static void assertRequiredSyncRejectsResponse(final TSStatusCode statusCode) + throws Exception { final ConfigNodeProcedureEnv env = Mockito.mock(ConfigNodeProcedureEnv.class); final Map responses = new LinkedHashMap<>(); - responses.put( - 2, - new TPullCommitProgressResp( - new TSStatus(TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode()))); + responses.put(2, new TPullCommitProgressResp(new TSStatus(statusCode.getStatusCode()))); Mockito.when(env.pullCommitProgressFromDataNodes()).thenReturn(responses); try { @@ -70,6 +74,54 @@ public void requiredSyncShouldRejectFailedResponseBeforeConsensusWrite() throws Mockito.verify(env, Mockito.never()).getConfigManager(); } + @Test + public void bestEffortSyncShouldSkipUnsuccessfulResponses() throws Exception { + final String progressKey = "successful_progress"; + final RegionProgress progress = + new RegionProgress( + Collections.singletonMap(new WriterId("DataRegion[1]", 1), new WriterProgress(200, 1))); + final Map responses = new LinkedHashMap<>(); + responses.put( + 1, + new TPullCommitProgressResp( + new TSStatus(TSStatusCode.UNSUPPORTED_OPERATION.getStatusCode()))); + responses.put( + 2, + new TPullCommitProgressResp(new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode())) + .setCommitRegionProgress(Collections.singletonMap(progressKey, serialize(progress)))); + responses.put( + 3, + new TPullCommitProgressResp( + new TSStatus(TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode()))); + responses.put(4, null); + responses.put(5, new TPullCommitProgressResp()); + + final ConfigNodeProcedureEnv env = Mockito.mock(ConfigNodeProcedureEnv.class); + final ConfigManager configManager = Mockito.mock(ConfigManager.class); + final ConsensusManager consensusManager = Mockito.mock(ConsensusManager.class); + Mockito.when(env.pullCommitProgressFromDataNodesBestEffort()).thenReturn(responses); + Mockito.when(env.getConfigManager()).thenReturn(configManager); + Mockito.when(configManager.getConsensusManager()).thenReturn(consensusManager); + Mockito.when(consensusManager.write(Mockito.any())) + .thenReturn(new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode())); + + final CommitProgressSyncProcedure procedure = + new CommitProgressSyncProcedure() { + { + subscriptionInfo = new AtomicReference<>(new SubscriptionInfo()); + } + }; + procedure.executeFromOperateOnConfigNodes(env); + + final ArgumentCaptor planCaptor = + ArgumentCaptor.forClass(CommitProgressHandleMetaChangePlan.class); + Mockito.verify(consensusManager).write(planCaptor.capture()); + assertEquals( + progress, + RegionProgress.deserialize(planCaptor.getValue().getRegionProgressMap().get(progressKey))); + Mockito.verify(env, Mockito.never()).pullCommitProgressFromDataNodes(); + } + @Test public void requiredSyncShouldPersistEmptyProgressAndMergeByMaximum() throws Exception { final String emptyProgressKey = "empty_progress";