kafka--基础知识点--5.3-事务
1 事务简介Kafka事务是Apache Kafka在流处理场景中实现Exactly-Once语义的核心机制。它允许生产者在跨多个分区和主题的操作中以原子性Atomicity的方式提交或回滚消息确保数据处理的最终一致性。例如在流处理中消费者读取消息后处理并生成新消息若处理失败事务可确保原始消息的消费偏移与新消息的发送同时回滚避免数据不一致。事务的核心作用:ACIDKafka是否具备程度原子性✅ 有全成功或全失败一致性⚠️ 有限无约束但保证可见性一致隔离性✅ 有两个级别read_committed / read_uncommitted持久性✅ 有acksall 副本持久化Kafka的事务是一个同时涉及生产者和消费者的综合机制。简单来说生产者事务确保写入的原子性多条消息要么全成功要么全失败。消费者确保读取的精确性消费、处理、提交位移要么全成功要么全失败。两者协同通过将“消费-处理-生产”绑定为一个事务实现端到端的恰好一次Exactly-Once 语义kafka事务指的是生产者事务消费者没有事务的概念。问题答案Consumer 有事务概念吗没有完全没有Consumer 有事务 API 吗没有read_committed 是事务吗不是只是读取过滤commit offset 是事务吗不是只是记录消费位置Consumer 怎么实现原子性借 Producer 的事务提交 offset2 事务原理详解了解即可kafka学习笔记四、生产者、消费者客户端深入研究(三)——事务详解及代码实例3 示例fromconfluent_kafkaimportProducer,Consumer,KafkaErrorimportjson# 事务生产者 deftransactional_producer():producerProducer({bootstrap.servers:localhost:9092,transactional.id:txn-producer-1,enable.idempotence:True,# 可以不用加transactional.id后默认Trueacks:all,# 可以不用加transactional.id后默认allretries:2147483647,# 可以不用加transactional.id后默认2147483647max.in.flight.requests.per.connection:5,# # 可以不用加transactional.id后默认5})producer.init_transactions()messages[{id:1,name:tom},{id:2,name:jerry},{id:3,name:lucy},]producer.begin_transaction()formsginmessages:producer.produce(test-topic,keystr(msg[id]).encode(utf-8),valuejson.dumps(msg).encode(utf-8),)producer.commit_transaction()producer.flush()print(事务提交完成)# 消费者手动提交 defmanual_commit_consumer():consumerConsumer({bootstrap.servers:localhost:9092,group.id:my-group,auto.offset.reset:earliest,enable.auto.commit:False,isolation.level:read_committed,})consumer.subscribe([test-topic])whileTrue:msgconsumer.poll(1.0)ifmsgisNone:continueifmsg.error():continuedatajson.loads(msg.value().decode(utf-8))print(f收到: partition{msg.partition()}, offset{msg.offset()}, data{data})consumer.commit(msg,asynchronousFalse)consumer.close()if__name____main__:importsysifsys.argv[1]producer:transactional_producer()elifsys.argv[1]consumer:manual_commit_consumer()