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

并行处理中如何实现无交错的任务日志输出?

解决并行任务日志交错问题

并行执行带日志输出的任务时,不同任务的日志会交错显示。如果可以接受牺牲单任务日志的实时性,让每个任务的日志连续输出,有个非常简便的实现方式。

简便实现方案

直接让每个任务先缓存自身的日志消息,等任务执行完毕后再批量输出:

import logging
from concurrent.futures import ThreadPoolExecutor
from time import sleep

logger = logging.getLogger(__name__)

def task(n):
    # 缓存当前任务的所有日志消息
    task_logs = []
    task_logs.append(f"Task {n} started")
    sleep(0.001)
    task_logs.append(f"Task {n} finished")
    
    # 任务完成后统一输出日志
    for msg in task_logs:
        logger.info(msg)
    return n * n

def main():
    with ThreadPoolExecutor(max_workers=2) as executor:
        futures = [executor.submit(task, i) for i in range(5)]
        _ = [f.result() for f in futures]

if __name__ == "__main__":
    logging.basicConfig(level=logging.INFO)
    main()

原理说明

每个任务执行过程中,先把要输出的日志暂存在本地列表里,不立即打印。等整个任务完成后,再一次性把所有日志消息输出。这样即使多个任务并行运行,每个任务的日志都是连续显示的,不会和其他任务的日志交错。

进阶优化(可选)

如果不想修改任务内部的日志调用逻辑,可以用装饰器封装日志缓存逻辑,让代码更整洁:

import logging
from concurrent.futures import ThreadPoolExecutor
from time import sleep
from io import StringIO

logger = logging.getLogger(__name__)

def batch_log(func):
    def wrapper(*args, **kwargs):
        # 创建内存流捕获日志
        log_stream = StringIO()
        handler = logging.StreamHandler(log_stream)
        handler.setLevel(logging.INFO)
        handler.setFormatter(logging.Formatter('%(levelname)s:%(name)s:%(message)s'))
        logger.addHandler(handler)
        
        try:
            # 执行任务
            result = func(*args, **kwargs)
        finally:
            # 移除临时handler
            logger.removeHandler(handler)
        
        # 输出捕获到的日志
        log_content = log_stream.getvalue().strip()
        if log_content:
            print(log_content)
        return result
    return wrapper

@batch_log
def task(n):
    logger.info(f"Task {n} started")
    sleep(0.001)
    logger.info(f"Task {n} finished")
    return n * n

def main():
    with ThreadPoolExecutor(max_workers=2) as executor:
        futures = [executor.submit(task, i) for i in range(5)]
        _ = [f.result() for f in futures]

if __name__ == "__main__":
    # 关闭默认控制台输出,避免重复打印
    logging.basicConfig(level=logging.INFO, handlers=[])
    main()

这种方式不需要改动原任务的日志代码,通过装饰器统一处理日志的缓存和批量输出。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 20:27:42