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**
This commit is contained in:
committed by
Artem Bilan
parent
ae0dec7720
commit
ec446d4da4
@@ -123,6 +123,8 @@ public class DefaultKafkaProducerFactory<K, V> implements ProducerFactory<K, V>,
|
||||
|
||||
private final Map<String, CloseSafeProducer<K, V>> consumerProducers = new HashMap<>();
|
||||
|
||||
private final AtomicInteger clientIdCounter = new AtomicInteger();
|
||||
|
||||
private Supplier<Serializer<K>> keySerializerSupplier;
|
||||
|
||||
private Supplier<Serializer<V>> valueSerializerSupplier;
|
||||
@@ -139,6 +141,8 @@ public class DefaultKafkaProducerFactory<K, V> implements ProducerFactory<K, V>,
|
||||
|
||||
private ThreadLocal<CloseSafeProducer<K, V>> threadBoundProducers;
|
||||
|
||||
private String clientIdPrefix;
|
||||
|
||||
private volatile CloseSafeProducer<K, V> producer;
|
||||
|
||||
/**
|
||||
@@ -182,6 +186,9 @@ public class DefaultKafkaProducerFactory<K, V> implements ProducerFactory<K, V>,
|
||||
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<K, V> implements ProducerFactory<K, V>,
|
||||
* @return the producer.
|
||||
*/
|
||||
protected Producer<K, V> 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<String, Object> 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<K, V> createTransactionalProducerForPartition() {
|
||||
@@ -460,6 +477,10 @@ public class DefaultKafkaProducerFactory<K, V> implements ProducerFactory<K, V>,
|
||||
Producer<K, V> newProducer;
|
||||
Map<String, Object> 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();
|
||||
|
||||
@@ -102,6 +102,7 @@ public class KafkaTemplateTransactionTests {
|
||||
Map<String, Object> 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<String, String> pf = new DefaultKafkaProducerFactory<>(senderProps);
|
||||
pf.setKeySerializer(new StringSerializer());
|
||||
KafkaTemplate<String, String> template = new KafkaTemplate<>(pf);
|
||||
@@ -116,6 +117,7 @@ public class KafkaTemplateTransactionTests {
|
||||
template.executeInTransaction(kt -> kt.send(LOCAL_TX_IN_TOPIC, "one"));
|
||||
ConsumerRecord<String, String> 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);
|
||||
|
||||
Reference in New Issue
Block a user