如何在Asyncio中避免上下文切换?aio-pika场景代码问询
问题解决:aio-pika消费回调中同步执行更新逻辑
问题分析
你遇到的核心问题是:_consume_callback触发_get_update_schedule时,异步上下文切换导致多个消息同时执行更新逻辑,引发并发冲突。之前尝试asyncio.Lock无效,主要是锁的使用位置错误,再加上代码存在几处基础逻辑错误,导致问题持续存在。
原代码的关键问题:
_receive方法缩进错误,被定义在__init__内部,无法被外部正确调用_check_match_in_schedule逻辑错误:遍历self._match_schedular.values()时,误将遍历到的value当作key去取值,触发KeyError,导致该方法始终返回False,频繁触发不必要的更新- 未在更新逻辑的关键路径上正确加锁,无法阻止并发执行
修复后的完整代码
import asyncio from datetime import datetime # 导入你的其他依赖:consumers, producers, fast_data, TRACER等 class GameMatchesWorker: def __init__( self, consumer: consumers.AbstractConsumer, producer: producers.AbstractProducer, provider_client: fast_data.FastDataClient, ) -> None: self._consumer = consumer self._producer = producer self._provider_client = provider_client self._matches = set() self._input_buffer = [] self._input_event = asyncio.Event() self._match_schedular = {} self._lock = asyncio.Lock() # 确保consumer的回调绑定到当前类的_consume_callback方法 # 示例:self._consumer.register_callback(self._consume_callback) async def _receive(self) -> dict: while not self._input_buffer: try: await asyncio.wait_for(self._input_event.wait(), timeout=10) self._input_event.clear() except asyncio.TimeoutError: # 处理超时场景,避免无限等待 return {} received_data = self._input_buffer.pop(0) if received_data.get("code") == fast_data.StatusCodes.SUCCESS: return {game_match_meta["game_id"]: game_match_meta for game_match_meta in received_data["list"]} return {} async def _get_update_schedule(self): # 加锁确保同一时间仅一个任务执行更新逻辑 async with self._lock: current_date = datetime.now() for source in range(1,6): await self._provider_client.emit_games_list( date_from=current_date, date_to=current_date, source=source ) matches_schedule = await self._receive() self._match_schedular[source] = matches_schedule def _check_match_in_schedule(self, match_id: int): # 修复遍历逻辑:直接检查每个数据源的game_id集合 for source_data in self._match_schedular.values(): if match_id in source_data: return True return False async def _consume_callback(self, game_match_message: consumers.GameMatchMessage) -> None: with TRACER.start_as_current_span("fastdata_workerhost_consume_message"): game_match_id = game_match_message.game_id # 先获取锁再检查匹配状态,避免多个请求同时触发更新 async with self._lock: if not self._check_match_in_schedule(game_match_id): await self._get_update_schedule() # 此处可添加匹配到game_id后的业务逻辑
关键修复点说明
- 修正
_receive方法缩进:将其从__init__内部移出,作为类的独立方法,确保能被正常调用。 - 修复匹配检查逻辑:直接遍历调度表的value集合,避免因错误使用key导致的
KeyError,确保匹配判断准确。 - 正确使用
asyncio.Lock:- 在
_get_update_schedule内部用async with self._lock包裹全量更新逻辑,强制同一时间仅一个任务执行更新。 - 在
_consume_callback中先获取锁再做匹配检查,避免多个请求同时进入“需要更新”分支,重复触发外部接口调用。
- 在
- 添加超时处理:给
_receive的等待逻辑添加TimeoutError捕获,避免因消息未到达导致的无限等待。
额外优化建议
- 给调度表添加更新时间戳,在
_get_update_schedule执行前判断是否在最近一段时间内已更新,减少不必要的外部接口调用。 - 确保填充
_input_buffer的逻辑会正确触发self._input_event.set(),否则_receive会一直处于等待状态。 - 若aio-pika的consumer配置了
prefetch_count>1,可适当降低该值,结合锁进一步控制消费并发度。
内容的提问来源于stack exchange,提问作者dovgalb
相关产品推荐
相关产品推荐

