
AutoMQ Kafka 客户端示例实战异步/同步生产消费与 Exactly-Once 事务处理【免费下载链接】automqDiskless Kafka® on S3. 10x Cost-Effective. No Cross-AZ Traffic Cost. Autoscale in seconds. Single-digit ms latency. Multi-AZ Availability.项目地址: https://gitcode.com/GitHub_Trending/au/automqexamples/README.md 是 AutoMQDiskless Kafka® on S3仓库中的 Kafka 客户端示例模块说明它用两个可直接运行的 Demo 演示了最核心的两类客户端能力基于KafkaProducer/KafkaConsumer的异步与同步生产消费以及基于事务 API 的 read-process-write 精确一次Exactly-Once简称 EOS处理。本文以该 README 为主线结合模块内全部源码完整还原这两个 Demo 的用法、执行流程与底层实现细节读完你可以在本地集群上复现运行并掌握生产消费者配置、事务边界、消费位移提交等关键编码范式。示例模块概览与源码结构整个示例模块位于仓库的 examples 目录除 README.md 外全部源码集中在src/main/java/kafka/examples/包下共 8 个 Java 文件职责划分清晰文件职责KafkaConsumerProducerDemo.java生产-消费 Demo 入口串起清理 Topic → 生产者线程发送 → 消费者线程拉取三个阶段Producer.java通用生产者线程支持异步默认与同步两种发送模式可配置幂等与事务Consumer.java通用消费者线程支持普通消费与read_committed隔离级别消费KafkaExactlyOnceDemo.java精确一次 Demo 入口编排输入 Topic、多个事务实例与校验消费者ExactlyOnceMessageProcessor.javaread-process-write 事务处理器EOS 的核心实现TransactionProducer.java最小化的事务生产者示例单独演示initTransactions/beginTransaction/commitTransaction/abortTransactionKafkaProperties.java集中配置连接地址BOOTSTRAP_SERVERS localhost:9092Utils.java公共工具日志打印、抽样输出、Topic 删除与重建这些示例全部基于标准org.apache.kafka.clients客户端 API 编写因此既可以用官方 Apache Kafka 2.5 集群运行也可以直接连接 AutoMQAutoMQ 完全兼容 Kafka 客户端协议运行。运行前置启动本地集群与连接配置README 的第一步要求是Start a Kafka 2.5 local cluster with a plain listener configured on port 9092.即本地需要先有一个监听9092端口的明文PLAINTEXT监听器的 Kafka 2.5 集群。这与示例代码中的连接配置严格对应——KafkaProperties.java 将 bootstrap servers 硬编码为localhost:9092所有 Producer、Consumer 和 Admin 客户端都统一使用该常量。若你的集群不在本机 9092 端口修改这一个文件即可。在本仓库中你可以直接使用 AutoMQ 自带的集群配置来搭建运行环境单机/简单部署可参考 config/kraft/server.propertiesAutoMQ 的 KRaft 模式 broker 配置以及 config/kraft/broker.properties、config/kraft/controller.properties 等独立节点配置需要容器化快速拉起完整 AutoMQ 集群时可参考 docker/README.md 与 devkit/README.md 中提供的 docker-compose 方式如 devkit/docker-compose.yml、docker/docker-compose.yaml。需要特别说明的版本前提精确一次示例依赖消费位移随事务一起提交的能力sendOffsetsToTransaction结合groupMetadata该 API 仅在 broker 2.5 及以上可用若连接更早版本的 brokerKafkaExactlyOnceDemo.java 会抛出org.apache.kafka.common.errors.UnsupportedVersionException。生产-消费 Demo异步与同步两种发送模式用法与参数README 给出了两条命令# 异步发送 10000 条记录到 topic1my-topic并消费 examples/bin/java-producer-consumer-demo.sh 10000 # 同步发送 10000 条记录到 topic1 并消费 examples/bin/java-producer-consumer-demo.sh 10000 sync脚本入口对应主类kafka.examples.KafkaConsumerProducerDemo。从 KafkaConsumerProducerDemo.java 的参数解析逻辑看它接收两个参数records必填要发送的记录总数即第一个参数10000mode可选传入sync则同步发送缺省或传其他值则走异步发送。源码中isAsync的判定是args.length 1 || !args[1].trim().equalsIgnoreCase(sync)也就是说只有显式传入sync才切到同步模式。如果不传任何参数程序会打印帮助信息后退出Utils.printHelp见 Utils.java。如果你在 IntelliJ IDEA 中运行主类源码注释KafkaConsumerProducerDemo.java明确提示这两个参数应填入Modify Run Configuration - Program Arguments并可通过Modify options - Save console output to file将日志输出保存到文件。执行流程三阶段编排Demo 内部通过CountDownLatch(2)同步两个线程主线程最长等待 5 分钟KafkaConsumerProducerDemo.java超时则强制关闭两个线程。整体分三阶段清理历史 Topic调用Utils.recreateTopics(KafkaProperties.BOOTSTRAP_SERVERS, -1, my-topic)删除上次运行残留的 Topic 再重建保证每次运行环境干净生产者线程发送创建名为producer的线程向my-topic发送numRecords条记录消费者线程拉取创建名为consumer的线程以消费者组my-group订阅my-topic拉取并打印全部记录。Topic 名与消费者组名来自 KafkaConsumerProducerDemo.java 的两个常量TOPIC_NAME my-topic、GROUP_NAME my-group。生产者实现要点Producer.java 的createKafkaProducer()L106-L128展示了最基础也最关键的生产者配置bootstrap.servers连接地址必填client.id非必填但推荐方便在服务端请求日志中按逻辑应用名追踪请求来源示例中生成随机 UUIDkey.serializer/value.serializerkey 使用IntegerSerializervalue 使用StringSerializertransaction.timeout.ms仅当transactionTimeoutMs 0时设置表示事务协调器主动中止未完成事务的最长时间transactional.id仅当传入非空transactionalId时设置该值必须静态且唯一用于进程重启后标识同一个生产者实例这是事务与幂等的基石enable.idempotence按示例传入的布尔值开启/关闭分区级幂等重复发送保护。发送逻辑上两种模式差异体现在 run() 与两个私有方法异步发送asyncSendproducer.send(new ProducerRecord(topic, key, value), callback)立即返回结果通过ProducerCallback.onCompletion回调通知。回调里对异常做了分类RetriableException可继续重试其余不可恢复异常则调用shutdown()终止Producer.java。源码注释特别提醒即使设置很小的batch.size且linger.ms0当buffer.memory已满或元数据不可用时send依然会被阻塞——异步不代表永不阻塞同步发送syncSendproducer.send(...).get()阻塞等待 broker 的 ack拿到RecordMetadata后打印抽样信息。它对AuthorizationException、UnsupportedVersionException、ProducerFencedException、FencedInstanceIdException、OutOfOrderSequenceException、SerializationException等不可恢复异常直接终止对其他KafkaException仅记录错误。消费者实现要点Consumer.java 是一个实现ConsumerRebalanceListener的消费者线程核心配置见createKafkaConsumer()L126-L149group.id使用subscribe(topics)做组管理时必填group.instance.id可选设置后启用静态成员static membership可提升滚动重启等场景下的可用性enable.auto.commit普通消费模式为trueEOS 模式下强制为false因为位移要随事务一起提交isolation.level仅在readCommitted模式下设置为read_committed跳过进行中与已中止的事务auto.offset.reset固定为earliest无有效位移时从头消费key/value 反序列化器与生产者对应IntegerDeserializer/StringDeserializer。消费主循环Consumer.java对异常做了细致的分层处理AuthorizationException、UnsupportedVersionException、RecordDeserializationException不可恢复直接关闭OffsetOutOfRangeException、NoOffsetForPartitionException无有效位移且无 reset 策略时抛出则seekToEnd后commitSync提交当前位移再继续其余KafkaException仅记录日志后继续。三个 rebalance 回调onPartitionsRevoked/onPartitionsAssigned/onPartitionsLost分别打印撤销、分配、丢失的分区方便理解组再平衡过程手动提交位移的应用可在onPartitionsRevoked中提交待处理位移。辅助逻辑Topic 清理与抽样输出Utils.java 的recreateTopicsL68-L105使用 Admin 客户端完成先删后建删除 Topic 时忽略UnknownTopicOrPartitionException首次运行无 Topic 属正常创建 Topic 放在重试循环里捕获TopicExistsException后等待 1 秒重试删除元数据尚未完全清理时会遇到创建时replicationFactor传-1使用 broker 默认副本数源码注释点明这是为了避免在minISR 1时触发NOT_ENOUGH_REPLICAS错误。而maybePrintRecordL53-L66在记录总数超过 20 时只抽样打印约 10 条按key % max(1, numRecords/10) 0采样避免海量日志刷屏。Exactly-Once 精确一次 Demoread-process-write 事务处理用法与参数README 给出的命令为examples/bin/exactly-once-demo.sh 6 3 10000对应主类kafka.examples.KafkaExactlyOnceDemo严格接收 3 个参数KafkaExactlyOnceDemo.javapartitioninput-topic 与 output-topic 的分区数示例为 6instances事务应用实例数示例为 3records总记录数示例为 10000。效果是创建各 6 分区的input-topic与output-topic启动 3 个事务处理实例把 10000 条记录从输入 Topic 读出来、处理后写入输出 Topic并用read_committed消费者校验结果。四阶段执行流程KafkaExactlyOnceDemo.java 用四个阶段完成整个编排每个阶段都有独立的 2 分钟超时保护清理并重建 TopicUtils.recreateTopics(..., numPartitions, input-topic, output-topic)注意这里显式指定了 6 个分区预灌数据启动一个同步生产者isAsyncfalse向input-topic写入 10000 条 key 为偶数的记录test0、test2…用CountDownLatch(1)阻塞等待数据装载完成启动事务实例用IntStream创建numInstances个ExactlyOnceMessageProcessor线程processor-0、processor-1、processor-2每个实例执行独立的 read-process-write 循环各自通过比较日志末尾偏移log end offset与已提交偏移来排空分配到自身分区的所有记录阻塞直到全部记录处理完成并写入output-topic校验输出创建一个read_committed消费者从output-topic消费全部记录验证分区级顺序保持。整个 Demo 的类注释KafkaExactlyOnceDemo.java给出了一个关键承诺每个记录被精确处理一次且分区级强有序——这正是事务 幂等 位移随事务提交三者共同保证的结果。事务处理核心实现ExactlyOnceMessageProcessor.java 是 EOS 的灵魂其构造与运行逻辑值得逐点拆解事务与组身份每个实例的transactionalId tid- threadName如tid-processor-0事务超时设为较短的 10 秒源码注释说明短事务超时有助于更快清理 pending offsets消费者使用组processor-group并设置groupInstanceId giid- threadName启用静态成员以尽量避免不必要的 rebalance消费者必须 read_committed读取输入 Topic 的消费者以read_committed隔离级别运行ExactlyOnceMessageProcessor.java这样它永远不会读到未提交的数据事务生命周期run()启动后先调用一次producer.initTransactions()作用是隔离僵尸生产者实例并中止任何挂起的事务只调用一次随后进入循环——poll(200ms)拉到一批记录后执行beginTransaction()→ 逐条处理并发送到output-topic示例中把 value 追加-ok后缀→producer.sendOffsetsToTransaction(getOffsetsToCommit(consumer), consumer.groupMetadata())把消费位移也纳入当前事务此 API 仅 broker 2.5 可用→commitTransaction()一次性提交消息与位移位移采集getOffsetsToCommitL203-L209遍历消费者当前分配的分区取consumer.position(topicPartition)即下一条待消费位置封装成OffsetAndMetadata完成判定getRemainingRecordsL211-L225通过endOffsets(assignment)对比各分区末尾偏移与当前消费位置求和得到剩余记录数若暂时拿不到末尾偏移则返回Long.MAX_VALUE表示继续拉取。事务失败处理与重试策略处理循环对异常同样做了分类ExactlyOnceMessageProcessor.java不可恢复异常AuthorizationException、UnsupportedVersionException、ProducerFencedException、FencedInstanceIdException、OutOfOrderSequenceException、SerializationException直接关闭实例无有效位移则seekToEnd后commitSync而一般KafkaException则先abortTransaction()中止当前事务再进入maybeRetry重试逻辑L236-L263重试次数上限为MAX_RETRIES 5L53重试时把消费者拉取位置回退到事务开始前的已提交偏移consumer.committed(assignment)后逐个seek没有已提交位移的分区则seekToBeginning从而保证重放的是同一批记录超过 5 次仍失败则调用commitSync()把当前位移当作已处理来跳过这批记录——源码注释L227-L235明确建议生产环境应当把这类记录投递到死信 TopicDLT做进一步处理。与之配套模块里还提供了一个最小化的事务生产者示例 TransactionProducer.java只需配置transactional.id示例为my-transactional-id即可用initTransactions → beginTransaction → send ×20 → commitTransaction完成一次事务性写入异常时走abortTransactionfinally中关闭生产者。它是理解单个生产者事务的最小可运行样例适合在跑通 EOS Demo 前先单独验证集群的事务能力。关键配置参数速查综合两个 Demo 的源码以下参数在连接 AutoMQ/Kafka 集群跑客户端示例时最常用参数示例取值作用与说明bootstrap.serverslocalhost:9092集群连接地址集中定义在 KafkaProperties.javaclient.idclient-uuid服务端请求日志中区分应用来源非必填但推荐key/value.serializerIntegerSerializer/StringSerializer消息 key/value 的序列化器enable.idempotencetrue/false分区级幂等重复发送保护事务场景必须开启transactional.idtid-processor-0事务生产者身份须静态唯一跨进程重启有效transaction.timeout.ms10000事务协调器中止挂起事务的最长等待时间group.idmy-group/processor-group消费者组subscribe组管理模式必填group.instance.idgiid-processor-0静态成员身份减少不必要的 rebalanceenable.auto.committrue/falseEOS 场景须为false位移改由事务提交isolation.levelread_committed只读已提交数据EOS 消费者必配auto.offset.resetearliest无有效位移时的重置策略总结与进一步阅读通过examples模块你可以在一套本地集群上依次验证 AutoMQ/Kafka 客户端的三种典型用法异步发送 消费java-producer-consumer-demo.sh 10000、同步发送 消费java-producer-consumer-demo.sh 10000 sync以及多实例 read-process-write 精确一次处理exactly-once-demo.sh 6 3 10000。其中 EOS Demo 的消息 位移同事务提交、read_committed 校验、有限次重试后跳过的容错组合是生产级流处理应用可直接借鉴的完整范式。想继续深入可以从这些路径入手阅读完整示例说明examples/README.md查看各示例实现Producer.java、Consumer.java、ExactlyOnceMessageProcessor.java、TransactionProducer.java了解 AutoMQ 集群的本地启动方式docker/README.md、devkit/README.md以及 KRaft 模式配置 config/kraft/server.properties参考项目对 Kafka 协议兼容性的工程实现可从 clients 与 server 模块入手。【免费下载链接】automqDiskless Kafka® on S3. 10x Cost-Effective. No Cross-AZ Traffic Cost. Autoscale in seconds. Single-digit ms latency. Multi-AZ Availability.项目地址: https://gitcode.com/GitHub_Trending/au/automq创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考