基于Kafka的日志传输链路断连恢复测试:能否通过Proxy Server实现?
用Proxy Server测试Kafka链路断连恢复场景的可行方案
完全可以通过搭建Proxy Server来实现这类测试,这是一种低成本、对现有环境无侵入的测试方式,下面具体讲实现思路和Kafka相关的关键细节:
一、代理服务器的作用逻辑
- 把代理部署在你的程序和Kafka集群之间,让程序的Kafka生产者客户端直接连接代理,再由代理转发请求到Kafka集群
- 你可以通过启停代理、阻断代理与Kafka的通信,轻松模拟链路断开/恢复的场景,完全复现真实网络故障
二、具体操作步骤
搭建代理
- 用现成工具就行,比如
mitmproxy(支持TCP模式转发),或者写个简单的Python TCP代理脚本(核心就是监听端口、转发数据包) - 修改你的程序配置:把Kafka生产者的
bootstrap.servers改成代理服务器的IP和端口,替代原来的Kafka集群地址
- 用现成工具就行,比如
模拟断连
- 最简单的方式:直接停止代理进程,此时程序和Kafka的链路直接中断
- 更贴近真实网络故障:用防火墙阻断代理到Kafka的通信,比如Linux下执行
iptables -A OUTPUT -d <Kafka集群IP> -j DROP - 此时观察程序日志,看生产者是否触发重试、消息是否被缓存,有没有抛出预期内的异常
模拟恢复
- 重启代理进程,或者删除防火墙规则(
iptables -D OUTPUT -d <Kafka集群IP> -j DROP) - 检查Kafka的消费端或日志系统,确认断连期间产生的日志消息是否全部被正常接收,验证消息完整性和顺序性
- 重启代理进程,或者删除防火墙规则(
三、Kafka生产者必须关注的配置
retries:设置足够大的重试次数(比如retries=10),确保链路恢复后能自动重试发送缓存的消息retry.backoff.ms:控制重试间隔(比如retry.backoff.ms=1000),避免频繁重试消耗过多资源acks=all:如果要求消息不丢失,一定要设这个值,确保消息被Kafka集群的ISR副本确认后才视为发送成功buffer.memory:设置合适的缓冲区大小(比如buffer.memory=33554432,即32M),断连期间生产者会把消息暂存到这里,避免消息丢失max.block.ms:可以设置这个参数(比如max.block.ms=60000),控制生产者发送消息时的最长阻塞时间,避免程序因长时间断连挂起
四、额外测试场景建议
- 测试不同断连时长:比如几秒的短断连、几十分钟的长断连,验证生产者缓存的承受能力
- 测试高并发场景:断连期间让程序持续产生大量日志,看缓冲区会不会溢出,生产者的容错逻辑是否正常
内容的提问来源于stack exchange,提问作者Olga Gontsova
相关产品推荐
相关产品推荐

