Как выполнить группу задач DAG с помощью async.io?Python

Программы на Python
Anonymous
Как выполнить группу задач DAG с помощью async.io?

Сообщение Anonymous »

Я написал следующий код:
`

Код: Выделить всё

    async def execute_task(self, task_id):
await self.tasks[task_id]()
self.task_status[task_id] = "done"

async with self.lock:
for successor in self.dag.successors(task_id):
if all(self.task_status[predecessor] == "done" for predecessor in self.dag.predecessors(successor)):
self.ready_queue.append(successor)

async def execute_dag(self):
for task in self.dag.nodes:
if self.dag.in_degree(task) == 0:
self.ready_queue.append(task)

while self.ready_queue:
current_tasks = []
async with self.lock:
while self.ready_queue:
task_id = self.ready_queue.popleft()
self.task_status[task_id] = "running"
current_tasks.append(self.execute_task(task_id))

await asyncio.gather(*current_tasks)

Для следующего тестового примера:

Код: Выделить всё

1 => 2
1 => 5
2 => 3
3 => 4
5 => 4
У меня есть результат:

Код: Выделить всё

Executing Task 1
Task 1 Done
Executing Task 2
Executing Task 5
Task 2 Done
Task 5 Done
Executing Task 3
Task 3 Done
Executing Task 4
Task 4 Done
Где задача 5 занимает значительно больше времени. Вот почему задача 3 ждет завершения задачи 5? asyncio.gather, похоже, запускает задачи пакетно. Есть ли обходной путь?
Мне нужен минимальный код для параллельного выполнения задач. Но несколько тестовых случаев кажутся узким местом в моей реализации.

Подробнее здесь: https://stackoverflow.com/questions/786 ... g-async-io

Вернуться в «Python»