You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Kotlin中如何让进程等待所有协程完成且不阻塞父线程?

实现进程等待所有协程完成且不阻塞父线程的方案

问题背景

应用中需对大量文件执行繁重任务,为每个任务创建协程。直接运行时进程会在协程完成前退出;使用job.join()会阻塞已有其他任务的父线程,不可行;需实现类似线程的进程等待机制,且不使用阻塞的async-await。

核心思路

用非守护线程托管asyncio事件循环:

  • 单独开一个非守护线程运行事件循环,进程会自动等待非守护线程结束后再退出
  • 所有文件处理协程通过线程安全的方式提交到这个循环
  • 追踪所有待完成的协程,最后一个任务结束时主动停止事件循环,让托管线程正常退出,进程随之结束

完整代码实现

import asyncio
import threading
from typing import List

# 全局变量:事件循环、待处理任务列表、线程锁
_task_loop: asyncio.AbstractEventLoop = None
_pending_tasks: List[asyncio.Task] = []
_task_lock = threading.Lock()

def _run_event_loop():
    """运行asyncio事件循环的线程函数"""
    global _task_loop
    _task_loop = asyncio.new_event_loop()
    asyncio.set_event_loop(_task_loop)
    _task_loop.run_forever()
    _task_loop.close()

# 启动事件循环线程(非守护线程,进程会等待它退出)
_loop_thread = threading.Thread(target=_run_event_loop, daemon=False)
_loop_thread.start()

def submit_coroutine(coro):
    """线程安全地提交协程到事件循环,并追踪任务状态"""
    async def _wrapped_task():
        try:
            await coro
        finally:
            # 任务完成后从待处理列表移除
            with _task_lock:
                current_task = asyncio.current_task()
                if current_task in _pending_tasks:
                    _pending_tasks.remove(current_task)
            # 所有任务完成时,停止事件循环
            with _task_lock:
                if not _pending_tasks:
                    _task_loop.call_soon_threadsafe(_task_loop.stop())

    # 线程安全地创建并提交任务
    task = _task_loop.call_soon_threadsafe(
        asyncio.create_task,
        _wrapped_task()
    ).result()
    with _task_lock:
        _pending_tasks.append(task)

# ------------------- 示例业务代码 -------------------
async def process_file(file_path):
    """模拟繁重的文件处理协程"""
    print(f"开始处理文件: {file_path}")
    # 替换为实际的文件读写、计算等操作
    await asyncio.sleep(2)
    print(f"完成文件处理: {file_path}")

def parent_thread_business():
    """父线程的原有任务,不受协程逻辑阻塞"""
    for i in range(5):
        print(f"父线程执行任务 {i}")
        threading.sleep(1)
    print("父线程任务全部完成")

if __name__ == "__main__":
    # 启动父线程的原有任务(实际场景中可能是主线程或其他工作线程)
    parent_thread = threading.Thread(target=parent_thread_business)
    parent_thread.start()

    # 提交多个文件处理协程
    for idx in range(3):
        submit_coroutine(process_file(f"document_{idx}.txt"))

    # 等待父线程任务完成(非必须,仅为模拟原有业务流程)
    parent_thread.join()

    # 进程会自动等待_loop_thread退出,也就是所有协程任务完成
    print("等待所有文件处理任务完成...")
    _loop_thread.join()
    print("所有任务执行完毕,进程退出")

关键细节说明

  • 非守护线程:_loop_thread设置daemon=False,进程必须等待该线程结束才会退出,这是实现进程等待协程的核心
  • 线程安全提交:用call_soon_threadsafe确保跨线程操作事件循环的安全性
  • 任务自动追踪:通过包装协程在任务结束时更新待处理列表,最后一个任务完成时主动停止事件循环,避免线程无限挂起
  • 父线程无阻塞:父线程的业务逻辑完全独立执行,不会被协程的等待逻辑卡住

内容的提问来源于stack exchange,提问作者Ahmed Eid

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.11 14:50:45