Аргумент django celery и asyncio - loop должен согласовываться с Future примерно каждые 3 минуты - PullRequest
0 голосов
/ 15 мая 2018

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

1 Ответ

0 голосов
/ 30 мая 2018

из этих URL, предложенных пользователем 4815162342

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

Когда использовать и когда не использовать Python 3.5 `await`?

При выполнении функций ascync любая операция ввода-вывода должна быть асинхронной, за исключением функций, которые выполняются в памяти.(в моем примере запрос регулярного выражения)

, т. е. любая функция, которой требуется собрать данные из другого источника (запрос django в моем примере), которая не является асинхронной, должна выполняться в исполнителе.

Я думаю, что теперь я исправил свои проблемы, выполнив все вызовы django DB в исполнителях, с тех пор у меня не было проблем с запуском сценария ad hoc.

Однако у меня есть проблема совместимости с celery и async (поскольку celery еще не совместим с asyncio, который выдает некоторые ошибки, но не те ошибки, которые я видел ранее)

...