【问题标题】:django celery and asyncio - loop argument must agree with Future approx every 3 minsdjango celery 和 asyncio - 循环参数必须大约每 3 分钟与 Future 一致
【发布时间】:2018-10-25 09:33:21
【问题描述】:

我正在使用 django celery 和 celery beat 来运行周期性任务。我每分钟运行一个任务,通过 SNMP 获取一些数据。

我的函数使用 asyncio,如下所示。我已经检查了代码以检查循环是否已关闭并创建一个新的。

但似乎发生的事情是每隔几个任务,我就会失败,并且在 Django-tasks-results db 中我有以下回溯。似乎每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'

【问题讨论】:

  • 这似乎与您already posted 然后删除的问题完全相同。这个设计似乎有缺陷,因为你的协程除了asyncio.sleep()之外没有任何await
  • 抱歉,我删除了原来的问题,因为脚本在经过一些调整后可以正常工作,我认为此时一切正常,可以正常工作 2/3 次。我是 Asyncio 的新手,所以我不太清楚你的意思?我需要该函数在 snmp 轮询之间等待 60 秒以收集数据
  • 我们已经在原始问题中进行了此对话,我回复说您应该等待任何可以阻止的操作。我还向您推荐了 a tutorialan answer 以了解更多信息。
  • 我至少读了这三遍,这对新手来说相当混乱,但据我了解,我认为我的脚本没问题,当它起作用时,我假设是这样。我不确定我在那里有任何阻塞操作。阻塞操作需要超过 50 毫秒的时间吗?那正确吗?如果我有其中之一,它应该在执行程序中运行?
  • 阻塞操作是任何可能等待任意时间的操作,例如任何取决于网络延迟或服务器响应的东西。如果您不使用与 asyncio 兼容的库并且您不 await 您的阻塞调用,那么使用 asyncio 不会给您带来任何好处。您没有并行性,并且您也遇到了异常,可能是由于与您的某些库在内部使用的另一个事件循环冲突。尝试使用普通的 def 而不是 async def 代替 retrieve_data 并将 await asyncio.sleep() 替换为 time.sleep()

标签: python django celery python-asyncio


【解决方案1】:

来自 user4815162342 建议的这些 url

https://medium.freecodecamp.org/a-guide-to-asynchronous-programming-in-python-with-asyncio-232e2afa44f6

When to use and when not to use Python 3.5 `await` ?

运行异步函数时,任何输入输出操作都需要异步兼容,但在内存中运行的函数除外。 (在我的例子中是一个正则表达式查询)

即任何需要从不兼容异步的其他来源(在我的示例中为 django 查询)收集数据的函数都必须在执行程序中运行。

我想我现在已经通过在执行程序中运行所有 django DB 调用来解决我的问题,从那以后我就没有遇到过运行临时脚本的问题。

但是我有 celery 和 async 的兼容性问题(因为 celery 还不兼容 asyncio 会引发一些错误,但不是我之前看到的错误)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2018-06-28
    • 1970-01-01
    • 2021-09-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-03-31
    相关资源
    最近更新 更多