From 36fdb2d31f99751bbee0c39a18bce7796766e3b6 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** # Conflicts: # spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java # spring-kafka/src/test/java/org/springframework/kafka/core/KafkaTemplateTransactionTests.java # Conflicts: # spring-kafka/src/main/java/org/springframework/kafka/core/DefaultKafkaProducerFactory.java --- .../core/DefaultKafkaProducerFactory.java | 23 +++++++++++++++++-- .../core/KafkaTemplateTransactionTests.java | 1 + 2 files changed, 22 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 9a6047e6..81b684e9 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 @@ -1,5 +1,5 @@ /* - * Copyright 2016-2018 the original author or authors. + * Copyright 2016-2019 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -93,6 +93,8 @@ public class DefaultKafkaProducerFactory implements ProducerFactory, private final Map> consumerProducers = new HashMap<>(); + private final AtomicInteger clientIdCounter = new AtomicInteger(); + private volatile CloseSafeProducer producer; private Serializer keySerializer; @@ -107,6 +109,7 @@ public class DefaultKafkaProducerFactory implements ProducerFactory, private boolean producerPerConsumerPartition = true; + private String clientIdPrefix; /** * Construct a factory with the provided configuration. * @param configs the configuration. @@ -117,9 +120,13 @@ public class DefaultKafkaProducerFactory implements ProducerFactory, public DefaultKafkaProducerFactory(Map configs, Serializer keySerializer, Serializer valueSerializer) { + this.configs = new HashMap<>(configs); this.keySerializer = keySerializer; this.valueSerializer = valueSerializer; + if (configs.get(ProducerConfig.CLIENT_ID_CONFIG) instanceof String) { + this.clientIdPrefix = (String) configs.get(ProducerConfig.CLIENT_ID_CONFIG); + } } public void setKeySerializer(Serializer keySerializer) { @@ -274,7 +281,15 @@ public class DefaultKafkaProducerFactory implements ProducerFactory, * @return the producer. */ protected Producer createKafkaProducer() { - return new KafkaProducer(this.configs, this.keySerializer, this.valueSerializer); + if (this.clientIdPrefix == null) { + return new KafkaProducer(this.configs, this.keySerializer, this.valueSerializer); + } + else { + Map newConfigs = new HashMap<>(this.configs); + newConfigs.put(ProducerConfig.CLIENT_ID_CONFIG, + this.clientIdPrefix + "-" + this.clientIdCounter.incrementAndGet()); + return new KafkaProducer<>(newConfigs, this.keySerializer, this.valueSerializer); + } } Producer createTransactionalProducerForPartition() { @@ -328,6 +343,10 @@ public class DefaultKafkaProducerFactory implements ProducerFactory, Producer newProducer; Map newProducerConfigs = new HashMap<>(this.configs); newProducerConfigs.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, this.transactionIdPrefix + suffix); + if (this.clientIdPrefix != null) { + newProducerConfigs.put(ProducerConfig.CLIENT_ID_CONFIG, + this.clientIdPrefix + "-" + this.clientIdCounter.incrementAndGet()); + } newProducer = new KafkaProducer(newProducerConfigs, this.keySerializer, this.valueSerializer); newProducer.initTransactions(); return new CloseSafeProducer(newProducer, this.cache, remover, 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 635e0fa8..c9aedc27 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 @@ -99,6 +99,7 @@ public class KafkaTemplateTransactionTests { public void testLocalTransaction() throws Exception { Map senderProps = KafkaTestUtils.producerProps(embeddedKafka); senderProps.put(ProducerConfig.RETRIES_CONFIG, 1); + senderProps.put(ProducerConfig.CLIENT_ID_CONFIG, "customClientId"); DefaultKafkaProducerFactory pf = new DefaultKafkaProducerFactory<>(senderProps); pf.setKeySerializer(new StringSerializer()); pf.setTransactionIdPrefix("my.transaction.");