使用Testcontainers和JUnit测试Kafka消费者恢复契约
一个 Kafka 消费者可能通过了单元测试,却仍然会在生产环境重新投递某条记录时失败。
发生这种情况,是因为单元测试通常检查的是处理器,而不是恢复契约。它验证一条消息产生一个预期副作用。它很少验证同一事件到达两次、格式错误的记录进入主题,或消费者应用了副作用后在提交偏移量之前崩溃时会发生什么。
对 Java 团队来说,Testcontainers 是一个实用的折中方案。它让你能够针对真实的 Kafka broker 运行测试,而无需维护共享的集成环境。但重要的不是容器,而是你测试的契约。
先定义恢复契约
在编写代码之前,先定义消费者承诺什么。一个有用的恢复契约会回答这些问题:
- 业务副作用何时被视为完成?
- Kafka 偏移量何时提交?
- 什么幂等键用于防止重复投递?
- 格式错误的记录会怎样?
- 什么证据表明失败记录是被有意处理的?
对于订单投影消费者,该契约可能是:
- 一个有效事件只更新投影一次。
- 一个重复事件不会应用两次副作用。
- 格式错误的事件会被路由到死信主题,并带有足够的上下文以便检查。
- 如果消费者在副作用之后、偏移量提交之前崩溃,重新投递不会产生重复的业务效果。
这些承诺比泛泛的“Kafka 集成测试”更有用。
将 Testcontainers 用作测试基座
当前 Testcontainers Kafka 模块使用 org.testcontainers.kafka 包下的类。旧的 org.testcontainers.containers.KafkaContainer 类在当前文档中已被弃用。
一个最小的 Maven 测试配置如下:
然后在 JUnit 5 测试中启动 Kafka:
这只是表明 broker 已启动。测试仍然需要一个受控的副作用和可观察的恢复结果。
测试重复投递
Kafka 消费者应假定会发生重复投递,除非整个处理路径被设计为不会如此。这并不意味着每个系统都需要精确一次处理。它意味着业务副作用应通过幂等键来保护。
在这个示例中,事件是一个简单的竖线分隔字符串:
第一个字段是事件 ID,也就是幂等键。
已消费记录的断言很重要。如果没有它,测试可能在第一条记录之后就通过,永远无法证明重复路径被执行过。消费者可能看到两条记录,但按照该契约,投影应只接受该事件 ID 一次。
测试格式错误的记录
无效记录不应消失。它们应进入由你负责的失败路径,并带有足够的上下文以供检查。
在生产测试中,如果这些字段是失败契约的一部分,还应断言原始分区、原始偏移量、错误类、消费者组和 schema 版本。
测试副作用边界
最棘手的 bug 往往存在于业务副作用和偏移量提交之间。
如果消费者在副作用之前提交,崩溃可能导致工作丢失。如果它在副作用之后提交,崩溃可能重新投递一条副作用已经发生过的记录。第二种模式很常见,但只有当副作用是幂等的时候才有效。
这个测试强制触发这个边界:
第一个消费者写入投影,并在提交偏移量之前模拟崩溃,因此对于该消费者组和主题,已提交的偏移量应仍然不存在。第二个消费者使用相同的 group ID,再次收到 evt-2,并依靠幂等键来避免重复的业务效果。
避免基于 sleep 的测试
Kafka 恢复测试常常变得不稳定,因为它们使用 Thread.sleep 并寄希望于消费者能够赶上。
优先使用可观察的条件:
- 投影行达到预期值
- 幂等键只存在一次
- 死信主题包含预期记录
- 已提交偏移量仅在成功处理后才变化
- 重试计数器达到预期上限
使用有界等待工具(例如 Awaitility),但要等待业务条件,而不是等待猜测的延迟。
不应宣称什么
不要宣称 Testcontainers 能完美证明生产行为。它不会复现你的 broker 集群、网络、分区数量、存储行为、安全配置或消费者组规模。
这没关系。这个测试的目标更窄,但仍然有价值:测试你的消费者代码是否能在真实 Kafka broker 上遵守其恢复契约。
检查清单
在称一个 Kafka 消费者已经过恢复测试之前,请验证:
- 重复投递是安全的
- 格式错误的记录有由你负责的失败路径
- 副作用是幂等的
- 围绕成功处理和有意失败路由测试偏移量提交行为
- 重启行为已测试
- 死信记录包含有用的上下文
- 测试断言结果,而不是处理器调用
结论
Kafka 并不承诺你的应用副作用会精确发生一次。它提供日志、偏移量、消费者组和客户端 API。围绕这些原语,恢复契约由你的应用负责。
Testcontainers 和 JUnit 很有用,因为它们让你能够针对真实 broker 测试该契约。但 broker 只是测试基座。真正的工作是定义消费者的承诺,并演练生产环境最终会暴露出来的失败情况。