list 队列丢消息的位置
一句话总结
BRPOP同时在队列中删除并删除信息。删除后不久,如果消费者死亡,该信息就会从世界上消失。
为什么需要这个?
用Redis列表创建队列只需两行代码。LPUSH q:jobs "..."放入BRPOP q:jobs 0打开。简单快速,大多数侧项目都用这个就足够了。
问题就在失败的路上。BRPOP在这个返回的瞬间,该消息已经从Redis中消失了。如果消费者处理完它后死亡,即使重新启动,也不会有该消息。在部署过程中,仅仅重新启动工作者,处理中的工作就会消失。
怎么行动
第一个解决办法是LMOVE(旧版本的RPOPLPUSH)。同时将其移到“处理中”列表中。消费者处理完成后从“处理中”列表中删除。死亡时,消息会留在“处理中”列表中,单独回收器将旧项恢复到原来的队列中。这是可见超时的手动版本。
第二个解决方案是Redis Stream。Stream是从一开始就考虑到这个问题而创建的。XADD添加到罗,创建消费者群组,XREADGROUP读作。读到的消息不会被删除,而是进入PEL(Pending Entries List)。处理结束后XACK消失。死亡后留在PEL中,XPENDING确认后XCLAIM可以让其他消费者带走。
整理三个性质的话是这样的。
| 性格 | 列表 | 流媒体 |
|---|---|---|
| 消费后保存 | 无(立即删除) | PEL 保存,XACK 删除 |
| 多个消费者群体 | 不可(取出一次就结束了) | 按群体独立消费 |
| 处理失败回收 | 直接实现 | XPENDING / XCLAIM |
| 内存 | 小 | 大(需要保留历史记录,需要MAXLEN) |
在现场相遇的样子
顺序能保证到什么程度呢?如果是一次性列表,单个消费者的话是FIFO。增加消费者的瞬间,顺序保证就会消失。因为两个消费者分别拿走消息A和B,不知道哪个先结束。需要顺序的大多是特定实体单位,所以用实体ID对队列进行分片,每片中只聚集一个消费者是实用的解决方案。
还有在使用流媒体的时候MAXLEN不能忘记。即使是XACK,流也会留下入口本身。XADD q:orders MAXLEN ~ 100000 * ...如果不设置上限,就会不断消耗内存。波浪号(~)通过近似修剪更小。
获得传递保障的三块
“不丢失信息”不仅仅是一个设定,而是三个方面都要具备。
생산자 → [브로커] → 소비자
①확인 ②지속성 ③확인 후 삭제
**①确认生产者。**经纪人没有收到发行呼叫回来的消息。经纪人的 需要等待确认(ack)。没有确认发送的话(fire-and-forget)经纪人死亡的那一刻 消失了。
**②经纪商持续性。**只要在内存中,在重新启动时就会消失。写在磁盘上,
如果可能的话,也确认一下副本。Kafka的acks=all,RabbitMQ的durable+
persistent这是这个。
**③消费者处理后确认。**一拿出来就删除的话,在处理过程中死亡时就会消失。
处理结束后必须确认。Redis List的BRPOP这个危险的原因是
这就是这个。
如果三个分支机构中任何一个缺失,即使其余的分支机构再坚固,也会失去。
没有确切的一次
在分布式系统中,无法获得“准确地传递一次”。可以获得的是 至少一次传递+优先处理,这个组合看起来正好是一次。
브로커가 "정확히 한 번" 을 광고하더라도
→ 그것은 브로커 안에서의 이야기다
→ 소비자가 처리하고 확인하기 전에 죽으면 다시 받는다
所以让消费者感到困惑是唯一的答案。处理过的消息ID 记录下来,如果同样的东西来了就跳过。
insert into processed(msg_id) values (%s) on conflict do nothing
-- 삽입된 행이 0이면 이미 처리한 것 → 건너뛴다
这个记录也会无限增长,所以设置了保存期限。最多可以重试的期限 比(一般几天)长一点。
只在需要顺序的地方遵守顺序。
为了遵守全局顺序,必须放弃并行处理。大部分情况下,如果按键单位顺序的话 足够了。
파티션 키 = 사용자 ID
→ 같은 사용자의 이벤트는 같은 파티션 → 순서 보장
→ 다른 사용자끼리는 병렬 처리
如果选择错误的键,就会集中在一个分区(hot partition)。先确认值的分布, 如果一个值占总值的很大比例,则不按其值除以总值。
下次实习要做的事情
用列表创建队列,确认FIFO,再复制丢失点,LMOVE然后,转到流媒体,直接处理消费者群体和PEL。最后自己写一张比较两者的表格。