使用 redis 和 python 进行故障安全消息广播,供特定接收者使用

fam*_*man 3 python ipc stream redis redis-py

所以redis 5.0新鲜推出了一个新功能,叫做Streams。它们似乎非常适合分发进程间通信的消息:

  • 它们在可靠性方面超越了 PUB/SUB 事件消息传递的能力:PUB/SUB 是即发即忘的,无法保证收件人会收到消息
  • redis 列表有点低级,但仍然可以使用。然而,流针对性能和上述用例进行了优化。

然而,由于这个功能相当新,几乎没有任何 Python(甚至通用的 redis)手册,而且我真的不知道如何使流系统适应我的用例。

我想要一个发布程序,将消息推送到流并包含收件人信息(如recipient: "user1")。然后我将有几个接收进程,所有进程都应该检查新的流消息并比较它们是否是目标接收者。如果是,他们应该处理该消息并将其标记为已处理(已确认)。

但是,我不太了解消费者群体、待处理状态等概念。有人能给我一个关于我的小伪代码的真实例子吗?

发件人.py

db = Redis(...)
db.the_stream.add({"recipient": "user1", "task": "be a python"})
Run Code Online (Sandbox Code Playgroud)

recipient.py(将有许多实例运行,每个实例都有一个唯一的收件人 ID)

recipient_id = "user1" # you get the idea...
db = Redis(...)
while True:
    message = db.the_stream.blocking_read("$") # "$" somehow means: just receive new messages
    if message.recipient == recipient_id:
        perform_task(message.task)
        message.acknowledge() # let the stream know it was processed
    else:
        pass # well, do nothing here since it's not our message. Another recipient instance should do the job.```

Run Code Online (Sandbox Code Playgroud)

sma*_*sey 5

根据您给出的示例和伪代码,让我们想象一下:

  • recipient.user1每分钟收到 60 条消息
  • 并且该perform_task()方法执行需要2秒。

这里发生的情况是显而易见的:新消息进入和处理之间的延迟只会随着时间的推移而增加,与“实时处理”的距离越来越远。

system throughput = 30 messages/minute

为了解决这个问题,您可能需要为user1. 在这里,您可以有 4 个不同的 python 进程并行运行,所有 4 个进程加入同一组中user1。现在,当 4 名工作人员之一收到一条消息时,user1他会接听它,然后perform_task()。

system throughput = 120 message/minute

在您的示例中,message.acknowledge()实际上并不存在,因为您的流读取器是单独的(XREAD 命令)。

如果它是一个组,则消息的确认变得至关重要,这就是 Redis 知道其中一个组成员实际上确实处理了该消息的方式,因此它可能会“继续前进”(它可能会忘记该消息正在等待确认) 。当您使用组时,会有一些服务器端逻辑来确保每条消息都一次传递给消费者组工作人员之一(XGROUPREAD 命令)。当客户端完成时,它会发出对该消息的确认(XACK 命令),以便服务器端“消费者组缓冲区”可以删除它并继续。

想象一下,如果一名工人去世了并且从未回复过该消息。通过消费者组,您可以留意这种情况(使用 XPENDING 命令)并通过例如重试在另一个消费者中处理相同的消息来采取行动。

当您不使用组时,redis 服务器不需要“继续”,“确认”变成 100% 客户端/业务逻辑。