Skip to content

Commit da4743d

Browse files
committed
NIFI-16097 Added enable.idemotence Option in Publish Kafka Processor
1 parent ec1f27f commit da4743d

4 files changed

Lines changed: 31 additions & 7 deletions

File tree

nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-processors/src/main/java/org/apache/nifi/kafka/processors/PublishKafka.java

Lines changed: 18 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -140,6 +140,19 @@ public class PublishKafka extends AbstractProcessor implements VerifiableProcess
140140
.defaultValue(DeliveryGuarantee.DELIVERY_REPLICATED)
141141
.build();
142142

143+
static final PropertyDescriptor ENABLE_IDEMPOTENCE = new PropertyDescriptor.Builder()
144+
.name("enable.idempotence")
145+
.displayName("Enable Idempotence")
146+
.description("Specifies whether the producer will ensure that exactly one copy of each message is written in the stream. " +
147+
"If set to ‘false’, producer retries due to broker failures, etc., may write duplicates of the retried message in the stream. " +
148+
"Corresponds to Kafka Client enable.idempotence property.")
149+
.expressionLanguageSupported(ExpressionLanguageScope.NONE)
150+
.dependsOn(DELIVERY_GUARANTEE, DeliveryGuarantee.DELIVERY_REPLICATED)
151+
.required(false)
152+
.allowableValues("true", "false")
153+
.defaultValue("true")
154+
.build();
155+
143156
static final PropertyDescriptor COMPRESSION_CODEC = new PropertyDescriptor.Builder()
144157
.name("compression.type")
145158
.displayName("Compression Type")
@@ -312,6 +325,7 @@ public class PublishKafka extends AbstractProcessor implements VerifiableProcess
312325
TOPIC_NAME,
313326
FAILURE_STRATEGY,
314327
DELIVERY_GUARANTEE,
328+
ENABLE_IDEMPOTENCE,
315329
COMPRESSION_CODEC,
316330
MAX_REQUEST_SIZE,
317331
TRANSACTIONS_ENABLED,
@@ -361,11 +375,12 @@ public List<ConfigVerificationResult> verify(final ProcessContext context, final
361375
final String transactionalIdPrefix = transactionsEnabled ? context.getProperty(TRANSACTIONAL_ID_PREFIX).evaluateAttributeExpressions().getValue() : null;
362376
final Supplier<String> transactionalIdSupplier = new TransactionIdSupplier(transactionalIdPrefix);
363377
final String deliveryGuarantee = context.getProperty(DELIVERY_GUARANTEE).getValue();
378+
final boolean enableIdempotence = context.getProperty(ENABLE_IDEMPOTENCE).asBoolean();
364379
final String compressionCodec = context.getProperty(COMPRESSION_CODEC).getValue();
365380
final String partitionClass = context.getProperty(PARTITION_CLASS).getValue();
366381
final int maxRequestSize = context.getProperty(MAX_REQUEST_SIZE).asDataSize(DataUnit.B).intValue();
367382
final ProducerConfiguration producerConfiguration = new ProducerConfiguration(
368-
transactionsEnabled, transactionalIdSupplier.get(), deliveryGuarantee, compressionCodec, partitionClass, maxRequestSize);
383+
transactionsEnabled, transactionalIdSupplier.get(), deliveryGuarantee, enableIdempotence, compressionCodec, partitionClass, maxRequestSize);
369384

370385
try (final KafkaProducerService producerService = connectionService.getProducerService(producerConfiguration)) {
371386
final ConfigVerificationResult.Builder verificationPartitions = new ConfigVerificationResult.Builder()
@@ -454,12 +469,13 @@ private KafkaProducerService createProducerService(final ProcessContext context)
454469
final boolean transactionsEnabled = context.getProperty(TRANSACTIONS_ENABLED).asBoolean();
455470
final String transactionalIdPrefix = transactionsEnabled ? context.getProperty(TRANSACTIONAL_ID_PREFIX).evaluateAttributeExpressions().getValue() : null;
456471
final String deliveryGuarantee = context.getProperty(DELIVERY_GUARANTEE).getValue();
472+
final boolean enableIdempotence = context.getProperty(ENABLE_IDEMPOTENCE).asBoolean();
457473
final String compressionCodec = context.getProperty(COMPRESSION_CODEC).getValue();
458474
final String partitionClass = context.getProperty(PARTITION_CLASS).getValue();
459475
final int maxRequestSize = context.getProperty(MAX_REQUEST_SIZE).asDataSize(DataUnit.B).intValue();
460476

461477
final ProducerConfiguration producerConfiguration = new ProducerConfiguration(
462-
transactionsEnabled, transactionalIdPrefix, deliveryGuarantee, compressionCodec, partitionClass, maxRequestSize);
478+
transactionsEnabled, transactionalIdPrefix, deliveryGuarantee, enableIdempotence, compressionCodec, partitionClass, maxRequestSize);
463479

464480
return connectionService.getProducerService(producerConfiguration);
465481
}

nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-service-api/src/main/java/org/apache/nifi/kafka/service/api/producer/ProducerConfiguration.java

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,19 +20,22 @@ public class ProducerConfiguration {
2020
private final boolean transactionsEnabled;
2121
private final String transactionIdPrefix;
2222
private final String deliveryGuarantee;
23+
private final boolean enableIdempotence;
2324
private final String compressionCodec;
2425
private final String partitionClass;
2526
private final int maxRequestSize;
2627

2728
public ProducerConfiguration(final boolean transactionsEnabled,
2829
final String transactionIdPrefix,
2930
final String deliveryGuarantee,
31+
final boolean enableIdempotence,
3032
final String compressionCodec,
3133
final String partitionClass,
3234
final int maxRequestSize) {
3335
this.transactionsEnabled = transactionsEnabled;
3436
this.transactionIdPrefix = transactionIdPrefix;
3537
this.deliveryGuarantee = deliveryGuarantee;
38+
this.enableIdempotence = enableIdempotence;
3639
this.compressionCodec = compressionCodec;
3740
this.partitionClass = partitionClass;
3841
this.maxRequestSize = maxRequestSize;
@@ -50,6 +53,10 @@ public String getDeliveryGuarantee() {
5053
return deliveryGuarantee;
5154
}
5255

56+
public boolean getEnableIdempotence() {
57+
return enableIdempotence;
58+
}
59+
5360
public String getCompressionCodec() {
5461
return compressionCodec;
5562
}

nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-service-shared/src/main/java/org/apache/nifi/kafka/service/Kafka3ConnectionService.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -281,6 +281,7 @@ public KafkaProducerService getProducerService(final ProducerConfiguration produ
281281
if (producerConfiguration.getDeliveryGuarantee() != null) {
282282
properties.put(ProducerConfig.ACKS_CONFIG, producerConfiguration.getDeliveryGuarantee());
283283
}
284+
properties.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, producerConfiguration.getEnableIdempotence());
284285
if (producerConfiguration.getCompressionCodec() != null) {
285286
properties.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, producerConfiguration.getCompressionCodec());
286287
}

nifi-extension-bundles/nifi-kafka-bundle/nifi-kafka-service-shared/src/test/java/org/apache/nifi/kafka/service/Kafka3ConnectionServiceBaseIT.java

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -241,7 +241,7 @@ void testAdminClient() throws ExecutionException, InterruptedException, TimeoutE
241241

242242
@Test
243243
void testProduceOneNoTransaction() {
244-
final ProducerConfiguration producerConfiguration = new ProducerConfiguration(false, null, null, null, null, 1_000_000);
244+
final ProducerConfiguration producerConfiguration = new ProducerConfiguration(false, null, null, false, null, null, 1_000_000);
245245
final KafkaProducerService producerService = service.getProducerService(producerConfiguration);
246246
final KafkaRecord kafkaRecord = new KafkaRecord(null, null, null, null, RECORD_VALUE, emptyList());
247247
final List<KafkaRecord> kafkaRecords = Collections.singletonList(kafkaRecord);
@@ -252,7 +252,7 @@ void testProduceOneNoTransaction() {
252252

253253
@Test
254254
void testProduceOneWithTransaction() {
255-
final ProducerConfiguration producerConfiguration = new ProducerConfiguration(true, "transaction-", null, null, null, 1_000_000);
255+
final ProducerConfiguration producerConfiguration = new ProducerConfiguration(true, "transaction-", null, false, null, null, 1_000_000);
256256
final KafkaProducerService producerService = service.getProducerService(producerConfiguration);
257257
final KafkaRecord kafkaRecord = new KafkaRecord(null, null, null, null, RECORD_VALUE, emptyList());
258258
final List<KafkaRecord> kafkaRecords = Collections.singletonList(kafkaRecord);
@@ -263,7 +263,7 @@ void testProduceOneWithTransaction() {
263263

264264
@Test
265265
void testProduceConsumeRecord() {
266-
final ProducerConfiguration producerConfiguration = new ProducerConfiguration(false, null, null, null, null, 1_000_000);
266+
final ProducerConfiguration producerConfiguration = new ProducerConfiguration(false, null, null, false, null, null, 1_000_000);
267267
final KafkaProducerService producerService = service.getProducerService(producerConfiguration);
268268

269269
final long timestamp = System.currentTimeMillis();
@@ -300,7 +300,7 @@ void testCurrentLag() {
300300
final String groupId = "Group_" + timestamp;
301301
final String topic = "Topic_" + timestamp;
302302
final int partition = 0;
303-
final ProducerConfiguration producerConfiguration = new ProducerConfiguration(false, null, null, null, null, 1_000_000);
303+
final ProducerConfiguration producerConfiguration = new ProducerConfiguration(false, null, null, false, null, null, 1_000_000);
304304
final KafkaProducerService producerService = service.getProducerService(producerConfiguration);
305305

306306
final KafkaRecord kafkaRecord = new KafkaRecord(null, partition, timestamp, RECORD_KEY, RECORD_VALUE, emptyList());
@@ -404,7 +404,7 @@ private void assertResultFound(
404404

405405
@Test
406406
void testGetProducerService() {
407-
final ProducerConfiguration producerConfiguration = new ProducerConfiguration(false, null, null, null, null, 1_000_000);
407+
final ProducerConfiguration producerConfiguration = new ProducerConfiguration(false, null, null, false, null, null, 1_000_000);
408408
final KafkaProducerService producerService = service.getProducerService(producerConfiguration);
409409
final List<PartitionState> partitionStates = producerService.getPartitionStates(TOPIC);
410410
assertPartitionStatesFound(partitionStates);

0 commit comments

Comments
 (0)