parent
0f97d359f8
commit
a4ae2b52a0
1 changed files with 2 additions and 2 deletions
|
|
@ -324,7 +324,7 @@ async def build_flow(
|
||||||
client_consumed_queue: asyncio.Queue,
|
client_consumed_queue: asyncio.Queue,
|
||||||
event_manager: "EventManager",
|
event_manager: "EventManager",
|
||||||
) -> None:
|
) -> None:
|
||||||
build_task = asyncio.create_task(await asyncio.to_thread(_build_vertex, vertex_id, graph, event_manager))
|
build_task = asyncio.create_task(asyncio.to_thread(asyncio.run, _build_vertex(vertex_id, graph, event_manager)))
|
||||||
try:
|
try:
|
||||||
await build_task
|
await build_task
|
||||||
except asyncio.CancelledError as exc:
|
except asyncio.CancelledError as exc:
|
||||||
|
|
@ -359,7 +359,7 @@ async def build_flow(
|
||||||
async def event_generator(event_manager: EventManager, client_consumed_queue: asyncio.Queue) -> None:
|
async def event_generator(event_manager: EventManager, client_consumed_queue: asyncio.Queue) -> None:
|
||||||
if not data:
|
if not data:
|
||||||
# using another thread since the DB query is I/O bound
|
# using another thread since the DB query is I/O bound
|
||||||
vertices_task = asyncio.create_task(await asyncio.to_thread(build_graph_and_get_order))
|
vertices_task = asyncio.create_task(asyncio.to_thread(asyncio.run, build_graph_and_get_order()))
|
||||||
try:
|
try:
|
||||||
await vertices_task
|
await vertices_task
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue