This is an automated email from the ASF dual-hosted git repository.
jt2594838 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new 320f0d5a839 [Subscription] Suppress expected commit progress warnings
(#18615)
320f0d5a839 is described below
commit 320f0d5a839879dd0c07b8042c049bf5b34168e8
Author: Caideyipi <[email protected]>
AuthorDate: Fri Sep 11 09:28:34 2026 +0800
[Subscription] Suppress expected commit progress warnings (#18615)
* [Subscription] Suppress expected commit progress warnings
* [Subscription] Remove Logback coupling from commit progress test
---
.../runtime/CommitProgressSyncProcedure.java | 17 ++++--
.../runtime/CommitProgressSyncProcedureTest.java | 60 ++++++++++++++++++++--
2 files changed, 69 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 efb89b4473d..39ace4419b5 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 @@ public class CommitProgressSyncProcedure extends
AbstractOperateSubscriptionProc
for (final Map.Entry<Integer, TPullCommitProgressResp> 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 @@ public class CommitProgressSyncProcedure extends
AbstractOperateSubscriptionProc
&& 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 3797563c850..2314a952086 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.WriterId;
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<Integer, TPullCommitProgressResp> 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 class CommitProgressSyncProcedureTest {
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<Integer, TPullCommitProgressResp> 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<CommitProgressHandleMetaChangePlan> 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";