Запуск несколько ассинхронных функций вместе на python
День добрый! Задача состоит в том, чтобы запустить несколько асинхронных функции разом. Сейчас у меня только две функции: первая открывает сокет и непрерывно получает данные из этого соединения; вторая периодически делает запросы к другому серверу и отправляет в тг канал ответ на сделанный запрос и данные, которые мы получали в первой функции. Проблема в том, что, возможно, нужно будет открывать еще сокеты и писать функции подобные второй. Можно ли выполнять асинхронные функции пареллельно? Пробовал совместить asyncio c multiprocessing и threading. Путного ничего не вышло у меня. Вот что сейчас написал
import asyncio, json, time, aiohttp, websockets async def func2(data): async with aiohttp.ClientSession() as session: async with session.post(url='. ', json=data) as response: response = await response.json() # Обрабатываю response async def func1(): url = f'. ' start_time = time.time() async with websockets.connect(url) as client: while True: data = json.loads(await client.recv()) # Обрабатываю data if time.time() - start_time > 5*60: await gen_tasks() start_time = time.time() async def gen_tasks(): tasks = [] for i in . data = tasks.append(asyncio.create_task(func2(data))) for task in tasks: await task asyncio.run(func1())
Отслеживать
задан 18 июл 2022 в 22:22
15 5 5 бронзовых знаков
И в чем проблема с кодом из вопроса? С точки зрения асинхронности в нем все выглядит правильно.
19 июл 2022 в 5:17
Основной вопрос в том, как запускать паралельно асинхронные функции
19 июл 2022 в 6:17
Ну так вот же вы запустили целую кучу асинхронных функций параллельно tasks.append(asyncio.create_task(func2(data))) . asyncio.create_task(f(data)) запускает f параллельно
19 июл 2022 в 6:18
В принципе в асинхронной функции можно собирать корутины просто в список без create_task: tasks.append(func2(data)) , а потом запускать через await asyncio.gather(tasks) (если нужен результат) или await asyncio.wait(tasks, return_when=asyncio.ALL_COMPLETED) (если результат не важен, нужно просто дождаться пока все выполнится).
19 июл 2022 в 7:47
Создал отдельную функцию, в которой создаются две задачи, и вызвал ее вместо func1(). Вывод, который и ожидал получить. Всем спасибо за помощь!
Запустить асинхронно функцию
Логи и асинхрон — это уже попытки решить проблему.
Проблема: time.sleep(1) ждёт больше 1 секунды. Ради проверки включал это приложение и через пару секунд включал секундомер на телефоне. Где-то ко второй минуте секундмер с телефона догонял. Не особо знаю за асинхронное программирование, но я подумал, что код с выводом в терминал времени и повышения на единицу заставляет программу ещё ждать, поэтому попытылся сделать функцию асинхронной. Вдохновлялся шаблоном из документации Логи при этом выглядят примерно так:
DEBUG:root:1699353460.6091332 DEBUG:asyncio:Using proactor: IocpProactor DEBUG:root:1699353461.6302617 DEBUG:root:1.0211284160614014 # Разница между старт и стоп больше 1 секунды DEBUG:root:------------------ DEBUG:root:1699353461.6313007 DEBUG:asyncio:Using proactor: IocpProactor DEBUG:root:1699353462.666385 DEBUG:root:1.0350842475891113 # Разница между старт и стоп больше 1 секунды DEBUG:root:------------------
Я почитал разные статьи, в том числе вопросы других людей с SO (Один из примеров). В них говорится о том, что погрешность — это допустимо, но там говорится о погрешности ~0.00001 и меньше. Я попробовал запустить без всяких таймеров цикл:
while True: start = time.time() logging.debug(start) time.sleep(1) end = time.time() logging.debug(end) logging.debug(end - start) logging.debug('------------------')
И в этом случае цифры в логе уже были более точными
DEBUG:root:1699355102.4560356 DEBUG:root:1699355103.4560683 DEBUG:root:1.000032663345337 DEBUG:root:------------------ DEBUG:root:1699355103.457032 DEBUG:root:1699355104.4571276 DEBUG:root:1.0000956058502197 DEBUG:root:------------------
Без логирования таймер всё равно чуть медленнее идёт.
Не могу понять, что конкретно тормозит мою программу?
Отслеживать
47.8k 17 17 золотых знаков 56 56 серебряных знаков 100 100 бронзовых знаков
задан 7 ноя в 11:18
Dark Space Dark Space
814 2 2 серебряных знака 17 17 бронзовых знаков
А никто не обещает наносекундную точность, тем более на питоне. все sleep работают примерно указанное количество времени. Операционная система не обязана отдавать управление процессу как только подошло время окончания sleep. Когда время окончания подходит ОС только отмечает процесс в своей таблице как готовый к выполнению. Позже планировщик, когда у него есть свободный квант времени запустит процесс. А после запуска еще питон будет выполнять какую нибудь внутреннюю работу, пока управление дойдет до time.time
7 ноя в 11:25
Хотите большей точности — делайте цикл и проверяйте, сколько времени прошло. А сон внутри цикла делайте маленькими промежутками, такими, какую хотите допустимую погрешность (и даже меньше). Либо какой-нибудь шедулер готовый возьмите, который всё это более оптимально делает под капотом. Синхронно или асинхронно запускаться — это не про точность вообще.
7 ноя в 11:39
Таймер вообще я бы проектировал по-другому..
7 ноя в 12:04
Тут не понятно, зачем из цикла вызывать асинхронную процедуру. С тем же успехом можно было бы просто print в том же самом цикле вызывать. И накладные расходы на вызов меньше были бы — у вас сейчас фактически при каждом вызове асинхронной функции заново новый асинхронный event loop создается. Если хочется асинхронности, то и вечный цикл имеет смысл внутрь асинхронной функции вынести.
7 ноя в 12:32
Всем спасибо за помощь
7 ноя в 12:52
2 ответа 2
Сортировка: Сброс на вариант по умолчанию
Вашу программу тормозит вызов os.system(«cls») в первом случае и logging.debug во втором.
Отслеживать
68.2k 5 5 золотых знаков 20 20 серебряных знаков 51 51 бронзовый знак
ответ дан 7 ноя в 11:40
33.6k 3 3 золотых знака 26 26 серебряных знаков 60 60 бронзовых знаков
В первом случае еще на каждой итерации новый event loop создается и убивается при вызове asyncio.run (docs.python.org/3/library/asyncio-runner.html#asyncio.run: If loop_factory is not None, it is used to create a new event loop; otherwise asyncio.new_event_loop() is used. The loop is closed at the end. ), тоже дополнительные накладные расходы.
7 ноя в 12:43
cls я думаю кушает 90% времени.
7 ноя в 14:39
time.sleep() добавляет задержку в вашем коде, но сам ваш код выполняется не нулевое время, из-за этого фактическое время выполнения одной итерации будет больше, чем указано в параметре sleep, и таймер постепенно будет отставать от реального времени.
Даже простой замер времени sleep
import time t = time.time() time.sleep(1.0) dt = time.time() - t print(dt)
покажет фактическую задержку больше 1 секунды (у меня показало 1.000610113143921 ). При большем количестве кода между замерами (и при наличии более медленного кода, например операций ввода-вывода, в том числе cls) — дополнительная задержка будет больше, и в цикле отставание будет постепенно увеличиваться.
Плюс операционная система не гарантирует абсолютную точность задержки, какой-то процесс может занять ядро надолго, возвращение из «сна» может произойти не в запланированное время.
Если предположить, что фоновых тяжелых процессов нет, можно пересчитывать задержку с учетом фактического времени выполнения итерации. Ну и выводить фактическое прошедшее время (по данным таймера компьютера), а не вручную посчитанное. Прототип таймера с пересчетом задержки, на каждую секунду должно выполняться 10 итераций:
import time target_dt = 0.1 # Целевой промежуток времени между итерациями (десять итераций в секунду) dt = target_dt # Целевое значение задержки между итерациями берем как начальное t = time.time() prev_time = t while True: print(f", ") time.sleep(dt) # time.sleep(0.01) # Даже если специально добавить дополнительную задержку, dt пересчитается так, чтобы ее учитывать current_time = time.time() # Фактический промежуток времени между итерациями dt_fact = current_time - prev_time # Пересчитываем задержку с поправкой на разницу между фактической задержкой и требуемой dt = target_dt - (dt_fact - dt) prev_time = current_time
0.000, 0.1 0.100, 0.09979171752929689 0.200, 0.09970550537109377 0.300, 0.099810266494751 0.400, 0.09973764419555667 0.500, 0.0997503757476807 0.600, 0.09978647232055668 0.700, 0.09971432685852055 0.800, 0.09980025291442876 0.900, 0.09976506233215338 1.000, 0.09972820281982428 1.100, 0.09979052543640143 1.200, 0.09976487159729011 1.300, 0.09973921775817879 1.400, 0.09974193572998055 1.500, 0.0997568130493165 1.600, 0.09974713325500498 1.700, 0.09976415634155283 1.800, 0.09980382919311534 1.900, 0.09972858428955089 2.000, 0.09971747398376477 2.100, 0.09978790283203137 2.200, 0.09974532127380384 2.300, 0.09975781440734877 2.400, 0.09979414939880385 2.500, 0.09974107742309585 2.600, 0.09973020553588882 2.700, 0.09979276657104508 2.800, 0.09974160194396989 2.900, 0.09979963302612321 3.000, 0.09978470802307146 3.100, 0.09977574348449725 3.200, 0.09972982406616229 3.300, 0.0997399330139162 3.400, 0.09973549842834492 3.500, 0.09977755546569844 3.600, 0.09981722831726095 3.700, 0.09977011680603048 3.800, 0.09973802566528342 3.900, 0.09972071647644065 4.000, 0.09967408180236839 .
Видно, что промежуток с точностью около +-0.001 между итерациям сохраняется.
Тестировалось на Linux, на Windows результат может быть другим (вроде бы там нельзя сделать задержку с точностью лучше 0.01 секунды, из-за этого точность таймера может быть хуже).
Помещать таймер в async функцию я не вижу особого смысла. Если в этом же эвент лупе будут выполняться другие функции, это будет добавлять дополнительную случайную погрешность. Если нужно — запускайте его в отдельном потоке.
Запуск асинхронных функций
Скажите, пожалуйста, каким правильно использовать запуска функций в асинхронном варианте? Есть список с вложенными списками. Где каждый список представляет из себя будущие значения для функции.
c = [[121, 'yes', 5],[345, 'no', 1]] sphere = 121 dart = 'yes' number = 5
Есть функция:
def test(sphere, dart, nubmer)
Я бы хотел совершить следующие действия
for element in c: sphere, dart, number = c[0], c[1], c[2] #дальше мне надо запускать функцию в асинхронном режиме. #Как мне сделать так, чтобы она сама завершалась? Это нужно прописывать в самой функции? #Или же надо останавливать поток непосредственно после каждого элемента итерации?
Введение в асинхронное программирование на Python
Асинхронное программирование – это вид параллельного программирования, в котором какая-либо единица работы может выполняться отдельно от основного потока выполнения приложения. Когда работа завершается, основной поток получает уведомление о завершении рабочего потока или произошедшей ошибке. У такого подхода есть множество преимуществ, таких как повышение производительности приложений и повышение скорости отклика.

В последние несколько лет асинхронное программирование привлекло к себе пристальное внимание, и на то есть причины. Несмотря на то, что этот вид программирования может быть сложнее традиционного последовательного выполнения, он гораздо более эффективен.
Например, вместо того, что ждать завершения HTTP-запроса перед продолжением выполнения, вы можете отправить запрос и выполнить другую работу, которая ждет своей очереди, с помощью асинхронных корутин в Python.
Асинхронность – это одна из основных причин популярности выбора Node.js для реализации бэкенда. Большое количество кода, который мы пишем, особенно в приложениях с тяжелым вводом-выводом, таком как на веб-сайтах, зависит от внешних ресурсов. В нем может оказаться все, что угодно, от удаленного вызова базы данных до POST-запросов в REST-сервис. Как только вы отправите запрос в один из этих ресурсов, ваш код будет просто ожидать ответа. С асинхронным программированием вы позволяете своему коду обрабатывать другие задачи, пока ждете ответа от ресурсов.
Как Python умудряется делать несколько вещей одновременно?

1. Множественные процессы
Самый очевидный способ – это использование нескольких процессов. Из терминала вы можете запустить свой скрипт два, три, четыре, десять раз, и все скрипты будут выполняться независимо и одновременно. Операционная система сама позаботится о распределении ресурсов процессора между всеми экземплярами. В качестве альтернативы вы можете воспользоваться библиотекой multiprocessing, которая умеет порождать несколько процессов, как показано в примере ниже.
from multiprocessing import Process def print_func(continent='Asia'): print('The name of continent is : ', continent) if __name__ == "__main__": # confirms that the code is under main function names = ['America', 'Europe', 'Africa'] procs = [] proc = Process(target=print_func) # instantiating without any argument procs.append(proc) proc.start() # instantiating process with arguments for name in names: # print(name) proc = Process(target=print_func, args=(name,)) procs.append(proc) proc.start() # complete the processes for proc in procs: proc.join()
The name of continent is : Asia The name of continent is : America The name of continent is : Europe The name of continent is : Africa
2. Множественные потоки
Еще один способ запустить несколько работ параллельно – это использовать потоки. Поток – это очередь выполнения, которая очень похожа на процесс, однако в одном процессе вы можете иметь несколько потоков, и у всех них будет общий доступ к ресурсам. Однако из-за этого написать код потока будет сложно. Аналогично, все тяжелую работу по выделению памяти процессора сделает операционная система, но глобальная блокировка интерпретатора (GIL) позволит только одному потоку Python запускаться в одну единицу времени, даже если у вас есть многопоточный код. Так GIL на CPython предотвращает многоядерную конкурентность. То есть вы насильно можете запуститься только на одном ядре, даже если у вас их два, четыре или больше.
import threading def print_cube(num): """ function to print cube of given num """ print("Cube: <>".format(num * num * num)) def print_square(num): """ function to print square of given num """ print("Square: <>".format(num * num)) if __name__ == "__main__": # creating thread t1 = threading.Thread(target=print_square, args=(10,)) t2 = threading.Thread(target=print_cube, args=(10,)) # starting thread 1 t1.start() # starting thread 2 t2.start() # wait until thread 1 is completely executed t1.join() # wait until thread 2 is completely executed t2.join() # both threads completely executed print("Done!")
Square: 100 Cube: 1000 Done!
3. Корутины и yield :
Корутины – это обобщение подпрограмм. Они используются для кооперативной многозадачности, когда процесс добровольно отдает контроль ( yield ) с какой-то периодичностью или в периоды ожидания, чтобы позволить нескольким приложениям работать одновременно. Корутины похожи на генераторы, но с дополнительными методами и небольшими изменениями в том, как мы используем оператор yield. Генераторы производят данные для итерации, в то время как корутины могут еще и потреблять данные.
def print_name(prefix): print("Searching prefix:<>".format(prefix)) try : while True: # yeild used to create coroutine name = (yield) if prefix in name: print(name) except GeneratorExit: print("Closing coroutine!!") corou = print_name("Dear") corou.__next__() corou.send("James") corou.send("Dear James") corou.close()
Searching prefix:Dear Dear James Closing coroutine!!
4. Асинхронное программирование
Четвертый способ – это асинхронное программирование, в котором не участвует операционная система. Со стороны операционной системы у вас останется один процесс, в котором будет всего один поток, но вы все еще сможете выполнять одновременно несколько задач. Так в чем тут фокус?
Asyncio – модуль асинхронного программирования, который был представлен в Python 3.4. Он предназначен для использования корутин и future для упрощения написания асинхронного кода и делает его почти таким же читаемым, как синхронный код, из-за отсутствия callback-ов.
Asyncio использует разные конструкции: event loop , корутины и future .
- event loop управляет и распределяет выполнение различных задач. Он регистрирует их и обрабатывает распределение потока управления между ними.
- Корутины (о которых мы говорили выше) – это специальные функции, работа которых схожа с работой генераторов в Python, с помощью await они возвращают поток управления обратно в event loop. Запуск корутины должен быть запланирован в event loop. Запланированные корутины будут обернуты в Tasks, что является типом Future.
- Future отражает результат таска, который может или не может быть выполнен. Результатом может быть exception.
Переключение контекста в asyncio представляет собой event loop , который передает поток управления от одной корутины к другой.
В следующем примере, мы запускаем 3 асинхронных таска, которые по-отдельности делают запросы к Reddit, извлекают и выводят содержимое JSON. Мы используем aiohttp – клиентскую библиотеку http, которая гарантирует, что даже HTTP-запрос будет выполнен асинхронно.
import signal import sys import asyncio import aiohttp import json loop = asyncio.get_event_loop() client = aiohttp.ClientSession(loop=loop) async def get_json(client, url): async with client.get(url) as response: assert response.status == 200 return await response.read() async def get_reddit_top(subreddit, client): data1 = await get_json(client, 'https://www.reddit.com/r/' + subreddit + '/top.json?sort=top&t=day&limit=5') j = json.loads(data1.decode('utf-8')) for i in j['data']['children']: score = i['data']['score'] title = i['data']['title'] link = i['data']['url'] print(str(score) + ': ' + title + ' (' + link + ')') print('DONE:', subreddit + '\n') def signal_handler(signal, frame): loop.stop() client.close() sys.exit(0) signal.signal(signal.SIGINT, signal_handler) asyncio.ensure_future(get_reddit_top('python', client)) asyncio.ensure_future(get_reddit_top('programming', client)) asyncio.ensure_future(get_reddit_top('compsci', client)) loop.run_forever()
50: Undershoot: Parsing theory in 1965 (http://jeffreykegler.github.io/Ocean-of-Awareness-blog/individual/2018/07/knuth_1965_2.html) 12: Question about best-prefix/failure function/primal match table in kmp algorithm (https://www.reddit.com/r/compsci/comments/8xd3m2/question_about_bestprefixfailure_functionprimal/) 1: Question regarding calculating the probability of failure of a RAID system (https://www.reddit.com/r/compsci/comments/8xbkk2/question_regarding_calculating_the_probability_of/) DONE: compsci 336: /r/thanosdidnothingwrong -- banning people with python (https://clips.twitch.tv/AstutePluckyCocoaLitty) 175: PythonRobotics: Python sample codes for robotics algorithms (https://atsushisakai.github.io/PythonRobotics/) 23: Python and Flask Tutorial in VS Code (https://code.visualstudio.com/docs/python/tutorial-flask) 17: Started a new blog on Celery - what would you like to read about? (https://www.python-celery.com) 14: A Simple Anomaly Detection Algorithm in Python (https://medium.com/@mathmare_/pyng-a-simple-anomaly-detection-algorithm-2f355d7dc054) DONE: python 1360: git bundle (https://dev.to/gabeguz/git-bundle-2l5o) 1191: Which hashing algorithm is best for uniqueness and speed? Ian Boyd's answer (top voted) is one of the best comments I've seen on Stackexchange. (https://softwareengineering.stackexchange.com/questions/49550/which-hashing-algorithm-is-best-for-uniqueness-and-speed) 430: ARM launches “Facts” campaign against RISC-V (https://riscv-basics.com/) 244: Choice of search engine on Android nuked by “Anonymous Coward” (2009) (https://android.googlesource.com/platform/packages/apps/GlobalSearch/+/592150ac00086400415afe936d96f04d3be3ba0c) 209: Exploiting freely accessible WhatsApp data or “Why does WhatsApp web know my phone’s battery level?” (https://medium.com/@juan_cortes/exploiting-freely-accessible-whatsapp-data-or-why-does-whatsapp-know-my-battery-level-ddac224041b4) DONE: programming
Использование Redis и Redis Queue RQ
Использование asyncio и aiohttp не всегда хорошая идея, особенно если вы пользуетесь более старыми версиями Python. К тому же, бывают моменты, когда вам нужно распределить задачи по разным серверам. В этом случае можно использовать RQ (Redis Queue). Это обычная библиотека Python для добавления работ в очередь и обработки их воркерами в фоновом режиме. Для организации очереди используется Redis – база данных ключей/значений.
В примере ниже мы добавили в очередь простую функцию count_words_at_url с помощью Redis.
from mymodule import count_words_at_url from redis import Redis from rq import Queue q = Queue(connection=Redis()) job = q.enqueue(count_words_at_url, 'http://nvie.com') ******mymodule.py****** import requests def count_words_at_url(url): """Just an example function that's called async.""" resp = requests.get(url) print( len(resp.text.split())) return( len(resp.text.split()))
15:10:45 RQ worker 'rq:worker:EMPID18030.9865' started, version 0.11.0 15:10:45 *** Listening on default. 15:10:45 Cleaning registries for queue: default 15:10:50 default: mymodule.count_words_at_url('http://nvie.com') (a2b7451e-731f-4f31-9232-2b7e3549051f) 322 15:10:51 default: Job OK (a2b7451e-731f-4f31-9232-2b7e3549051f) 15:10:51 Result is kept for 500 seconds
Заключение
В качестве примера возьмем шахматную выставку, где один из лучших шахматистов соревнуется с большим количеством людей. У нас есть 24 игры и 24 человека, с которыми можно сыграть, и, если шахматист будет играть с ними синхронно, это займет не менее 12 часов (при условии, что средняя игра занимает 30 ходов, шахматист продумывает ход в течение 5 секунд, а противник – примерно 55 секунд.) Однако в асинхронном режиме шахматист сможет делать ход и оставлять противнику время на раздумья, тем временем переходя к следующему противнику и деля ход. Таким образом, сделать ход во всех 24 играх можно за 2 минуты, и выиграны они все могут быть всего за один час.
Это и подразумевается, когда говорят о том, что асинхронность ускоряет работу. О такой быстроте идет речь. Хороший шахматист не начинает играть в шахматы быстрее, просто время более оптимизировано, и оно не тратится впустую на ожидание. Так это работает.
По этой аналогии шахматист будет процессором, а основная идея будет заключаться в том, чтобы процессор простаивал как можно меньше времени. Речь о том, чтобы у него всегда было занятие.
На практике асинхронность определяется как стиль параллельного программирования, в котором одни задачи освобождают процессор в периоды ожидания, чтобы другие задачи могли им воспользоваться. В Python есть несколько способов достижения параллелизма, отвечающих вашим требованиям, потоку кода, обработке данных, архитектуре и вариантам использования, и вы можете выбрать любой из них.
