📚 消息队列系列 · 第 5 篇(共 6 篇)📌 上一篇:Python 上手 Kafka,绕不开分区和消费者组📌 本篇主角:消息队列的三大经典难题——丢失、重复、乱序📌 你将学到:每个问题为什么发生、测试怎么发现、怎么用 Python 验证📌 难度指数:⭐⭐(需要 Python 基础)
做过消息队列测试的同学都懂——测 MQ 系统,翻来覆去就是三件事:
消息丢了没? (丢失)消息重复没? (重复消费)消息顺序对没? (乱序)
这三个问题不解决,轻则数据对不上,重则订单错乱、钱算错、库存对不上。
本篇不讲空洞的理论,直接讲:每个问题怎么发生 → 测试怎么发现 → 用 Python 怎么验。
难题一:消息丢失 😱
为什么消息会丢?
消息从生产到消费要经过 3 个环节,每个环节都可能丢:
生产者 → (环节1) → Broker → (环节2) → 消费者 发送失败 存储丢失 处理失败
测试怎么发现?
核心思路:发 N 条,收 N 条,对数量。
"""消息丢失测试:发 100 条,看能收到几条"""from kafka import KafkaProducer, KafkaConsumerBOOTSTRAP = 'localhost:9092'TOPIC = 'loss_test'TOTAL = 100# 1. 发 100 条带编号的消息producer = KafkaProducer(bootstrap_servers=BOOTSTRAP)for i in range(TOTAL): producer.send(TOPIC, f"msg-{i}".encode())producer.flush()producer.close()print(f"📤 已发送 {TOTAL} 条")# 2. 收消息,统计收到了哪些编号consumer = KafkaConsumer( TOPIC, bootstrap_servers=BOOTSTRAP, auto_offset_reset='earliest', group_id='loss_check', consumer_timeout_ms=8000, # 8秒没新消息结束)received = []for msg in consumer: received.append(msg.value.decode()) if len(received) >= TOTAL: break# 3. 对比:有没有丢?sent = {f"msg-{i}" for i in range(TOTAL)}got = set(received)lost = sent - got # 应该收到但没收到的duplicated = [m for m in received if received.count(m) > 1] # 重复的print(f"📥 收到 {len(received)} 条")print(f"❌ 丢失 {len(lost)} 条: {sorted(lost)[:5]}...")
💡 这个脚本就是"消息完整性测试"的原型——发编号消息,收回来比对,丢没丢、重复没重复一目了然。测试工具的核心就这么简单。
怎么防丢失?(生产环境)
生产端:开启确认机制(Kafka acks=all / RabbitMQ confirm)Broker:持久化 + 副本(durable / replication)消费端:手动 ack,处理成功才确认
测试要点:重点验证"消费端处理失败时,消息会不会被重新投递"——这是防丢失最关键的一环。
难题二:重复消费 🔄
为什么消息会重复?
这是最常被问倒的问题。核心原因就一个:
消费者处理完了消息,但在"告诉 MQ 我处理完了"之前,挂了。
消费者:处理消息 ✅ → 准备 ack → 💀 进程挂了MQ:没收到 ack → 以为没处理 → 重新投递消费者:又收到同一条消息 → 重复处理!
再比如 Kafka 里:处理完消息但没来得及提交 offset,重启后从旧 offset 重新读,也会重复。
测试怎么发现?
核心思路:发重复消息,看消费方处理了几次。
"""重复消费测试:同一条消息发两次,看消费者是否处理两次"""import timefrom kafka import KafkaProducer, KafkaConsumerBOOTSTRAP = 'localhost:9092'TOPIC = 'duplicate_test'MSG_ID = 'order-10086' # 唯一的业务ID# 1. 发送同一条消息两次(模拟 MQ 重投)producer = KafkaProducer(bootstrap_servers=BOOTSTRAP)producer.send(TOPIC, MSG_ID.encode())producer.flush()producer.send(TOPIC, MSG_ID.encode()) # 模拟重复投递producer.flush()producer.close()# 2. 消费者统计处理次数processed = {}consumer = KafkaConsumer( TOPIC, bootstrap_servers=BOOTSTRAP, auto_offset_reset='earliest', group_id='dup_check', consumer_timeout_ms=5000,)for msg in consumer: mid = msg.value.decode() processed[mid] = processed.get(mid, 0) + 1 print(f"📥 收到 {mid},这是第 {processed[mid]} 次")# 3. 判断:处理了 2 次 = 有重复消费问题for mid, count in processed.items(): if count > 1: print(f"❌ 消息 {mid} 被消费了 {count} 次!存在重复消费") else: print(f"✅ 消息 {mid} 只消费了 1 次")
怎么解决?(幂等性)
幂等:同一个操作做一次和做一百次,结果一样。
方案1:数据库唯一键 订单ID建唯一索引,重复插入直接报错跳过方案2:Redis 去重 处理前先 SETNX(msg_id),已存在就说明处理过方案3:业务判断 订单状态已经是"已发货",就不再重复发货
# 幂等处理示例:用 Redis 去重import redisr = redis.Redis(decode_responses=True)def process_order(order_id): # SETNX:不存在才设置成功,返回 True if r.setnx(f"order:{order_id}", "processed"): print(f"✅ 第一次处理订单 {order_id}") # ... 真正的业务处理 ... else: print(f"⏭️ 订单 {order_id} 已处理过,跳过")
测试要点:专门构造"同一条消息投递两次"的场景,验证消费方幂等逻辑是否生效。
难题三:消息乱序 🔀
为什么消息会乱序?
同一个业务的多条消息,处理顺序不能反。
比如订单状态:待支付 → 已支付 → 已发货,如果消费者先处理了"已发货"再处理"已支付",状态就乱了。
乱序的根源:
原因1:多个消费者并行处理 消息1给消费者A,消息2给消费者B → 可能B先处理完原因2:网络/重试 消息1处理失败重试时,消息2已经处理完了原因3:Kafka 多分区 同一订单的消息进了不同分区 → 顺序无法保证
测试怎么发现?
核心思路:发有序消息,检查接收顺序。
"""乱序测试:发 1-10 条有序消息,看接收顺序"""from kafka import KafkaProducer, KafkaConsumerBOOTSTRAP = 'localhost:9092'TOPIC = 'order_test'# 1. 发 10 条有序消息producer = KafkaProducer(bootstrap_servers=BOOTSTRAP)for i in range(1, 11): producer.send(TOPIC, f"step-{i:02d}".encode())producer.flush()producer.close()# 2. 接收并记录顺序consumer = KafkaConsumer( TOPIC, bootstrap_servers=BOOTSTRAP, auto_offset_reset='earliest', group_id='order_check', consumer_timeout_ms=5000,)received = []for msg in consumer: received.append(msg.value.decode()) if len(received) >= 10: break# 3. 检查是否有序expected = [f"step-{i:02d}" for i in range(1, 11)]if received == expected: print("✅ 顺序完全正确")else: print("❌ 消息乱序!") print(f" 期望: {expected}") print(f" 实际: {received}") # 找出乱序的位置 for i, (e, a) in enumerate(zip(expected, received)): if e != a: print(f" 👉 第 {i+1} 条开始乱序: 期望 {e},实际 {a}")
怎么保证有序?
方案1:单分区 + 单消费者 Kafka 一个分区内天然有序,只用一个消费者消费 (顺序要求极高的场景)方案2:同一个业务ID 路由到同一分区 Kafka 按 key 哈希决定分区 → 同一订单永远进同一分区方案3:消费者串行处理 并发消费改成串行(或加锁),天然有序
# 方案2:按订单ID 路由到同一分区(Kafka)producer = KafkaProducer(bootstrap_servers=BOOTSTRAP)for step in range(3): producer.send( TOPIC, key=b'order-10086', # ← 同一订单固定 key value=f"step-{step}".encode() # → 永远进同一分区,有序! )
测试要点:重点验证"同一业务ID的消息是否进同一分区",这是 Kafka 保序的关键。
三大难题速查表
实战:一个完整的 MQ 可靠性测试
把三个测试合成一个脚本,测完一轮就知道系统靠不靠谱:
"""消息队列可靠性测试:丢失 + 重复 + 乱序 一把测"""from kafka import KafkaProducer, KafkaConsumerfrom collections import CounterBOOTSTRAP = 'localhost:9092'TOPIC = 'reliability_test'TOTAL = 50# ===== 1. 发送带序号的消息 =====producer = KafkaProducer(bootstrap_servers=BOOTSTRAP)for i in range(TOTAL): producer.send(TOPIC, key=b'fixed-key', value=f"msg-{i:03d}".encode())producer.flush()producer.close()print(f"📤 已发送 {TOTAL} 条(同一 key,保证同分区有序)\n")# ===== 2. 接收全部消息 =====consumer = KafkaConsumer( TOPIC, bootstrap_servers=BOOTSTRAP, auto_offset_reset='earliest', group_id='reliability_check', consumer_timeout_ms=8000,)received = []for msg in consumer: received.append(msg.value.decode()) if len(received) >= TOTAL: breakconsumer.close()# ===== 3. 三项检查 =====print("=" * 40)# 检查1:丢失got_ids = set(received)lost = [f"msg-{i:03d}" for i in range(TOTAL) if f"msg-{i:03d}" not in got_ids]print(f"🔍 丢失检查: {'✅ 无丢失' if not lost else f'❌ 丢失 {len(lost)} 条: {lost[:5]}'}")# 检查2:重复counts = Counter(received)dups = {k: v for k, v in counts.items() if v > 1}print(f"🔍 重复检查: {'✅ 无重复' if not dups else f'❌ 重复: {dups}'}")# 检查3:乱序expected = [f"msg-{i:03d}" for i in range(TOTAL)]print(f"🔍 乱序检查: {'✅ 顺序正确' if received == expected else '❌ 消息乱序'}")# 最终结论if not lost and not dups and received == expected: print("\n🎉 全部通过!MQ 系统可靠")else: print("\n⚠️ 存在问题,请对照速查表定位修复")
测试工程师必问的三个问题
问自己这三个问题,能答上就说明真的懂了:
1. 消息丢了,怎么定位是哪个环节丢的? → 分环节测:生产端用 confirm 确认,Broker 看日志,消费端手动 ack2. 重复消费了,业务会不会出问题? → 看有没有幂等:有唯一键/去重 → 安全;没有 → 必须修3. 顺序要求高的消息,怎么保证? → 同业务固定 key 进同分区,消费者串行处理
总结
| |
|---|
| 丢失 | 发 N 收 N 对数量,重点测"消费失败会不会重投" |
| 重复 | 投递两次验证幂等,Redis SETNX / 数据库唯一键 |
| 乱序 | |
| 测试本质 | |
MQ 测试的精髓就一句话:把系统往坏处想,再验证它能不能扛住。
📝 下篇预告:Python 消息队列测试实战:验证、mock、造数据
怎么造 MQ 测试数据?怎么 mock 消息?测试环境怎么搭?下篇一次讲透。
关注我,软件测试实战干货持续更新 🚀