LabHub
学习 学习路径 课程

队列与异步 API

list 队列丢消息的位置

在 LabHub 中继续学习

一句话总结

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。最后自己写一张比较两者的表格。