Я использую сельдерей Джанго и ритм сельдерея для выполнения периодических заданий. Я запускаю задачу каждую минуту, чтобы получить данные через SNMP.
Моя функция использует asyncio, как показано ниже. Я поставил проверку в коде, чтобы проверить, замкнут ли цикл, и создать новый.
но, похоже, что происходит через каждые несколько задач, я получаю сбой, и в db Dango-tasks-results у меня есть приведенная ниже трассировка. каждые 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'