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