From 720f558421819ee5c50419946f7c465b1ceffe40 Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Thu, 10 Sep 2026 17:20:50 +0800 Subject: [PATCH 1/2] [Subscription] Suppress expected commit progress warnings --- .../runtime/CommitProgressSyncProcedure.java | 17 +++- .../CommitProgressSyncProcedureTest.java | 83 ++++++++++++++++++- 2 files changed, 92 insertions(+), 8 deletions(-) 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..d148ea68cdd4 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 @@ -32,8 +32,14 @@ import org.apache.iotdb.rpc.subscription.payload.poll.WriterId; import org.apache.iotdb.rpc.subscription.payload.poll.WriterProgress; +import ch.qos.logback.classic.Level; +import ch.qos.logback.classic.Logger; +import ch.qos.logback.classic.spi.ILoggingEvent; +import ch.qos.logback.core.read.ListAppender; import org.junit.Test; +import org.mockito.ArgumentCaptor; import org.mockito.Mockito; +import org.slf4j.LoggerFactory; import java.io.ByteArrayOutputStream; import java.io.DataOutputStream; @@ -51,12 +57,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 +79,72 @@ public void requiredSyncShouldRejectFailedResponseBeforeConsensusWrite() throws Mockito.verify(env, Mockito.never()).getConfigManager(); } + @Test + public void bestEffortSyncShouldOnlyWarnForUnexpectedResponses() 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()); + } + }; + final Logger logger = (Logger) LoggerFactory.getLogger(CommitProgressSyncProcedure.class); + final Level originalLevel = logger.getLevel(); + final ListAppender appender = new ListAppender<>(); + appender.start(); + logger.addAppender(appender); + logger.setLevel(Level.WARN); + try { + procedure.executeFromOperateOnConfigNodes(env); + + assertEquals(3, appender.list.size()); + for (int i = 0; i < appender.list.size(); i++) { + assertEquals(Level.WARN, appender.list.get(i).getLevel()); + assertEquals(i + 3, appender.list.get(i).getArgumentArray()[0]); + } + } finally { + logger.detachAppender(appender); + appender.stop(); + logger.setLevel(originalLevel); + } + + 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"; From 56b08ec30d1f0609f8302535e8d3252c89c48682 Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Thu, 10 Sep 2026 18:36:51 +0800 Subject: [PATCH 2/2] [Subscription] Remove Logback coupling from commit progress test --- .../CommitProgressSyncProcedureTest.java | 27 ++----------------- 1 file changed, 2 insertions(+), 25 deletions(-) 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 d148ea68cdd4..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 @@ -32,14 +32,9 @@ import org.apache.iotdb.rpc.subscription.payload.poll.WriterId; import org.apache.iotdb.rpc.subscription.payload.poll.WriterProgress; -import ch.qos.logback.classic.Level; -import ch.qos.logback.classic.Logger; -import ch.qos.logback.classic.spi.ILoggingEvent; -import ch.qos.logback.core.read.ListAppender; import org.junit.Test; import org.mockito.ArgumentCaptor; import org.mockito.Mockito; -import org.slf4j.LoggerFactory; import java.io.ByteArrayOutputStream; import java.io.DataOutputStream; @@ -80,7 +75,7 @@ private static void assertRequiredSyncRejectsResponse(final TSStatusCode statusC } @Test - public void bestEffortSyncShouldOnlyWarnForUnexpectedResponses() throws Exception { + public void bestEffortSyncShouldSkipUnsuccessfulResponses() throws Exception { final String progressKey = "successful_progress"; final RegionProgress progress = new RegionProgress( @@ -116,25 +111,7 @@ public void bestEffortSyncShouldOnlyWarnForUnexpectedResponses() throws Exceptio subscriptionInfo = new AtomicReference<>(new SubscriptionInfo()); } }; - final Logger logger = (Logger) LoggerFactory.getLogger(CommitProgressSyncProcedure.class); - final Level originalLevel = logger.getLevel(); - final ListAppender appender = new ListAppender<>(); - appender.start(); - logger.addAppender(appender); - logger.setLevel(Level.WARN); - try { - procedure.executeFromOperateOnConfigNodes(env); - - assertEquals(3, appender.list.size()); - for (int i = 0; i < appender.list.size(); i++) { - assertEquals(Level.WARN, appender.list.get(i).getLevel()); - assertEquals(i + 3, appender.list.get(i).getArgumentArray()[0]); - } - } finally { - logger.detachAppender(appender); - appender.stop(); - logger.setLevel(originalLevel); - } + procedure.executeFromOperateOnConfigNodes(env); final ArgumentCaptor planCaptor = ArgumentCaptor.forClass(CommitProgressHandleMetaChangePlan.class);