`
Код: Выделить всё
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
Мне нужен минимальный код для параллельного выполнения задач. Но несколько тестовых случаев кажутся узким местом в моей реализации.
Подробнее здесь: https://stackoverflow.com/questions/786 ... g-async-io