8.4 How to use transactions
(This is the practical payoff of KIP-447's group-metadata fencing.)
4.1 The recommended way — don't use them directly
"The most common and most recommended way to use transactions is to enable exactly-once guarantees in Kafka Streams. This way, we will not use transactions directly at all, but rather Kafka Streams will use them for us behind the scenes... Transactions were designed with this use case in mind, so using them via Kafka Streams is the easiest and most likely to work as expected."
processing.guarantee=exactly_once
# or
processing.guarantee=exactly_once_beta"That's it."
NOTE —
exactly_once_beta"a slightly different method of handling application instances that crash or hang with in-flight transactions. Introduced in release 2.5 to Kafka brokers, and release 2.6 to Kafka Streams. The main benefit of this method is the ability to handle MANY PARTITIONS WITH A SINGLE TRANSACTIONAL PRODUCER and therefore create MORE SCALABLE Kafka Streams applications."
(This is the practical payoff of KIP-447's group-metadata fencing.)
4.2 Using the transactional API directly
Properties producerProps = new Properties();
producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
producerProps.put(ProducerConfig.CLIENT_ID_CONFIG, "DemoProducer");
producerProps.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, transactionalId); // ①
producer = new KafkaProducer<>(producerProps);
Properties consumerProps = new Properties();
consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); // ②
consumerProps.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed"); // ③
consumer = new KafkaConsumer<>(consumerProps);
producer.initTransactions(); // ④
consumer.subscribe(Collections.singleton(inputTopic)); // ⑤
while (true) {
try {
ConsumerRecords<Integer, String> records =
consumer.poll(Duration.ofMillis(200));
if (records.count() > 0) {
producer.beginTransaction(); // ⑥
for (ConsumerRecord<Integer, String> record : records) {
ProducerRecord<Integer, String> customizedRecord = transform(record); // ⑦
producer.send(customizedRecord);
}
Map<TopicPartition, OffsetAndMetadata> offsets = consumerOffsets();
producer.sendOffsetsToTransaction(offsets, consumer.groupMetadata()); // ⑧
producer.commitTransaction(); // ⑨
}
} catch (ProducerFencedException|InvalidProducerEpochException e) { // ⑩
throw new KafkaException(String.format(
"The transactional.id %s is used by another process", transactionalId));
} catch (KafkaException e) { // ⑪
producer.abortTransaction();
resetToLastCommittedPositions(consumer);
}}Annotations:
① "Configuring a producer with transactional.id makes it a transactional producer capable of producing atomic multipartition writes. The transactional ID must be UNIQUE and LONG-LIVED. Essentially it defines an INSTANCE OF THE APPLICATION."
② "Consumers that are part of the transactions DON'T COMMIT THEIR OWN OFFSETS — the PRODUCER writes offsets as part of the transaction. So offset commit should be disabled."
③ "To read transactions cleanly (i.e., ignore in-flight and aborted transactions), we set the consumer isolation level to read_committed. Note that the consumer will still read NONTRANSACTIONAL WRITES, in addition to reading committed transactions." (The input needn't be transactional — "there is no such requirement for the input.")
④ "The first thing a transactional producer must do is initialize. This:
- REGISTERS the transactional ID.
- BUMPS UP THE EPOCH to guarantee that other producers with the same ID will be considered zombies.
- ABORTS older in-flight transactions from the same transactional ID.
⑤ "Here we are using the subscribe consumer API, which means that partitions assigned to this instance can change at any point as a result of rebalance. Prior to release 2.5 (KIP-447), this was much more challenging..." (see §3.7). "When using this method, it also makes sense to commit transactions whenever the related partitions are revoked."
⑥ "This method guarantees that everything that is produced from the time it was called, until the transaction is either committed or aborted, is part of a single atomic transaction."
⑦ "This is where we process the records — all our business logic goes here."
⑧ ⚠️ "it is important to commit the offsets as part of the transaction. This guarantees that if we fail to produce results, we won't commit the offsets for records that were not, in fact, processed. Note that it is important NOT to commit offsets in ANY OTHER WAY — disable offset auto-commit, and DON'T CALL ANY OF THE CONSUMER COMMIT APIs. Committing offsets by any other method does not provide transactional guarantees."
⑨ "Once this method returns successfully, the entire transaction has made it through, and we can continue to read and process the next batch."
⑩ "If we got this exception, it means WE ARE THE ZOMBIE. Somehow our application froze or disconnected, and there is a newer instance of the app with our transactional ID running. Most likely the transaction we started has already been aborted and someone else is processing those records. NOTHING TO DO BUT DIE GRACEFULLY."
⑪ "If we got an error while writing a transaction, we can abort the transaction, set the consumer position back, and try again."
The three-part contract — get any one wrong and you lose the guarantee:
- Producer:
transactional.idset +initTransactions(). - Consumer:
enable.auto.commit=falseand never callconsumer.commitSync/Async— offsets go viaproducer.sendOffsetsToTransaction(). - Downstream consumer:
isolation.level=read_committed— the default is wrong for this.
(Full example: Apache Kafka GitHub, "which includes a demo driver and a simple exactly-once processor that runs in separate threads.")