Django Celery+Asyncio周期性错误:loop参数与Future不匹配
排查Celery Beat周期性任务中的asyncio循环不匹配错误
问题背景
我正在用Django + Celery + Celery Beat运行每分钟一次的SNMP数据采集任务,代码里用了asyncio,还加了事件循环关闭检查和重建逻辑,但每隔大约3分钟就会出现一次失败,其余时间任务正常运行。错误信息如下:
Traceback (most recent call last): File "/usr/local/lib/python3.6/site-packages/celery/app/trace.py", line 374, in trace_task R = retval = fun(*args, **kwargs) File "/usr/local/lib/python3.6/site-packages/celery/app/trace.py", line 629, in __protected_call__ return self.run(*args, **kwargs) File "/itapp/itapp/monitoring/tasks.py", line 32, in link_data return get_link_data() File "/itapp/itapp/monitoring/jobs/link_monitoring.py", line 209, in get_link_data done, pending = loop.run_until_complete(asyncio.wait(tasks)) File "/usr/local/lib/python3.6/asyncio/base_events.py", line 468, in run_until_complete return future.result() File "/usr/local/lib/python3.6/asyncio/tasks.py", line 311, in wait fs = {ensure_future(f, loop=loop) for f in set(fs)} File "/usr/local/lib/python3.6/asyncio/tasks.py", line 311, in <setcomp> fs = {ensure_future(f, loop=loop) for f in set(fs)} File "/usr/local/lib/python3.6/asyncio/tasks.py", line 514, in ensure_future raise ValueError('loop argument must agree with Future') ValueError: loop argument must agree with Future
相关代码片段
异步采集函数
async def retrieve_data(link): poll_interval = 60 results = [] # credentials: link_mgmt_ip = link.mgmt_ip link_index = link.interface_index snmp_user = link.device_circuit_subnet.device.snmp_data.name snmp_auth = link.device_circuit_subnet.device.snmp_data.auth snmp_priv = link.device_circuit_subnet.device.snmp_data.priv hostname = link.device_circuit_subnet.device.hostname print('polling data for {} on {}'.format(hostname,link_mgmt_ip)) # first poll for speeds download_speed_data_poll1 = snmp_get(link_mgmt_ip, down_speed_oid % link_index ,snmp_user, snmp_auth, snmp_priv) # check we were able to poll if 'timeout' in str(get_snmp_value(download_speed_data_poll1)).lower(): return 'timeout trying to poll {} - {}'.format(hostname ,link_mgmt_ip) upload_speed_data_poll1 = snmp_get(link_mgmt_ip, up_speed_oid % link_index, snmp_user, snmp_auth, snmp_priv) # wait for poll interval await asyncio.sleep(poll_interval) # second poll for speeds download_speed_data_poll2 = snmp_get(link_mgmt_ip, down_speed_oid % link_index, snmp_user, snmp_auth, snmp_priv) upload_speed_data_poll2 = snmp_get(link_mgmt_ip, up_speed_oid % link_index, snmp_user, snmp_auth, snmp_priv) # create deltas for speed down_delta = int(get_snmp_value(download_speed_data_poll2)) - int(get_snmp_value(download_speed_data_poll1)) up_delta = int(get_snmp_value(upload_speed_data_poll2)) - int(get_snmp_value(upload_speed_data_poll1)) # set speed results download_speed = round((down_delta * 8 / poll_interval) / 1048576) upload_speed = round((up_delta * 8 / poll_interval) / 1048576) # get description and interface state int_desc = snmp_get(link_mgmt_ip, int_desc_oid % link_index, snmp_user, snmp_auth, snmp_priv) int_state = snmp_get(link_mgmt_ip, int_state_oid % link_index, snmp_user, snmp_auth, snmp_priv) ... return results
任务入口函数
def get_link_data(): mgmt_ip = Subquery( DeviceCircuitSubnets.objects.filter(device_id=OuterRef('device_circuit_subnet__device_id'),subnet__subnet_type__poll=True).values('subnet__subnet')[:1]) link_data = LinkTargets.objects.all() \ .select_related('device_circuit_subnet') \ .select_related('device_circuit_subnet__device') \ .select_related('device_circuit_subnet__device__snmp_data') \ .select_related('device_circuit_subnet__subnet') \ .select_related('device_circuit_subnet__circuit') \ .annotate(mgmt_ip=mgmt_ip) tasks = [] loop = asyncio.get_event_loop() if asyncio.get_event_loop().is_closed(): loop = asyncio.new_event_loop() asyncio.set_event_loop(asyncio.new_event_loop()) for link in link_data: tasks.append(asyncio.ensure_future(retrieve_data(link))) if tasks: start = time.time() done, pending = loop.run_until_complete(asyncio.wait(tasks)) loop.close() results = [] for completed_task in done: results.append(completed_task.result()[0]) end = time.time() print("Poll time: {}".format(end - start)) return 'Link data updated for {}'.format(' \n '.join(results)) else: return 'no tasks defined'
错误原因分析
最核心的问题出在事件循环的初始化逻辑上:
在get_link_data里,当检测到当前事件循环已关闭时,你创建了两个完全独立的新循环:
loop = asyncio.new_event_loop() asyncio.set_event_loop(asyncio.new_event_loop())
- 第一行创建了一个新循环并赋值给
loop变量 - 第二行又创建了另一个全新的循环,并设置为当前全局循环
这就导致:
- 你用
asyncio.ensure_future(retrieve_data(link))创建任务时,任务会绑定到当前全局循环(也就是第二行创建的那个) - 但你调用
loop.run_until_complete(...)时,用的是第一行赋值的loop变量(另一个循环) - 两个循环不匹配,就触发了
ValueError: loop argument must agree with Future
另外还有一个潜在问题:retrieve_data里用了await asyncio.sleep(60),而Celery Beat是每分钟触发一次任务,这会导致前一个任务还在执行(sleep中),新的任务就已经启动,多个任务实例共享或竞争事件循环,进一步加剧了循环混乱的概率,这也是错误间歇性出现的原因之一。
解决方法
1. 修复事件循环初始化逻辑
把循环初始化的代码改成这样,确保只创建一个新循环,并且同时赋值给变量和设置为当前全局循环:
def get_link_data(): # ... 其他查询代码保持不变 ... tasks = [] # 正确初始化事件循环 try: loop = asyncio.get_event_loop() if loop.is_closed(): raise RuntimeError("Event loop is closed") except RuntimeError: loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) for link in link_data: # 现在ensure_future会用和run_until_complete相同的循环 tasks.append(asyncio.ensure_future(retrieve_data(link))) if tasks: start = time.time() done, pending = loop.run_until_complete(asyncio.wait(tasks)) loop.close() # ... 后续结果处理代码保持不变 ...
或者更简洁的写法:
loop = asyncio.get_event_loop() if loop.is_closed(): loop = asyncio.new_event_loop() asyncio.set_event_loop(loop)
这样就保证了创建任务和运行任务用的是同一个事件循环,彻底解决循环不匹配的问题。
2. 调整任务调度与采集逻辑
retrieve_data里的await asyncio.sleep(60)会让任务至少运行60秒,而Celery Beat每分钟触发一次,必然会导致任务重叠。建议重构采集逻辑:
- 把两次SNMP采集的间隔逻辑从任务内部移出去:每次任务只采集一次数据,存储到数据库或缓存中,下次任务执行时读取上一次的采集值来计算差值
- 这样任务执行时间会大幅缩短,避免任务重叠,也能让Celery Beat的调度更稳定
3. 替换同步SNMP调用为异步实现
目前的snmp_get看起来是同步函数,在asyncio事件循环里调用会阻塞整个循环,导致并发效率极低。建议改用异步SNMP库(比如aiosnmp),这样多个采集任务可以真正并发执行,减少任务总耗时。
内容的提问来源于stack exchange,提问作者AlexW
相关产品推荐
相关产品推荐

