diff --git a/libraries/pom.xml b/libraries/pom.xml
index 7402d88ef3..83d78af84f 100644
--- a/libraries/pom.xml
+++ b/libraries/pom.xml
@@ -795,8 +795,8 @@
http://dl.bintray.com/cuba-platform/main
- Apache Staging
- https://repository.apache.org/content/groups/staging
+ Maven Central
+ https://repo.maven.apache.org/maven2/
@@ -912,6 +912,41 @@
+
+
+
+
+ org.eclipse.m2e
+ lifecycle-mapping
+ 1.0.0
+
+
+
+
+
+
+ org.apache.maven.plugins
+
+
+ maven-pmd-plugin
+
+
+ [3.8,)
+
+
+ check
+
+
+
+
+
+
+
+
+
+
+
+
@@ -974,7 +1009,7 @@
2.5.5
1.23.0
v4-rev493-1.21.0
- 1.0.0
+ 2.0.0
1.7.0
3.0.14
8.5.24
diff --git a/libraries/src/main/java/com/baeldung/kafka/TransactionalApp.java b/libraries/src/main/java/com/baeldung/kafka/TransactionalApp.java
new file mode 100644
index 0000000000..1e95041a0d
--- /dev/null
+++ b/libraries/src/main/java/com/baeldung/kafka/TransactionalApp.java
@@ -0,0 +1,102 @@
+package com.baeldung.kafka;
+
+import org.apache.kafka.clients.consumer.ConsumerRecord;
+import org.apache.kafka.clients.consumer.ConsumerRecords;
+import org.apache.kafka.clients.consumer.KafkaConsumer;
+import org.apache.kafka.clients.consumer.OffsetAndMetadata;
+import org.apache.kafka.clients.producer.KafkaProducer;
+import org.apache.kafka.clients.producer.ProducerConfig;
+import org.apache.kafka.clients.producer.ProducerRecord;
+import org.apache.kafka.common.KafkaException;
+import org.apache.kafka.common.TopicPartition;
+
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Properties;
+
+import static java.time.Duration.ofSeconds;
+import static java.util.Collections.singleton;
+import static org.apache.kafka.clients.consumer.ConsumerConfig.*;
+import static org.apache.kafka.clients.consumer.ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG;
+import static org.apache.kafka.clients.producer.ProducerConfig.*;
+
+public class TransactionalApp {
+
+ private static final String CONSUMER_GROUP_ID = "test";
+ private static final String OUTPUT_TOPIC = "output";
+ private static final String INPUT_TOPIC = "input";
+
+ public static void main(String[] args) {
+
+ KafkaConsumer consumer = initConsumer();
+ KafkaProducer producer = initProducer();
+
+ producer.initTransactions();
+
+ try {
+
+ while (true) {
+
+ ConsumerRecords records = consumer.poll(ofSeconds(20));
+
+ producer.beginTransaction();
+
+ for (ConsumerRecord record : records)
+ producer.send(new ProducerRecord(OUTPUT_TOPIC, record));
+
+ Map offsetsToCommit = new HashMap<>();
+
+ for (TopicPartition partition : records.partitions()) {
+ List> partitionedRecords = records.records(partition);
+ long offset = partitionedRecords.get(partitionedRecords.size() - 1).offset();
+
+ offsetsToCommit.put(partition, new OffsetAndMetadata(offset));
+ }
+
+ producer.sendOffsetsToTransaction(offsetsToCommit, CONSUMER_GROUP_ID);
+ producer.commitTransaction();
+
+ }
+
+ } catch (KafkaException e) {
+
+ producer.abortTransaction();
+
+ }
+
+
+ }
+
+ private static KafkaConsumer initConsumer() {
+ Properties props = new Properties();
+ props.put(BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
+ props.put(GROUP_ID_CONFIG, CONSUMER_GROUP_ID);
+ props.put(ENABLE_AUTO_COMMIT_CONFIG, "false");
+ props.put(KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
+ props.put(VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
+
+ KafkaConsumer consumer = new KafkaConsumer<>(props);
+ consumer.subscribe(singleton(INPUT_TOPIC));
+ return consumer;
+ }
+
+ private static KafkaProducer initProducer() {
+
+ Properties props = new Properties();
+ props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
+ props.put(ACKS_CONFIG, "all");
+ props.put(RETRIES_CONFIG, 3);
+ props.put(BATCH_SIZE_CONFIG, 16384);
+ props.put(LINGER_MS_CONFIG, 1);
+ props.put(BUFFER_MEMORY_CONFIG, 33554432);
+ props.put(ENABLE_IDEMPOTENCE_CONFIG, "true");
+ props.put(TRANSACTIONAL_ID_CONFIG, "prod-1");
+ props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
+ props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
+
+ return new KafkaProducer(props);
+
+ }
+
+}