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 541337579d..6b1c051eeb 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 @@ -81,14 +81,24 @@ public class TbKafkaConsumerTemplate implements TbQueueCon @Override public void subscribe() { - partitions = Collections.singleton(new TopicPartitionInfo(topic, null, null, true)); - subscribed = false; + consumerLock.lock(); + try { + partitions = Collections.singleton(new TopicPartitionInfo(topic, null, null, true)); + subscribed = false; + } finally { + consumerLock.unlock(); + } } @Override public void subscribe(Set partitions) { - this.partitions = partitions; - subscribed = false; + consumerLock.lock(); + try { + this.partitions = partitions; + subscribed = false; + } finally { + consumerLock.unlock(); + } } @Override @@ -100,13 +110,11 @@ public class TbKafkaConsumerTemplate implements TbQueueCon log.debug("Failed to await subscription", e); } } else { + consumerLock.lock(); try { - consumerLock.lock(); - if (!subscribed) { List topicNames = partitions.stream().map(TopicPartitionInfo::getFullTopicName).collect(Collectors.toList()); topicNames.forEach(admin::createTopicIfNotExists); - consumer.unsubscribe(); consumer.subscribe(topicNames); subscribed = true; } @@ -132,8 +140,8 @@ public class TbKafkaConsumerTemplate implements TbQueueCon @Override public void commit() { + consumerLock.lock(); try { - consumerLock.lock(); consumer.commitAsync(); } finally { consumerLock.unlock(); @@ -142,8 +150,8 @@ public class TbKafkaConsumerTemplate implements TbQueueCon @Override public void unsubscribe() { + consumerLock.lock(); try { - consumerLock.lock(); if (consumer != null) { consumer.unsubscribe(); consumer.close(); diff --git a/docker/docker-compose.yml b/docker/docker-compose.yml index 2f0d0fb3b3..9061f3e2de 100644 --- a/docker/docker-compose.yml +++ b/docker/docker-compose.yml @@ -28,7 +28,7 @@ services: ZOO_SERVERS: server.1=zookeeper:2888:3888;zookeeper:2181 kafka: restart: always - image: "wurstmeister/kafka:2.12-2.2.1" + image: "wurstmeister/kafka:2.12-2.3.0" ports: - "9092:9092" env_file: