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

healchow pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/inlong.git


The following commit(s) were added to refs/heads/master by this push:
     new 14cff90e3 [INLONG-5621][SDK] Support multi-topic fetcher for Pulsar 
(#5625)
14cff90e3 is described below

commit 14cff90e373897c6dad4f27f2ad9e3282bf4625a
Author: vernedeng <[email protected]>
AuthorDate: Wed Aug 24 16:02:38 2022 +0800

    [INLONG-5621][SDK] Support multi-topic fetcher for Pulsar (#5625)
---
 .../inlong/sdk/sort/api/MultiTopicsFetcher.java    |  68 ++++
 .../inlong/sdk/sort/api/SortClientConfig.java      |  11 +
 .../sdk/sort/fetcher/pulsar/PulsarConsumer.java    |  96 +++++
 .../fetcher/pulsar/PulsarMultiTopicsFetcher.java   | 434 +++++++++++++++++++++
 4 files changed, 609 insertions(+)

diff --git 
a/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/api/MultiTopicsFetcher.java
 
b/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/api/MultiTopicsFetcher.java
new file mode 100644
index 000000000..7f7ab6ec4
--- /dev/null
+++ 
b/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/api/MultiTopicsFetcher.java
@@ -0,0 +1,68 @@
+/*
+ * 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.inlong.sdk.sort.api;
+
+import org.apache.inlong.sdk.sort.entity.InLongTopic;
+
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.locks.ReentrantReadWriteLock;
+import java.util.stream.Collectors;
+
+/**
+ * Basic class of multi topic fetchers.
+ * The main differences between this and {@link SingleTopicFetcher} is that:
+ * 1. MultiTopicFetcher maintains a list of topics while {@link 
SingleTopicFetcher} only maintains one;
+ * 2. All topics share the same properties;
+ * 3. The joining and removing of topics will result in new creation of 
consumer,
+ *      and the old ones will be put in a list, waiting to be cleaned by a 
scheduled thread.
+ */
+public abstract class MultiTopicsFetcher implements TopicFetcher {
+    protected final ReentrantReadWriteLock mainLock = new 
ReentrantReadWriteLock(true);
+    protected final ScheduledExecutorService executor;
+    protected Map<String, InLongTopic> onlineTopics;
+    protected ClientContext context;
+    protected Deserializer deserializer;
+    protected volatile Thread fetchThread;
+    protected volatile boolean closed = false;
+    protected volatile boolean stopConsume = false;
+    // use for empty topic to sleep
+    protected long sleepTime = 0L;
+    protected int emptyFetchTimes = 0;
+    // for rollback
+    protected Interceptor interceptor;
+    protected Seeker seeker;
+
+    public MultiTopicsFetcher(
+            List<InLongTopic> topics,
+            ClientContext context,
+            Interceptor interceptor,
+            Deserializer deserializer) {
+        this.onlineTopics = topics.stream()
+                .filter(Objects::nonNull)
+                .collect((Collectors.toMap(InLongTopic::getTopic, t -> t)));
+        this.context = context;
+        this.interceptor = interceptor;
+        this.deserializer = deserializer;
+        this.executor = Executors.newSingleThreadScheduledExecutor();
+    }
+
+}
diff --git 
a/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/api/SortClientConfig.java
 
b/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/api/SortClientConfig.java
index cba627705..473075f89 100644
--- 
a/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/api/SortClientConfig.java
+++ 
b/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/api/SortClientConfig.java
@@ -58,6 +58,7 @@ public class SortClientConfig implements Serializable {
     private int emptyPollSleepStepMs = 50;
     private int maxEmptyPollSleepMs = 500;
     private int emptyPollTimes = 10;
+    private int cleanOldConsumerIntervalSec = 60;
 
     public SortClientConfig(String sortTaskId, String sortClusterName, 
InLongTopicChangeListener assignmentsListener,
             ConsumeStrategy consumeStrategy, String localIp) {
@@ -314,6 +315,14 @@ public class SortClientConfig implements Serializable {
         this.emptyPollTimes = emptyPollTimes;
     }
 
+    public int getCleanOldConsumerIntervalSec() {
+        return cleanOldConsumerIntervalSec;
+    }
+
+    public void setCleanOldConsumerIntervalSec(int 
cleanOldConsumerIntervalSec) {
+        this.cleanOldConsumerIntervalSec = cleanOldConsumerIntervalSec;
+    }
+
     /**
      * ConsumeStrategy
      */
@@ -367,6 +376,8 @@ public class SortClientConfig implements Serializable {
         this.updateMetaDataIntervalSec = 
NumberUtils.toInt(sortSdkParams.get("updateMetaDataIntervalSec"),
                 updateMetaDataIntervalSec);
         this.ackTimeoutSec = 
NumberUtils.toInt(sortSdkParams.get("ackTimeoutSec"), ackTimeoutSec);
+        this.cleanOldConsumerIntervalSec = 
NumberUtils.toInt(sortSdkParams.get("cleanOldConsumerIntervalSec"),
+                cleanOldConsumerIntervalSec);
 
         String strPrometheusEnabled = 
sortSdkParams.getOrDefault("isPrometheusEnabled", Boolean.TRUE.toString());
         this.isPrometheusEnabled = 
StringUtils.equalsIgnoreCase(strPrometheusEnabled, Boolean.TRUE.toString());
diff --git 
a/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/fetcher/pulsar/PulsarConsumer.java
 
b/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/fetcher/pulsar/PulsarConsumer.java
new file mode 100644
index 000000000..5aa2f423b
--- /dev/null
+++ 
b/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/fetcher/pulsar/PulsarConsumer.java
@@ -0,0 +1,96 @@
+/*
+ * 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.inlong.sdk.sort.fetcher.pulsar;
+
+import org.apache.inlong.sdk.sort.entity.InLongTopic;
+import org.apache.inlong.tubemq.corebase.utils.Tuple2;
+import org.apache.pulsar.client.api.Consumer;
+import org.apache.pulsar.client.api.MessageId;
+import org.apache.pulsar.client.api.Messages;
+import org.apache.pulsar.client.api.PulsarClientException;
+
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.ConcurrentHashMap;
+
+/**
+ * Wrapper of pulsar consumer.
+ */
+public class PulsarConsumer {
+    private final ConcurrentHashMap<String, Tuple2<InLongTopic, MessageId>> 
offsetCache = new ConcurrentHashMap<>();
+    private final Consumer<byte[]> consumer;
+    private long stopTime = -1;
+
+    public PulsarConsumer(Consumer<byte[]> consumer) {
+        this.consumer = consumer;
+    }
+
+    public void close() throws PulsarClientException {
+        this.consumer.close();
+        this.offsetCache.clear();
+    }
+
+    public void pause() {
+        this.consumer.pause();
+    }
+
+    public void resume() {
+        this.consumer.resume();
+    }
+
+    public Messages<byte[]> batchReceive() throws PulsarClientException {
+        return this.consumer.batchReceive();
+    }
+
+    public CompletableFuture<Void> acknowledgeAsync(MessageId messageId) {
+        return this.consumer.acknowledgeAsync(messageId);
+    }
+
+    public long getStopTime() {
+        return stopTime;
+    }
+
+    public void setStopTime(long stopTime) {
+        this.stopTime = stopTime;
+    }
+
+    public InLongTopic getTopic(String msgOffset) {
+        Tuple2<InLongTopic, MessageId> tuple = offsetCache.get(msgOffset);
+        return tuple == null ? null : tuple.getF0();
+    }
+
+    public MessageId getMessageId(String msgOffset) {
+        Tuple2<InLongTopic, MessageId> tuple = offsetCache.get(msgOffset);
+        return tuple == null ? null : tuple.getF1();
+    }
+
+    public boolean remove(String offsetKey) {
+        return offsetCache.remove(offsetKey) != null;
+    }
+
+    public void put(String offsetKey, InLongTopic topic, MessageId messageId) {
+        offsetCache.put(offsetKey, new Tuple2<>(topic, messageId));
+    }
+
+    public boolean isEmpty() {
+        return offsetCache.isEmpty();
+    }
+
+    public boolean isConnected() {
+        return consumer.isConnected();
+    }
+}
diff --git 
a/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/fetcher/pulsar/PulsarMultiTopicsFetcher.java
 
b/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/fetcher/pulsar/PulsarMultiTopicsFetcher.java
new file mode 100644
index 000000000..3244e9d5b
--- /dev/null
+++ 
b/inlong-sdk/sort-sdk/src/main/java/org/apache/inlong/sdk/sort/fetcher/pulsar/PulsarMultiTopicsFetcher.java
@@ -0,0 +1,434 @@
+/*
+ * 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.inlong.sdk.sort.fetcher.pulsar;
+
+import com.google.common.base.Preconditions;
+import org.apache.commons.collections.CollectionUtils;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.inlong.sdk.sort.api.ClientContext;
+import org.apache.inlong.sdk.sort.api.Deserializer;
+import org.apache.inlong.sdk.sort.api.Interceptor;
+import org.apache.inlong.sdk.sort.api.MultiTopicsFetcher;
+import org.apache.inlong.sdk.sort.api.Seeker;
+import org.apache.inlong.sdk.sort.api.SeekerFactory;
+import org.apache.inlong.sdk.sort.api.SortClientConfig;
+import org.apache.inlong.sdk.sort.entity.InLongMessage;
+import org.apache.inlong.sdk.sort.entity.InLongTopic;
+import org.apache.inlong.sdk.sort.entity.MessageRecord;
+import org.apache.pulsar.client.api.Consumer;
+import org.apache.pulsar.client.api.Message;
+import org.apache.pulsar.client.api.MessageId;
+import org.apache.pulsar.client.api.Messages;
+import org.apache.pulsar.client.api.PulsarClient;
+import org.apache.pulsar.client.api.PulsarClientException;
+import org.apache.pulsar.client.api.Schema;
+import org.apache.pulsar.client.api.SubscriptionInitialPosition;
+import org.apache.pulsar.client.api.SubscriptionType;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.ArrayList;
+import java.util.Base64;
+import java.util.Collection;
+import java.util.LinkedList;
+import java.util.List;
+import java.util.Objects;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+
+/**
+ * MultiTopicsFetcher for pulsar.
+ */
+public class PulsarMultiTopicsFetcher extends MultiTopicsFetcher {
+    private static final Logger LOGGER = 
LoggerFactory.getLogger(PulsarMultiTopicsFetcher.class);
+    private PulsarConsumer currentConsumer;
+    private List<PulsarConsumer> toBeRemovedConsumers = new LinkedList<>();
+    private PulsarClient pulsarClient;
+
+    public PulsarMultiTopicsFetcher(
+            List<InLongTopic> topics,
+            ClientContext context,
+            Interceptor interceptor,
+            Deserializer deserializer,
+            PulsarClient pulsarClient) {
+        super(topics, context, interceptor, deserializer);
+        this.pulsarClient = Preconditions.checkNotNull(pulsarClient);
+    }
+
+    @Override
+    public boolean init() {
+        Consumer<byte[]> newConsumer = createConsumer(onlineTopics.values());
+        if (Objects.isNull(newConsumer)) {
+            LOGGER.error("create new consumer is null");
+            return false;
+        }
+        this.currentConsumer = new PulsarConsumer(newConsumer);
+        InLongTopic firstTopic = 
onlineTopics.values().stream().findFirst().get();
+        this.seeker = SeekerFactory.createPulsarSeeker(newConsumer, 
firstTopic);
+        String threadName = 
String.format("sort_sdk_pulsar_multi_topic_fetch_thread_%d", this.hashCode());
+        this.fetchThread = new Thread(new PulsarMultiTopicsFetcher.Fetcher(), 
threadName);
+        this.fetchThread.start();
+        this.executor.scheduleWithFixedDelay(this::clearRemovedConsumerList,
+                context.getConfig().getCleanOldConsumerIntervalSec(),
+                context.getConfig().getCleanOldConsumerIntervalSec(),
+                TimeUnit.SECONDS);
+        return true;
+    }
+
+    private void clearRemovedConsumerList() {
+        long cur = System.currentTimeMillis();
+        List<PulsarConsumer> newList = new LinkedList<>();
+        toBeRemovedConsumers.forEach(consumer -> {
+            long diff = cur - consumer.getStopTime();
+            if (diff > context.getConfig().getCleanOldConsumerIntervalSec() * 
1000L || consumer.isEmpty()) {
+                try {
+                    consumer.close();
+                } catch (PulsarClientException e) {
+                    LOGGER.warn("got exception in close old consumer", e);
+                }
+                return;
+            }
+            newList.add(consumer);
+        });
+        LOGGER.info("after clear old consumers, the old size is {}, current 
size is {}",
+                toBeRemovedConsumers.size(), newList.size());
+        this.toBeRemovedConsumers = newList;
+    }
+
+    private boolean updateAll(Collection<InLongTopic> newTopics) {
+        if (CollectionUtils.isEmpty(newTopics)) {
+            LOGGER.error("new topics is empty or null");
+            return false;
+        }
+        // stop old;
+        this.setStopConsume(true);
+        this.currentConsumer.pause();
+        // create new;
+        Consumer<byte[]> newConsumer = createConsumer(newTopics);
+        if (Objects.isNull(newConsumer)) {
+            currentConsumer.resume();
+            this.setStopConsume(false);
+            LOGGER.error("create new consumer failed, use the old one");
+            return false;
+        }
+        PulsarConsumer newConsumerWrapper = new PulsarConsumer(newConsumer);
+        InLongTopic firstTopic = newTopics.stream().findFirst().get();
+        final Seeker newSeeker = SeekerFactory.createPulsarSeeker(newConsumer, 
firstTopic);
+        // save
+        currentConsumer.setStopTime(System.currentTimeMillis());
+        toBeRemovedConsumers.add(currentConsumer);
+        // replace
+        this.currentConsumer = newConsumerWrapper;
+        this.seeker = newSeeker;
+        this.interceptor.configure(firstTopic);
+        this.onlineTopics = 
newTopics.stream().collect(Collectors.toMap(InLongTopic::getTopic, t -> t));
+        // resume
+        this.setStopConsume(false);
+        return true;
+    }
+
+    private Consumer<byte[]> createConsumer(Collection<InLongTopic> newTopics) 
{
+        if (CollectionUtils.isEmpty(newTopics)) {
+            LOGGER.error("new topic is empty or null");
+            return null;
+        }
+        try {
+            SubscriptionInitialPosition position = 
SubscriptionInitialPosition.Latest;
+            SortClientConfig.ConsumeStrategy offsetResetStrategy = 
context.getConfig().getOffsetResetStrategy();
+            if (offsetResetStrategy == 
SortClientConfig.ConsumeStrategy.earliest
+                    || offsetResetStrategy == 
SortClientConfig.ConsumeStrategy.earliest_absolutely) {
+                LOGGER.info("the subscription initial position is earliest!");
+                position = SubscriptionInitialPosition.Earliest;
+            }
+
+            List<String> topicNames = newTopics.stream()
+                    .map(InLongTopic::getTopic)
+                    .collect(Collectors.toList());
+            Consumer<byte[]> consumer = pulsarClient.newConsumer(Schema.BYTES)
+                    .topics(topicNames)
+                    .subscriptionName(context.getConfig().getSortTaskId())
+                    .subscriptionType(SubscriptionType.Shared)
+                    .startMessageIdInclusive()
+                    .subscriptionInitialPosition(position)
+                    .ackTimeout(context.getConfig().getAckTimeoutSec(), 
TimeUnit.SECONDS)
+                    
.receiverQueueSize(context.getConfig().getPulsarReceiveQueueSize())
+                    .subscribe();
+            LOGGER.info("create consumer for topics {}", topicNames);
+            return consumer;
+        } catch (Exception e) {
+            LOGGER.error("failed to create pulsar consumer", e);
+            return null;
+        }
+    }
+
+    @Override
+    public void ack(String msgOffset) throws Exception {
+        if (StringUtils.isBlank(msgOffset)) {
+            LOGGER.error("ack failed, msg offset should not be blank");
+            return;
+        }
+        if (Objects.isNull(currentConsumer)) {
+            LOGGER.error("ack failed, consumer is null");
+            return;
+        }
+        // if this ack belongs to current consumer
+        MessageId messageId = currentConsumer.getMessageId(msgOffset);
+        if (!Objects.isNull(messageId)) {
+            doAck(msgOffset, this.currentConsumer, messageId);
+            return;
+        }
+
+        // if this ack doesn't belong to current consumer, find in to be 
removed ones.
+        for (PulsarConsumer oldConsumer : toBeRemovedConsumers) {
+            MessageId id = oldConsumer.getMessageId(msgOffset);
+            if (Objects.isNull(id)) {
+                continue;
+            }
+            doAck(msgOffset, oldConsumer, id);
+            LOGGER.info("ack an old consumer message");
+            return;
+        }
+        context.getDefaultStateCounter().addAckFailTimes(1L);
+        LOGGER.error("in pulsar multi topic fetcher, messageId == null");
+    }
+
+    private void doAck(String msgOffset, PulsarConsumer consumer, MessageId 
messageId) {
+        if (!consumer.isConnected()) {
+            return;
+        }
+        InLongTopic topic = consumer.getTopic(msgOffset);
+        consumer.acknowledgeAsync(messageId)
+                .thenAccept(ctx -> ackSucc(msgOffset, topic, 
this.currentConsumer))
+                .exceptionally(exception -> {
+                    LOGGER.error("topic " + topic + " ack failed for offset " 
+ msgOffset + ", error: ", exception);
+                    context.getStateCounterByTopic(topic).addAckFailTimes(1L);
+                    return null;
+                });
+    }
+
+    private void ackSucc(String offset, InLongTopic topic, PulsarConsumer 
consumer) {
+        consumer.remove(offset);
+        context.getStateCounterByTopic(topic).addAckSuccTimes(1L);
+    }
+
+    @Override
+    public void pause() {
+        if (Objects.nonNull(currentConsumer)) {
+            currentConsumer.pause();
+        }
+    }
+
+    @Override
+    public void resume() {
+        if (Objects.nonNull(currentConsumer)) {
+            currentConsumer.resume();
+        }
+    }
+
+    @Override
+    public boolean close() {
+        mainLock.writeLock().lock();
+        try {
+            this.setStopConsume(true);
+            LOGGER.info("closed online topics {}", onlineTopics);
+            try {
+                if (currentConsumer != null) {
+                    currentConsumer.close();
+                }
+                if (fetchThread != null) {
+                    fetchThread.interrupt();
+                }
+            } catch (PulsarClientException e) {
+                LOGGER.warn("close pulsar client: ", e);
+            } catch (Throwable t) {
+                LOGGER.warn("got exception in close multi topic fetcher: ", t);
+            }
+            toBeRemovedConsumers.stream()
+                    .filter(Objects::nonNull)
+                    .forEach(c -> {
+                        try {
+                            c.close();
+                        } catch (PulsarClientException e) {
+                            LOGGER.warn("close pulsar client: ", e);
+                        }
+                    });
+            toBeRemovedConsumers.clear();
+            return true;
+        } finally {
+            this.closed = true;
+            mainLock.writeLock().unlock();
+        }
+    }
+
+    @Override
+    public boolean isClosed() {
+        return closed;
+    }
+
+    @Override
+    public void setStopConsume(boolean stopConsume) {
+        this.stopConsume = stopConsume;
+    }
+
+    @Override
+    public boolean isStopConsume() {
+        return stopConsume;
+    }
+
+    @Override
+    public List<InLongTopic> getTopics() {
+        return new ArrayList<>(onlineTopics.values());
+    }
+
+    @Override
+    public boolean updateTopics(List<InLongTopic> topics) {
+        if (needUpdate(topics)) {
+            return updateAll(topics);
+        }
+        LOGGER.info("no need to update multi topic fetcher");
+        return false;
+    }
+
+    private boolean needUpdate(Collection<InLongTopic> newTopics) {
+        if (newTopics.size() != onlineTopics.size()) {
+            return true;
+        }
+        // all topic should share the same properties in one task
+        if (Objects.equals(newTopics.stream().findFirst(), 
onlineTopics.values().stream().findFirst())) {
+            return true;
+        }
+        for (InLongTopic topic : newTopics) {
+            if (!onlineTopics.containsKey(topic.getTopic())) {
+                return true;
+            }
+        }
+        return false;
+    }
+
+    public class Fetcher implements Runnable {
+
+        /**
+         * put the received msg to onFinished method
+         *
+         * @param messageRecords {@link List}
+         */
+        private void handleAndCallbackMsg(List<MessageRecord> messageRecords) {
+            long start = System.currentTimeMillis();
+            try {
+                context.getDefaultStateCounter().addCallbackTimes(1L);
+                
context.getConfig().getCallback().onFinishedBatch(messageRecords);
+                context.getDefaultStateCounter()
+                        .addCallbackTimeCost(System.currentTimeMillis() - 
start).addCallbackDoneTimes(1L);
+            } catch (Exception e) {
+                context.getDefaultStateCounter().addCallbackErrorTimes(1L);
+                LOGGER.error("failed to handle callback: ", e);
+            }
+        }
+
+        private String getOffset(MessageId msgId) {
+            return Base64.getEncoder().encodeToString(msgId.toByteArray());
+        }
+
+        private List<MessageRecord> processPulsarMsg(Messages<byte[]> 
messages) throws Exception {
+            List<MessageRecord> msgs = new ArrayList<>();
+            for (Message<byte[]> msg : messages) {
+                String topicName = msg.getTopicName();
+                InLongTopic topic = onlineTopics.get(topicName);
+                if (Objects.isNull(topic)) {
+                    LOGGER.error("got a message with topic {}, which is not 
subscribe", topicName);
+                    continue;
+                }
+                // if need seek
+                if (msg.getPublishTime() < seeker.getSeekTime()) {
+                    seeker.seek();
+                    break;
+                }
+                String offsetKey = getOffset(msg.getMessageId());
+                currentConsumer.put(offsetKey, topic, msg.getMessageId());
+
+                //deserialize
+                List<InLongMessage> inLongMessages = deserializer
+                        .deserialize(context, topic, msg.getProperties(), 
msg.getData());
+                // intercept
+                inLongMessages = interceptor.intercept(inLongMessages);
+                if (inLongMessages.isEmpty()) {
+                    ack(offsetKey);
+                    continue;
+                }
+
+                msgs.add(new MessageRecord(topic.getTopicKey(),
+                        inLongMessages,
+                        offsetKey, System.currentTimeMillis()));
+                
context.getStateCounterByTopic(topic).addConsumeSize(msg.getData().length);
+            }
+            context.getDefaultStateCounter().addMsgCount(msgs.size());
+            return msgs;
+        }
+
+        @Override
+        public void run() {
+            boolean hasPermit;
+            while (true) {
+                hasPermit = false;
+                try {
+                    if (context.getConfig().isStopConsume() || stopConsume) {
+                        TimeUnit.MILLISECONDS.sleep(50);
+                        continue;
+                    }
+
+                    if (sleepTime > 0) {
+                        TimeUnit.MILLISECONDS.sleep(sleepTime);
+                    }
+
+                    context.acquireRequestPermit();
+                    hasPermit = true;
+                    
context.getDefaultStateCounter().addMsgCount(1L).addFetchTimes(1L);
+
+                    long startFetchTime = System.currentTimeMillis();
+                    Messages<byte[]> pulsarMessages = 
currentConsumer.batchReceive();
+
+                    
context.getDefaultStateCounter().addFetchTimeCost(System.currentTimeMillis() - 
startFetchTime);
+                    if (null != pulsarMessages && pulsarMessages.size() != 0) {
+                        List<MessageRecord> msgs = 
this.processPulsarMsg(pulsarMessages);
+                        handleAndCallbackMsg(msgs);
+                        sleepTime = 0L;
+                    } else {
+                        
context.getDefaultStateCounter().addEmptyFetchTimes(1L);
+                        emptyFetchTimes++;
+                        if (emptyFetchTimes >= 
context.getConfig().getEmptyPollTimes()) {
+                            sleepTime = Math.min((sleepTime += 
context.getConfig().getEmptyPollSleepStepMs()),
+                                    
context.getConfig().getMaxEmptyPollSleepMs());
+                            emptyFetchTimes = 0;
+                        }
+                    }
+                } catch (Exception e) {
+                    context.getDefaultStateCounter().addFetchErrorTimes(1L);
+                    LOGGER.error("failed to fetch msg ", e);
+                } finally {
+                    if (hasPermit) {
+                        context.releaseRequestPermit();
+                    }
+                }
+
+                if (closed) {
+                    break;
+                }
+            }
+        }
+    }
+}

Reply via email to