Aborted records stay in the log; consumers decide whether to see them. With isolation.level=read_committed the broker returns records only up to the last stable offset, the first offset of the oldest open transaction, and the client drops records of aborted transactions. read_uncommitted returns everything. Java consumers default to read_uncommitted, librdkafka to read_committed.
from confluent_kafka import Consumer, Producer
B = "localhost:33092"
p = Producer({"bootstrap.servers": B, "transactional.id": "booknest-isolation-demo"})
p.init_transactions()
for outcome in ("commit", "abort", "commit"): # three transactions of two records
p.begin_transaction()
for n in (1, 2):
p.produce("booknest.tx-demo", f"{outcome} {n}")
p.flush() # the records reach the broker either way
p.commit_transaction() if outcome == "commit" else p.abort_transaction()
for level in ("read_uncommitted", "read_committed"):
c = Consumer({"bootstrap.servers": B, "group.id": f"demo-{level}",
"isolation.level": level, "auto.offset.reset": "earliest"})
c.subscribe(["booknest.tx-demo"])
values = [m.value().decode() for m in c.consume(100, timeout=5)]
print(f"{level:<17}", values)
c.close()Output
read_uncommitted ['commit 1', 'commit 2', 'commit 1', 'commit 2', 'abort 1', 'abort 2'] read_committed ['commit 1', 'commit 2', 'commit 1', 'commit 2']
The aborted pair was written and only filtered out on reading. A long transaction holds back every read_committed consumer of its partitions until it ends, so keep transactions short, seconds rather than minutes.