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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 10:08:28