如何用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
相关产品推荐
相关产品推荐

