You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用aioredis异步处理Redis消息?解决串行处理问题

解决Redis PubSub消息异步并行处理问题

你的问题核心在于reader函数在处理消息时被阻塞:原代码中await process_message(data)会让reader暂停执行,直到当前消息处理完成,才能继续调用get_message获取下一条消息,导致所有消息串行处理。

修改方案:将消息处理任务异步提交

只需要修改reader函数,把消息处理逻辑包装成独立的异步任务提交给事件循环,而不是直接await,这样reader可以立刻回到循环中继续接收新消息:

import asyncio
import aioredis


async def reader(channel: aioredis.client.PubSub):
    while True:
        data = None
        try:
            message = await channel.get_message(ignore_subscribe_messages=True)
            if message is not None:
                print(f"(Reader) Message Received: {message}")
                data = message["data"]
        except asyncio.TimeoutError:
            pass

        if data is not None:
            # 关键改动:创建独立任务执行消息处理,不阻塞reader
            asyncio.create_task(process_message(data))


async def process_message(message):
    print(f"start process {message=}")
    await asyncio.sleep(10)
    print(f"+processed {message=}")


async def publish(redis, channel, message):
    print(f"-->publish {message=} to {channel=}")
    result = await redis.publish(channel, message)
    print("     +published")
    return result


async def main():
    redis = aioredis.from_url("redis://localhost")
    pubsub = redis.pubsub()
    await pubsub.subscribe("channel:1", "channel:2")

    future = asyncio.create_task(reader(pubsub))

    await publish(redis, "channel:1", "Hello")
    await publish(redis, "channel:2", "World")

    await future


if __name__ == "__main__":
    asyncio.run(main())

优化:控制并发处理数量

如果消息量极大,无限制创建任务可能导致资源耗尽,可以用asyncio.Semaphore来限制同时运行的处理任务数:

# 在main中初始化信号量,传递给reader
async def main():
    redis = aioredis.from_url("redis://localhost")
    pubsub = redis.pubsub()
    await pubsub.subscribe("channel:1", "channel:2")
    
    # 限制最多同时处理3条消息
    semaphore = asyncio.Semaphore(3)
    future = asyncio.create_task(reader(pubsub, semaphore))

    await publish(redis, "channel:1", "Hello")
    await publish(redis, "channel:2", "World")
    # 可以多发布几条测试并发
    await publish(redis, "channel:1", "Msg3")
    await publish(redis, "channel:2", "Msg4")

    await future

# 修改reader和process_message,使用信号量
async def reader(channel: aioredis.client.PubSub, semaphore: asyncio.Semaphore):
    while True:
        data = None
        try:
            message = await channel.get_message(ignore_subscribe_messages=True)
            if message is not None:
                print(f"(Reader) Message Received: {message}")
                data = message["data"]
        except asyncio.TimeoutError:
            pass

        if data is not None:
            asyncio.create_task(process_message(data, semaphore))


async def process_message(message, semaphore: asyncio.Semaphore):
    async with semaphore:
        print(f"start process {message=}")
        await asyncio.sleep(10)
        print(f"+processed {message=}")

这样既保证了消息能并行处理,又能避免并发过高导致的系统问题。

内容的提问来源于stack exchange,提问作者SKulibin

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.03 04:55:16