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 6fa33f21ba0 Fix duplicate subscription auto commits (#18710)
6fa33f21ba0 is described below

commit 6fa33f21ba0f5b2af8dfdda398e7674fdbfd1689
Author: Caideyipi <[email protected]>
AuthorDate: Thu Sep 24 14:55:02 2026 +0800

    Fix duplicate subscription auto commits (#18710)
---
 .../base/AbstractSubscriptionConsumer.java         |   7 ++
 .../base/AbstractSubscriptionProvider.java         |   2 +-
 .../base/AbstractSubscriptionPullConsumer.java     |  19 ++++
 .../base/SubscriptionConsumerLifecycleTest.java    | 119 ++++++++++++++++++++-
 .../broker/SubscriptionPrefetchingQueue.java       |   4 +-
 .../receiver/SubscriptionReceiverV1.java           |  19 ++++
 6 files changed, 164 insertions(+), 6 deletions(-)

diff --git 
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumer.java
 
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumer.java
index a4f5d36a324..e2ab3476606 100644
--- 
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumer.java
+++ 
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumer.java
@@ -1526,6 +1526,7 @@ abstract class AbstractSubscriptionConsumer implements 
AutoCloseable {
       acceptedCount += acceptedCommitContexts.size();
       if (!nack) {
         overlayCommittedPositions(commitResp.getCommittedProgressByTopic());
+        onAckedCommitContexts(acceptedCommitContexts);
       }
       if (acceptedCommitContexts.size() != groupedCommitContexts.size()) {
         final List<SubscriptionCommitContext> failedInGroup =
@@ -1547,6 +1548,12 @@ abstract class AbstractSubscriptionConsumer implements 
AutoCloseable {
     }
   }
 
+  /** Invoked after the server accepts commit contexts from an acknowledgement 
request. */
+  protected void onAckedCommitContexts(
+      final Collection<SubscriptionCommitContext> acceptedCommitContexts) {
+    // Do nothing by default.
+  }
+
   protected Set<SubscriptionMessage> ackWithPartialProgress(
       final Iterable<SubscriptionMessage> messages) throws 
SubscriptionException {
     final List<SubscriptionMessage> bufferedMessages = new ArrayList<>();
diff --git 
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionProvider.java
 
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionProvider.java
index 02a6de8a9bb..477a7ae45cb 100644
--- 
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionProvider.java
+++ 
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionProvider.java
@@ -674,7 +674,7 @@ public abstract class AbstractSubscriptionProvider {
 
     private final Map<String, TopicProgress> committedProgressByTopic;
 
-    private CommitResult(
+    CommitResult(
         final List<SubscriptionCommitContext> acceptedCommitContexts,
         final Map<String, TopicProgress> committedProgressByTopic) {
       this.acceptedCommitContexts =
diff --git 
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionPullConsumer.java
 
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionPullConsumer.java
index b565d0d9293..8ad85e2ba9e 100644
--- 
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionPullConsumer.java
+++ 
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionPullConsumer.java
@@ -37,6 +37,7 @@ import org.slf4j.LoggerFactory;
 
 import java.time.Duration;
 import java.util.ArrayList;
+import java.util.Collection;
 import java.util.Collections;
 import java.util.List;
 import java.util.Map;
@@ -550,6 +551,24 @@ public abstract class AbstractSubscriptionPullConsumer 
extends AbstractSubscript
     super.commitAsync(messages, callback);
   }
 
+  @Override
+  protected void onAckedCommitContexts(
+      final Collection<SubscriptionCommitContext> acceptedCommitContexts) {
+    if (!autoCommit
+        || Objects.isNull(uncommittedCommitContexts)
+        || acceptedCommitContexts.isEmpty()) {
+      return;
+    }
+
+    for (final Map.Entry<Long, Set<SubscriptionCommitContext>> entry :
+        uncommittedCommitContexts.entrySet()) {
+      entry.getValue().removeAll(acceptedCommitContexts);
+      if (entry.getValue().isEmpty()) {
+        uncommittedCommitContexts.remove(entry.getKey(), entry.getValue());
+      }
+    }
+  }
+
   private List<SubscriptionMessage> filterUserVisibleMessages(
       final List<SubscriptionMessage> messages) {
     if (messages.isEmpty()) {
diff --git 
a/iotdb-client/subscription/src/test/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionConsumerLifecycleTest.java
 
b/iotdb-client/subscription/src/test/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionConsumerLifecycleTest.java
index 00db8d31145..f67d8829263 100644
--- 
a/iotdb-client/subscription/src/test/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionConsumerLifecycleTest.java
+++ 
b/iotdb-client/subscription/src/test/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionConsumerLifecycleTest.java
@@ -22,6 +22,7 @@ package org.apache.iotdb.session.subscription.consumer.base;
 import org.apache.iotdb.common.rpc.thrift.TEndPoint;
 import 
org.apache.iotdb.rpc.subscription.exception.SubscriptionConsumerFencedException;
 import org.apache.iotdb.rpc.subscription.exception.SubscriptionException;
+import 
org.apache.iotdb.rpc.subscription.exception.SubscriptionRuntimeNonCriticalException;
 import 
org.apache.iotdb.rpc.subscription.payload.poll.SubscriptionCommitContext;
 import org.apache.iotdb.rpc.subscription.payload.poll.SubscriptionPollResponse;
 import 
org.apache.iotdb.rpc.subscription.payload.poll.SubscriptionPollResponseType;
@@ -40,10 +41,12 @@ import java.lang.reflect.Field;
 import java.util.ArrayList;
 import java.util.Arrays;
 import java.util.Collections;
+import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
 import java.util.Objects;
 import java.util.Set;
+import java.util.SortedMap;
 import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.CountDownLatch;
 import java.util.concurrent.ExecutionException;
@@ -279,6 +282,92 @@ public class SubscriptionConsumerLifecycleTest {
     }
   }
 
+  @Test
+  public void testSyncCommitClearsAcceptedContextsFromAutoCommitBuffer() 
throws Exception {
+    final TestPullConsumer consumer = new TestPullConsumer(true);
+    final SubscriptionCommitContext commitContext =
+        new SubscriptionCommitContext(0, 0, "topic", CONSUMER_GROUP_ID, 1L);
+    try {
+      consumer.open();
+      addUncommittedCommitContexts(consumer, commitContext);
+
+      consumer.commitSync(new SubscriptionMessage(commitContext, 1L));
+
+      Assert.assertTrue(getUncommittedCommitContexts(consumer).isEmpty());
+      Assert.assertEquals(1, consumer.commitRequestCount);
+    } finally {
+      consumer.close();
+    }
+  }
+
+  @Test
+  public void testAsyncCommitClearsAcceptedContextsFromAutoCommitBuffer() 
throws Exception {
+    final TestPullConsumer consumer = new TestPullConsumer(true);
+    final SubscriptionCommitContext commitContext =
+        new SubscriptionCommitContext(0, 0, "topic", CONSUMER_GROUP_ID, 1L);
+    try {
+      consumer.open();
+      addUncommittedCommitContexts(consumer, commitContext);
+
+      consumer.commitAsync(new SubscriptionMessage(commitContext, 1L)).get(5, 
TimeUnit.SECONDS);
+
+      Assert.assertTrue(getUncommittedCommitContexts(consumer).isEmpty());
+      Assert.assertEquals(1, consumer.commitRequestCount);
+    } finally {
+      consumer.close();
+    }
+  }
+
+  @Test
+  public void testPartialCommitKeepsRejectedContextsInAutoCommitBuffer() 
throws Exception {
+    final TestPullConsumer consumer = new TestPullConsumer(true);
+    final SubscriptionCommitContext acceptedCommitContext =
+        new SubscriptionCommitContext(0, 0, "topic", CONSUMER_GROUP_ID, 1L);
+    final SubscriptionCommitContext rejectedCommitContext =
+        new SubscriptionCommitContext(0, 0, "topic", CONSUMER_GROUP_ID, 2L);
+    consumer.rejectedCommitContexts.add(rejectedCommitContext);
+    try {
+      consumer.open();
+      addUncommittedCommitContexts(consumer, acceptedCommitContext, 
rejectedCommitContext);
+
+      try {
+        consumer.commitSync(
+            Arrays.asList(
+                new SubscriptionMessage(acceptedCommitContext, 1L),
+                new SubscriptionMessage(rejectedCommitContext, 2L)));
+        Assert.fail("A partially accepted commit must fail");
+      } catch (final SubscriptionRuntimeNonCriticalException expected) {
+        Assert.assertTrue(expected.getMessage().contains("partially 
accepted"));
+      }
+
+      final SortedMap<Long, Set<SubscriptionCommitContext>> 
uncommittedCommitContexts =
+          getUncommittedCommitContexts(consumer);
+      Assert.assertEquals(1, uncommittedCommitContexts.size());
+      Assert.assertEquals(
+          Collections.singleton(rejectedCommitContext), 
uncommittedCommitContexts.get(0L));
+      Assert.assertEquals(1, consumer.commitRequestCount);
+    } finally {
+      consumer.close();
+    }
+  }
+
+  private static void addUncommittedCommitContexts(
+      final TestPullConsumer consumer, final SubscriptionCommitContext... 
commitContexts)
+      throws Exception {
+    getUncommittedCommitContexts(consumer)
+        .computeIfAbsent(0L, ignored -> new HashSet<>())
+        .addAll(Arrays.asList(commitContexts));
+  }
+
+  @SuppressWarnings("unchecked")
+  private static SortedMap<Long, Set<SubscriptionCommitContext>> 
getUncommittedCommitContexts(
+      final TestPullConsumer consumer) throws Exception {
+    final Field field =
+        
AbstractSubscriptionPullConsumer.class.getDeclaredField("uncommittedCommitContexts");
+    field.setAccessible(true);
+    return (SortedMap<Long, Set<SubscriptionCommitContext>>) 
field.get(consumer);
+  }
+
   private AbstractSubscriptionProviders getProviders(final 
AbstractSubscriptionConsumer consumer)
       throws Exception {
     final Field field = 
AbstractSubscriptionConsumer.class.getDeclaredField("providers");
@@ -366,6 +455,7 @@ public class SubscriptionConsumerLifecycleTest {
           () -> {},
           () -> {},
           () -> {},
+          Collections.emptySet(),
           null,
           null);
     }
@@ -383,15 +473,27 @@ public class SubscriptionConsumerLifecycleTest {
     private int closeRequestCount;
     private int sessionCloseCount;
     private int commitRequestCount;
+    private final Set<SubscriptionCommitContext> rejectedCommitContexts = new 
HashSet<>();
     private final CountDownLatch providerCloseStarted;
     private final CountDownLatch allowProviderClose;
 
     private TestPullConsumer() {
-      this(null, null);
+      this(false, null, null);
+    }
+
+    private TestPullConsumer(final boolean autoCommit) {
+      this(autoCommit, null, null);
     }
 
     private TestPullConsumer(
         final CountDownLatch providerCloseStarted, final CountDownLatch 
allowProviderClose) {
+      this(false, providerCloseStarted, allowProviderClose);
+    }
+
+    private TestPullConsumer(
+        final boolean autoCommit,
+        final CountDownLatch providerCloseStarted,
+        final CountDownLatch allowProviderClose) {
       super(
           new AbstractSubscriptionPullConsumerBuilder()
               .host(HOST)
@@ -400,7 +502,8 @@ public class SubscriptionConsumerLifecycleTest {
               .consumerGroupId(CONSUMER_GROUP_ID)
               .heartbeatIntervalMs(LONG_INTERVAL_MS)
               .endpointsSyncIntervalMs(LONG_INTERVAL_MS)
-              .autoCommit(false));
+              .autoCommit(autoCommit)
+              .autoCommitIntervalMs(LONG_INTERVAL_MS));
       this.providerCloseStarted = providerCloseStarted;
       this.allowProviderClose = allowProviderClose;
     }
@@ -441,6 +544,7 @@ public class SubscriptionConsumerLifecycleTest {
               () -> closeRequestCount++,
               () -> commitRequestCount++,
               () -> sessionCloseCount++,
+              rejectedCommitContexts,
               providerCloseStarted,
               allowProviderClose);
       createdProviders.add(provider);
@@ -460,6 +564,7 @@ public class SubscriptionConsumerLifecycleTest {
     private final Runnable closeRequest;
     private final Runnable commitRequest;
     private final Runnable sessionClose;
+    private final Set<SubscriptionCommitContext> rejectedCommitContexts;
     private final CountDownLatch providerCloseStarted;
     private final CountDownLatch allowProviderClose;
 
@@ -485,6 +590,7 @@ public class SubscriptionConsumerLifecycleTest {
         final Runnable closeRequest,
         final Runnable commitRequest,
         final Runnable sessionClose,
+        final Set<SubscriptionCommitContext> rejectedCommitContexts,
         final CountDownLatch providerCloseStarted,
         final CountDownLatch allowProviderClose) {
       super(
@@ -509,6 +615,7 @@ public class SubscriptionConsumerLifecycleTest {
       this.closeRequest = closeRequest;
       this.commitRequest = commitRequest;
       this.sessionClose = sessionClose;
+      this.rejectedCommitContexts = rejectedCommitContexts;
       this.providerCloseStarted = providerCloseStarted;
       this.allowProviderClose = allowProviderClose;
     }
@@ -603,7 +710,13 @@ public class SubscriptionConsumerLifecycleTest {
     CommitResult commit(
         final List<SubscriptionCommitContext> subscriptionCommitContexts, 
final boolean nack) {
       commitRequest.run();
-      return CommitResult.empty();
+      final List<SubscriptionCommitContext> acceptedCommitContexts = new 
ArrayList<>();
+      for (final SubscriptionCommitContext commitContext : 
subscriptionCommitContexts) {
+        if (!rejectedCommitContexts.contains(commitContext)) {
+          acceptedCommitContexts.add(commitContext);
+        }
+      }
+      return new CommitResult(acceptedCommitContexts, Collections.emptyMap());
     }
   }
 }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionPrefetchingQueue.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionPrefetchingQueue.java
index 801fc03d3b8..2f37ac6f64b 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionPrefetchingQueue.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionPrefetchingQueue.java
@@ -872,7 +872,7 @@ public abstract class SubscriptionPrefetchingQueue {
         new Pair<>(consumerId, commitContext),
         (key, ev) -> {
           if (Objects.isNull(ev)) {
-            LOGGER.warn(
+            LOGGER.debug(
                 DataNodePipeMessages
                     
.PIPE_LOG_SUBSCRIPTION_SUBSCRIPTION_COMMIT_CONTEXT_DOES_NOT_EXIST_0E4EF990,
                 commitContext,
@@ -949,7 +949,7 @@ public abstract class SubscriptionPrefetchingQueue {
         new Pair<>(consumerId, commitContext),
         (key, ev) -> {
           if (Objects.isNull(ev)) {
-            LOGGER.warn(
+            LOGGER.debug(
                 DataNodePipeMessages
                     
.PIPE_LOG_SUBSCRIPTION_SUBSCRIPTION_COMMIT_CONTEXT_DOES_NOT_EXIST_DE907E05,
                 commitContext,
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
index 4c8c1c08d30..697645cde13 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
@@ -953,6 +953,25 @@ public class SubscriptionReceiverV1 implements 
SubscriptionReceiver {
             nack,
             commitContexts);
       }
+    } else if (acceptedCommitContexts.isEmpty()) {
+      LOGGER.debug(
+          DataNodePipeMessages
+              
.PIPE_LOG_SUBSCRIPTION_CONSUMER_COMMIT_NACK_PARTIALLY_ACCEPTED_REQUESTED_87D0C038,
+          consumerConfig,
+          nack,
+          summarizeCommitContexts(commitContexts),
+          summarizeCommitContexts(acceptedCommitContexts),
+          summarizeCommitContexts(staleUnsubscribedCommitContexts));
+      if (LOGGER.isDebugEnabled()) {
+        LOGGER.debug(
+            DataNodePipeMessages
+                
.PIPE_LOG_SUBSCRIPTION_CONSUMER_COMMIT_NACK_FULL_REQUESTED_COMMIT_1E67E8A3,
+            consumerConfig,
+            nack,
+            commitContexts,
+            acceptedCommitContexts,
+            staleUnsubscribedCommitContexts);
+      }
     } else {
       LOGGER.warn(
           DataNodePipeMessages

Reply via email to