【发布时间】: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 tutorial 和 an answer 以了解更多信息。
-
我至少读了这三遍,这对新手来说相当混乱,但据我了解,我认为我的脚本没问题,当它起作用时,我假设是这样。我不确定我在那里有任何阻塞操作。阻塞操作需要超过 50 毫秒的时间吗?那正确吗?如果我有其中之一,它应该在执行程序中运行?
-
阻塞操作是任何可能等待任意时间的操作,例如任何取决于网络延迟或服务器响应的东西。如果您不使用与 asyncio 兼容的库并且您不
await您的阻塞调用,那么使用 asyncio 不会给您带来任何好处。您没有并行性,并且您也遇到了异常,可能是由于与您的某些库在内部使用的另一个事件循环冲突。尝试使用普通的def而不是async def代替retrieve_data并将await asyncio.sleep()替换为time.sleep()。
标签: python django celery python-asyncio