From b39328c989a36711c133de67d753cbdedc2a4467 Mon Sep 17 00:00:00 2001 From: Yevhen Bondarenko <56396344+YevhenBondarenko@users.noreply.github.com> Date: Tue, 31 Mar 2020 16:47:43 +0300 Subject: [PATCH] [2.5]fix ConcurrentModificationException (#2560) * fix ConcurrentModificationException * kafka consumer improvements * kafka consumer improvements * refactored kafka consumer * refactored kafka consumer --- .../queue/kafka/TBKafkaConsumerTemplate.java | 62 +++++++++++++------ 1 file changed, 43 insertions(+), 19 deletions(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TBKafkaConsumerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TBKafkaConsumerTemplate.java index 68b9cea7fa..d882153177 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TBKafkaConsumerTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TBKafkaConsumerTemplate.java @@ -33,6 +33,8 @@ import java.util.Collections; import java.util.List; import java.util.Properties; import java.util.Set; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; import java.util.stream.Collectors; /** @@ -46,6 +48,7 @@ public class TBKafkaConsumerTemplate implements TbQueueCon private final TbKafkaDecoder decoder; private volatile boolean subscribed; private volatile Set partitions; + private final Lock consumerLock; @Getter private final String topic; @@ -71,6 +74,7 @@ public class TBKafkaConsumerTemplate implements TbQueueCon this.consumer = new KafkaConsumer<>(props); this.decoder = decoder; this.topic = topic; + this.consumerLock = new ReentrantLock(); } @Override @@ -94,23 +98,30 @@ public class TBKafkaConsumerTemplate implements TbQueueCon log.debug("Failed to await subscription", e); } } else { - if (!subscribed) { - List topicNames = partitions.stream().map(TopicPartitionInfo::getFullTopicName).collect(Collectors.toList()); - topicNames.forEach(admin::createTopicIfNotExists); - consumer.subscribe(topicNames); - subscribed = true; - } - ConsumerRecords records = consumer.poll(Duration.ofMillis(durationInMillis)); - if (records.count() > 0) { - List result = new ArrayList<>(); - records.forEach(record -> { - try { - result.add(decode(record)); - } catch (IOException e) { - log.error("Failed decode record: [{}]", record); - } - }); - return result; + try { + consumerLock.lock(); + + if (!subscribed) { + List topicNames = partitions.stream().map(TopicPartitionInfo::getFullTopicName).collect(Collectors.toList()); + topicNames.forEach(admin::createTopicIfNotExists); + consumer.subscribe(topicNames); + subscribed = true; + } + + ConsumerRecords records = consumer.poll(Duration.ofMillis(durationInMillis)); + if (records.count() > 0) { + List result = new ArrayList<>(); + records.forEach(record -> { + try { + result.add(decode(record)); + } catch (IOException e) { + log.error("Failed decode record: [{}]", record); + } + }); + return result; + } + } finally { + consumerLock.unlock(); } } return Collections.emptyList(); @@ -118,12 +129,25 @@ public class TBKafkaConsumerTemplate implements TbQueueCon @Override public void commit() { - consumer.commitAsync(); + try { + consumerLock.lock(); + consumer.commitAsync(); + } finally { + consumerLock.unlock(); + } } @Override public void unsubscribe() { - consumer.unsubscribe(); + try { + consumerLock.lock(); + if (consumer != null) { + consumer.unsubscribe(); + consumer.close(); + } + } finally { + consumerLock.unlock(); + } } public T decode(ConsumerRecord record) throws IOException {