From ec446d4da4a843c5c4481eb4d78401006e15dc42 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Mon, 28 Oct 2019 13:34:51 -0400 Subject: [PATCH] GH-1283: Unique client.id for each producer Resolves https://github.com/spring-projects/spring-kafka/issues/1283 Avoid `InstanceAlreadyExistsException` s when the user supplies a custom `client.id`. **cherry-pick to all supported branches** --- .../core/DefaultKafkaProducerFactory.java | 23 ++++++++++++++++++- .../core/KafkaTemplateTransactionTests.java | 4 +++- 2 files changed, 25 insertions(+), 2 deletions(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java index 7657ad4c..5b9cdd3e 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java @@ -123,6 +123,8 @@ public class DefaultKafkaProducerFactory implements ProducerFactory, private final Map> consumerProducers = new HashMap<>(); + private final AtomicInteger clientIdCounter = new AtomicInteger(); + private Supplier> keySerializerSupplier; private Supplier> valueSerializerSupplier; @@ -139,6 +141,8 @@ public class DefaultKafkaProducerFactory implements ProducerFactory, private ThreadLocal> threadBoundProducers; + private String clientIdPrefix; + private volatile CloseSafeProducer producer; /** @@ -182,6 +186,9 @@ public class DefaultKafkaProducerFactory implements ProducerFactory, this.configs = new HashMap<>(configs); this.keySerializerSupplier = keySerializerSupplier == null ? () -> null : keySerializerSupplier; this.valueSerializerSupplier = valueSerializerSupplier == null ? () -> null : valueSerializerSupplier; + if (this.clientIdPrefix == null && configs.get(ProducerConfig.CLIENT_ID_CONFIG) instanceof String) { + this.clientIdPrefix = (String) configs.get(ProducerConfig.CLIENT_ID_CONFIG); + } String txId = (String) this.configs.get(ProducerConfig.TRANSACTIONAL_ID_CONFIG); if (StringUtils.hasText(txId)) { @@ -393,7 +400,17 @@ public class DefaultKafkaProducerFactory implements ProducerFactory, * @return the producer. */ protected Producer createKafkaProducer() { - return new KafkaProducer<>(this.configs, this.keySerializerSupplier.get(), this.valueSerializerSupplier.get()); + if (this.clientIdPrefix == null) { + return new KafkaProducer<>(this.configs, this.keySerializerSupplier.get(), + this.valueSerializerSupplier.get()); + } + else { + Map newConfigs = new HashMap<>(this.configs); + newConfigs.put(ProducerConfig.CLIENT_ID_CONFIG, + this.clientIdPrefix + "-" + this.clientIdCounter.incrementAndGet()); + return new KafkaProducer<>(newConfigs, this.keySerializerSupplier.get(), + this.valueSerializerSupplier.get()); + } } protected Producer createTransactionalProducerForPartition() { @@ -460,6 +477,10 @@ public class DefaultKafkaProducerFactory implements ProducerFactory, Producer newProducer; Map newProducerConfigs = new HashMap<>(this.configs); newProducerConfigs.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, prefix + suffix); + if (this.clientIdPrefix != null) { + newProducerConfigs.put(ProducerConfig.CLIENT_ID_CONFIG, + this.clientIdPrefix + "-" + this.clientIdCounter.incrementAndGet()); + } newProducer = new KafkaProducer<>(newProducerConfigs, this.keySerializerSupplier .get(), this.valueSerializerSupplier.get()); newProducer.initTransactions(); diff --git a/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTransactionTests.java b/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTransactionTests.java index aa2cc78c..8cd41795 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTransactionTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTransactionTests.java @@ -102,6 +102,7 @@ public class KafkaTemplateTransactionTests { Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); senderProps.put(ProducerConfig.RETRIES_CONFIG, 1); senderProps.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "my.transaction."); + senderProps.put(ProducerConfig.CLIENT_ID_CONFIG, "customClientId"); DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>(senderProps); pf.setKeySerializer(new StringSerializer()); KafkaTemplate template = new KafkaTemplate<>(pf); @@ -116,6 +117,7 @@ public class KafkaTemplateTransactionTests { template.executeInTransaction(kt -> kt.send(LOCAL_TX_IN_TOPIC, "one")); ConsumerRecord singleRecord = KafkaTestUtils.getSingleRecord(consumer, LOCAL_TX_IN_TOPIC); template.executeInTransaction(t -> { + pf.createProducer("testCustomClientIdIsUnique").close(); t.sendDefault("foo", "bar"); t.sendDefault("baz", "qux"); t.sendOffsetsToTransaction(Collections.singletonMap( @@ -144,7 +146,7 @@ public class KafkaTemplateTransactionTests { template.executeInTransaction(t -> { assertThat(KafkaTestUtils.getPropertyValue( KafkaTestUtils.getPropertyValue(template, "producers", ThreadLocal.class).get(), - "delegate.transactionManager.transactionalId")).isEqualTo("tx.template.override.1"); + "delegate.transactionManager.transactionalId")).isEqualTo("tx.template.override.2"); return null; }); assertThat(pf.getCache("tx.template.override.")).hasSize(1);