Redis Streams 引入消费者组之后,消息的读取和确认就变成了两个独立动作。XREADGROUP 负责读取消息,而 XACK 负责告诉 Redis:这条消息已经被当前消费者处理完毕,可以从该消费者的待处理列表中移除。很多人只调用 XREADGROUP 读取数据,却忽略了 XACK,结果发现同一条消息会被重复投递。本文将围绕 XACK 的机制、语法和实际使用场景展开。

消费者组中的 pending 消息是如何产生的
在 Redis Streams 中,每个消费者组都有自己的游标和待处理列表。当某个消费者通过 XREADGROUP 读取消息时,Redis 会把消息 ID 写入该消费者的 pending entries list,同时记录投递次数和最后投递时间。这个动作并不表示消息已完成,而是表示消息已经被投递给某个消费者,等待处理结果。
如果没有调用 XACK,这条消息会一直停留在 pending 列表里。即使消费者已经处理完业务逻辑并正常结束,Redis 仍然认为该消息处于未确认状态。一旦消费者因为异常退出、网络断开或消费者组重新分配分区,pending 中的消息就可能再次被投递。因此,理解 pending 的产生与清除机制,是理解 XACK 价值的前提。
# 创建消费者组,从当前末尾开始监听 XGROUP CREATE mystream mygroup $ MKSTREAM # 消费者读取两条消息 XREADGROUP GROUP mygroup consumer-1 COUNT 2 STREAMS mystream > # 查看所有待处理消息 XPENDING mystream mygroup
上面的命令中,> 表示只读取从未投递给任何消费者的新消息。执行 XREADGROUP 后,返回的消息 ID 会进入 consumer-1 的 pending 列表。接下来若不对这些 ID 执行 XACK,它们会长期存在。
XACK 命令的具体用法
XACK 命令的语法比较简单:XACK key group id [id ...]。第一个参数是 Stream 的键名,第二个参数是消费者组名称,后面的参数是一个或多个消息 ID。执行成功后,命令返回被成功确认的消息数量。
例如,确认一条消息:
XACK mystream mygroup 1718000000000-0 # (integer) 1
如果传入多个 ID,返回的是实际确认成功的数量。如果某个 ID 原就不属于该消费者组,或者已经被确认过,则该 ID 不会计入返回值。比如第二次确认同一个 ID,返回值会是 0。
XACK mystream mygroup 1718000000000-0 # (integer) 0
可以通过 XPENDING 在确认前后对比。确认前,XPENDING mystream mygroup 会显示 pending 数量以及最早、最晚的空闲时间等信息。执行 XACK 后,pending 数量减少,消费者级别的 pending 列表中也不再包含该消息 ID。
不确认消息会造成什么影响
最直接的影响就是重复消费。消费者从组里读取消息后,如果业务处理成功但忘记调用 XACK,这条消息不会自动消失。下次有新的消费者启动,或者当前消费者重新加入组时,pending 列表中的消息可能被重新分配给其他消费者。
例如,consumer-1 读取消息后崩溃,消息在 pending 中保留。此时 consumer-2 可以通过 XCLAIM 或 XAUTOCLAIM 接管这些消息。接管后的消息被 consumer-2 重新处理。如果业务逻辑不具备幂等性,就会出现数据重复写入、重复扣款等问题。
# consumer-2 接管空闲超过 60 秒的消息 XAUTOCLAIM mystream mygroup consumer-2 60000 0-0 # 查看 consumer-2 当前 pending 的消息 XPENDING mystream mygroup - + 10 consumer-2
另一个容易被忽略的问题是 pending 列表增长。若长期不确认,每个消费者的 pending 列表会不断膨胀,XPENDING 查询会变慢,也会增加内存压力。对于消费量大的系统,建议在监控中加入 pending 数量和最长空闲时间指标。
如何设计可靠的消息确认流程
最常见的方式是在业务处理成功后再调用 XACK。这样可以避免消息在处理前就被确认,导致业务尚未完成时消息丢失。比如在处理数据库事务时,先提交事务,再执行 XACK。如果 XACK 失败,消息会留在 pending 中等待重新处理,此时业务逻辑需要保证幂等。
批量确认是另一种常用策略。当一次读取多条消息时,可以先把所有成功处理的消息 ID 收集起来,最后一次调用 XACK。这样做能减少网络往返次数,但要注意某些 ID 失败时不要影响整体。如果某条消息处理失败,可以跳过该 ID,让它继续留在 pending 中,后续通过重试机制处理。
def handle_messages():
data = redis.xreadgroup('mygroup', 'consumer-1', {'mystream': '>'}, count=10)
if not data:
return
ack_ids = []
for stream, entries in data:
for msg_id, fields in entries:
try:
process_message(fields)
ack_ids.append(msg_id)
except Exception as exc:
# 处理失败时不确认,等待重试
continue
if ack_ids:
print(redis.xack('mystream', 'mygroup', *ack_ids))
需要注意的是,XACK 只对消费者组有效,它不会删除 Stream 中的消息。如果需要删除消息本身,应使用 XDEL。同时,确认后消息虽然从 pending 列表中移除,但仍然可以被 XRANGE、XREVRANGE 等命令读取,这部分属于 Stream 的存储语义,不要混淆。
常见误区与排查
一个常见误区是认为消息一旦被读取,Redis 就会自动确认。实际上,XREADGROUP 只负责投递,不负责确认。另一个误区是把 XACK 理解成删除消息,其实它只确认消费者组的处理状态。
排查问题时,可以先执行 XPENDING mystream mygroup 查看 pending 总数。如果 pending 持续增长,基本可以判断 XACK 调用缺失或失败。再通过 XPENDING mystream mygroup - + 10 consumer-1 查看具体消息 ID 和投递时间,结合业务日志判断问题发生在处理阶段还是确认阶段。
# 查看 pending 总数和边界空闲时间 XPENDING mystream mygroup # 查看某个消费者的 pending 详情 XPENDING mystream mygroup - + 10 consumer-1
还有一种情况是 XACK 返回 0,但开发人员没检查返回值,导致问题被掩盖。应当在代码中判断返回数量,并与期望确认的消息数做对比。若返回数量小于入参数量,说明部分 ID 已经不存在或不属于当前组,需要记录日志并排查。
Redis XACKRedis Streams消费者组修改时间:2026-09-27 12:30:13