并行处理中如何实现无交错的任务日志输出?
解决并行任务日志交错问题
并行执行带日志输出的任务时,不同任务的日志会交错显示。如果可以接受牺牲单任务日志的实时性,让每个任务的日志连续输出,有个非常简便的实现方式。
简便实现方案
直接让每个任务先缓存自身的日志消息,等任务执行完毕后再批量输出:
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
相关产品推荐
相关产品推荐

