This is an automated email from the ASF dual-hosted git repository.

Caideyipi 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 5db9f416143 fix(subscription): avoid per-provider poll backoff (#18448)
5db9f416143 is described below

commit 5db9f416143af6bd51518990152c22da100cf4c7
Author: Caideyipi <[email protected]>
AuthorDate: Tue Aug 18 12:15:44 2026 +0800

    fix(subscription): avoid per-provider poll backoff (#18448)
    
    * fix(subscription): avoid per-provider poll backoff
    
    * fix(subscription): remove added TODO marker
---
 .../base/AbstractSubscriptionConsumer.java         |  24 +-
 .../base/AbstractSubscriptionProviders.java        |   8 +
 .../SubscriptionConsumerMultiProviderPollTest.java | 281 +++++++++++++++++++++
 3 files changed, 310 insertions(+), 3 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 fa48b6a099a..38b16ed7cf9 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
@@ -824,6 +824,9 @@ abstract class AbstractSubscriptionConsumer implements 
AutoCloseable {
     final List<SubscriptionMessage> messages = new ArrayList<>();
     List<SubscriptionPollResponse> currentResponses = new ArrayList<>();
     final PollTimer timer = new PollTimer(System.currentTimeMillis(), 
timeoutMs);
+    // Poll every available provider before backing off. Otherwise an idle 
provider adds the random
+    // backoff latency even when the next provider already has data ready.
+    int remainingProvidersBeforeBackoff = getAvailableProviderCount();
 
     try {
       do {
@@ -901,9 +904,10 @@ abstract class AbstractSubscriptionConsumer implements 
AutoCloseable {
         // update timer
         timer.update();
 
-        // TODO: associated with timeoutMs instead of hardcoding
-        // random sleep time within the range [SLEEP_DELTA_MS, SLEEP_DELTA_MS 
+ SLEEP_MS)
-        Thread.sleep(((long) (Math.random() * SLEEP_MS)) + SLEEP_DELTA_MS);
+        if (--remainingProvidersBeforeBackoff <= 0) {
+          sleepAfterEmptyPollRound();
+          remainingProvidersBeforeBackoff = getAvailableProviderCount();
+        }
 
         // the use of TIMER_DELTA_MS here slightly reduces the timeout to 
avoid being interrupted as
         // much as possible
@@ -933,6 +937,20 @@ abstract class AbstractSubscriptionConsumer implements 
AutoCloseable {
     return messages;
   }
 
+  private int getAvailableProviderCount() {
+    providers.acquireReadLock();
+    try {
+      return providers.getAvailableProviderCount();
+    } finally {
+      providers.releaseReadLock();
+    }
+  }
+
+  void sleepAfterEmptyPollRound() throws InterruptedException {
+    // Randomize the pause between fully empty provider rounds to avoid 
synchronized polling.
+    Thread.sleep(((long) (Math.random() * SLEEP_MS)) + SLEEP_DELTA_MS);
+  }
+
   private Optional<SubscriptionMessage> pollFile(
       final SubscriptionPollResponse response, final PollTimer timer) throws 
SubscriptionException {
     final SubscriptionCommitContext commitContext = 
response.getCommitContext();
diff --git 
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionProviders.java
 
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionProviders.java
index d866f50749c..be164a7da6a 100644
--- 
a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionProviders.java
+++ 
b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionProviders.java
@@ -201,6 +201,14 @@ final class AbstractSubscriptionProviders {
         .noneMatch(AbstractSubscriptionProvider::isAvailable);
   }
 
+  /** Caller should ensure that the method is called in the lock {@link 
#acquireReadLock()}. */
+  int getAvailableProviderCount() {
+    return (int)
+        subscriptionProviders.values().stream()
+            .filter(AbstractSubscriptionProvider::isAvailable)
+            .count();
+  }
+
   /** Caller should ensure that the method is called in the lock {@link 
#acquireReadLock()}. */
   boolean containsProvider(final int dataNodeId) {
     return subscriptionProviders.containsKey(dataNodeId);
diff --git 
a/iotdb-client/subscription/src/test/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionConsumerMultiProviderPollTest.java
 
b/iotdb-client/subscription/src/test/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionConsumerMultiProviderPollTest.java
new file mode 100644
index 00000000000..34ef7a8fc48
--- /dev/null
+++ 
b/iotdb-client/subscription/src/test/java/org/apache/iotdb/session/subscription/consumer/base/SubscriptionConsumerMultiProviderPollTest.java
@@ -0,0 +1,281 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.session.subscription.consumer.base;
+
+import org.apache.iotdb.common.rpc.thrift.TEndPoint;
+import org.apache.iotdb.rpc.subscription.config.TopicConfig;
+import org.apache.iotdb.rpc.subscription.exception.SubscriptionException;
+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;
+import org.apache.iotdb.rpc.subscription.payload.poll.TabletsPayload;
+import org.apache.iotdb.rpc.subscription.payload.poll.TopicProgress;
+import 
org.apache.iotdb.rpc.subscription.payload.response.PipeSubscribeHeartbeatResp;
+import org.apache.iotdb.session.AbstractSessionBuilder;
+import org.apache.iotdb.session.subscription.SubscriptionTreeSessionBuilder;
+import org.apache.iotdb.session.subscription.payload.SubscriptionMessage;
+
+import org.apache.tsfile.enums.TSDataType;
+import org.apache.tsfile.write.record.Tablet;
+import org.apache.tsfile.write.schema.IMeasurementSchema;
+import org.apache.tsfile.write.schema.MeasurementSchema;
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Set;
+
+public class SubscriptionConsumerMultiProviderPollTest {
+
+  private static final String HOST = "127.0.0.1";
+  private static final int FIRST_PORT = 10_001;
+  private static final String TOPIC = "topic1";
+  private static final String CONSUMER_ID = "test_consumer";
+  private static final String CONSUMER_GROUP_ID = "test_consumer_group";
+  private static final long LONG_INTERVAL_MS = 86_400_000L;
+
+  @Test
+  public void testPollTriesNextProviderBeforeBackingOff() throws 
SubscriptionException {
+    final TestPullConsumer consumer = new TestPullConsumer(2, 0);
+    try {
+      consumer.open();
+      consumer.subscribeTopic();
+
+      final List<SubscriptionMessage> messages = consumer.pollForTest(1_000L);
+
+      Assert.assertEquals(1, messages.size());
+      Assert.assertEquals(1, consumer.getPollCount(1));
+      Assert.assertEquals(1, consumer.getPollCount(2));
+      Assert.assertEquals(0, consumer.getPollCount(3));
+      Assert.assertEquals(0, consumer.getBackoffCount());
+    } finally {
+      consumer.close();
+    }
+  }
+
+  @Test
+  public void testPollBacksOffAfterAllProvidersAreEmpty() throws 
SubscriptionException {
+    final TestPullConsumer consumer = new TestPullConsumer(1, 1);
+    try {
+      consumer.open();
+      consumer.subscribeTopic();
+
+      final List<SubscriptionMessage> messages = consumer.pollForTest(1_000L);
+
+      Assert.assertEquals(1, messages.size());
+      Assert.assertEquals(2, consumer.getPollCount(1));
+      Assert.assertEquals(1, consumer.getPollCount(2));
+      Assert.assertEquals(1, consumer.getPollCount(3));
+      Assert.assertEquals(1, consumer.getBackoffCount());
+    } finally {
+      consumer.close();
+    }
+  }
+
+  private static class TestPullConsumer extends 
AbstractSubscriptionPullConsumer {
+
+    private final Map<Integer, Integer> pollCounts = new HashMap<>();
+    private final int dataProviderId;
+    private final int emptyPollsBeforeData;
+    private int backoffCount;
+
+    private TestPullConsumer(final int dataProviderId, final int 
emptyPollsBeforeData) {
+      super(
+          new AbstractSubscriptionPullConsumerBuilder()
+              .host(HOST)
+              .port(FIRST_PORT)
+              .consumerId(CONSUMER_ID)
+              .consumerGroupId(CONSUMER_GROUP_ID)
+              .heartbeatIntervalMs(LONG_INTERVAL_MS)
+              .endpointsSyncIntervalMs(LONG_INTERVAL_MS)
+              .autoCommit(false));
+      this.dataProviderId = dataProviderId;
+      this.emptyPollsBeforeData = emptyPollsBeforeData;
+    }
+
+    @Override
+    protected AbstractSubscriptionProvider constructSubscriptionProvider(
+        final TEndPoint endPoint,
+        final String username,
+        final String password,
+        final String encryptedPassword,
+        final String consumerId,
+        final String consumerGroupId,
+        final String ownerId,
+        final Long ownerEpoch,
+        final int thriftMaxFrameSize,
+        final long heartbeatIntervalMs,
+        final int connectionTimeoutInMs) {
+      return new TestSubscriptionProvider(
+          endPoint,
+          username,
+          password,
+          encryptedPassword,
+          consumerId,
+          consumerGroupId,
+          ownerId,
+          ownerEpoch,
+          thriftMaxFrameSize,
+          heartbeatIntervalMs,
+          connectionTimeoutInMs,
+          pollCounts,
+          dataProviderId,
+          emptyPollsBeforeData);
+    }
+
+    private void subscribeTopic() {
+      subscribedTopics = Collections.singletonMap(TOPIC, new TopicConfig());
+    }
+
+    private List<SubscriptionMessage> pollForTest(final long timeoutMs)
+        throws SubscriptionException {
+      return poll(timeoutMs);
+    }
+
+    private int getPollCount(final int dataNodeId) {
+      return pollCounts.getOrDefault(dataNodeId, 0);
+    }
+
+    private int getBackoffCount() {
+      return backoffCount;
+    }
+
+    @Override
+    void sleepAfterEmptyPollRound() {
+      backoffCount++;
+    }
+  }
+
+  private static class TestSubscriptionProvider extends 
AbstractSubscriptionProvider {
+
+    private final int dataNodeId;
+    private final Map<Integer, Integer> pollCounts;
+    private final int dataProviderId;
+    private final int emptyPollsBeforeData;
+
+    private TestSubscriptionProvider(
+        final TEndPoint endPoint,
+        final String username,
+        final String password,
+        final String encryptedPassword,
+        final String consumerId,
+        final String consumerGroupId,
+        final String ownerId,
+        final Long ownerEpoch,
+        final int thriftMaxFrameSize,
+        final long heartbeatIntervalMs,
+        final int connectionTimeoutInMs,
+        final Map<Integer, Integer> pollCounts,
+        final int dataProviderId,
+        final int emptyPollsBeforeData) {
+      super(
+          endPoint,
+          username,
+          password,
+          encryptedPassword,
+          consumerId,
+          consumerGroupId,
+          ownerId,
+          ownerEpoch,
+          thriftMaxFrameSize,
+          heartbeatIntervalMs,
+          connectionTimeoutInMs);
+      this.dataNodeId = endPoint.port - FIRST_PORT + 1;
+      this.pollCounts = pollCounts;
+      this.dataProviderId = dataProviderId;
+      this.emptyPollsBeforeData = emptyPollsBeforeData;
+    }
+
+    @Override
+    protected AbstractSessionBuilder constructSubscriptionSessionBuilder(
+        final String host,
+        final int port,
+        final String username,
+        final String password,
+        final String encryptedPassword,
+        final int thriftMaxFrameSize,
+        final int connectionTimeoutInMs) {
+      final boolean useEncryptedPassword = Objects.nonNull(encryptedPassword);
+      return new SubscriptionTreeSessionBuilder()
+          .host(host)
+          .port(port)
+          .username(username)
+          .password(useEncryptedPassword ? encryptedPassword : password)
+          .useEncryptedPassword(useEncryptedPassword)
+          .thriftMaxFrameSize(thriftMaxFrameSize)
+          .connectionTimeoutInMs(connectionTimeoutInMs);
+    }
+
+    @Override
+    synchronized void handshake() {
+      setAvailable();
+    }
+
+    @Override
+    synchronized void close() {
+      setUnavailable();
+    }
+
+    @Override
+    int getDataNodeId() {
+      return dataNodeId;
+    }
+
+    @Override
+    PipeSubscribeHeartbeatResp heartbeat(
+        final List<SubscriptionCommitContext> processorBufferedCommitContexts) 
{
+      final PipeSubscribeHeartbeatResp response = new 
PipeSubscribeHeartbeatResp();
+      response.getTopics().put(TOPIC, new TopicConfig());
+      response.getEndPoints().put(1, new TEndPoint(HOST, FIRST_PORT));
+      response.getEndPoints().put(2, new TEndPoint(HOST, FIRST_PORT + 1));
+      response.getEndPoints().put(3, new TEndPoint(HOST, FIRST_PORT + 2));
+      return response;
+    }
+
+    @Override
+    List<SubscriptionPollResponse> poll(
+        final Set<String> topicNames,
+        final long timeoutMs,
+        final Map<String, TopicProgress> progressByTopic)
+        throws SubscriptionException {
+      pollCounts.merge(dataNodeId, 1, Integer::sum);
+      if (dataNodeId != dataProviderId || pollCounts.get(dataNodeId) <= 
emptyPollsBeforeData) {
+        return Collections.emptyList();
+      }
+      final SubscriptionCommitContext commitContext =
+          new SubscriptionCommitContext(dataNodeId, 0, TOPIC, 
CONSUMER_GROUP_ID, 0L);
+      final List<IMeasurementSchema> schemas =
+          Collections.singletonList(new MeasurementSchema("s1", 
TSDataType.INT64));
+      final Tablet tablet = new Tablet("root.sg.d1", schemas, 1);
+      tablet.setTimestamps(new long[] {1L});
+      ((long[]) tablet.getValues()[0])[0] = 1L;
+      tablet.setRowSize(1);
+      return Collections.singletonList(
+          new SubscriptionPollResponse(
+              SubscriptionPollResponseType.TABLETS.getType(),
+              new TabletsPayload(Collections.singletonList(tablet), -1),
+              commitContext));
+    }
+  }
+}

Reply via email to