ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Python多进程异步日志实现:告别FileHandler同步写入卡顿

Python多进程异步日志实现:告别FileHandler同步写入卡顿 我先说一个三年前的真实场景一个爬虫系统8个worker进程并发往同一个日志文件里写东西跑了一个下午日志文件开头出现连续的空行、错位、甚至半截消息。当时我第一个反应是给FileHandler加锁结果业务线程的耗时不涨反跳单次日志等待从不到1毫秒直接飙到几百毫秒。也是从那时候起我下定决心不再干“在生产环境里继续给同步日志打补丁”这件事转而在Python 3.12环境里实现了一个专门用于多进程场景的 MultiProcessAsyncLogger。这篇文章不打算讲asyncio这里的Async和协程没有关系核心思路是每个业务进程把日志记录丢进一个跨进程Queue由一个独立的消费者进程统一负责写盘。听起来很简单但是等到真正落地的时候你会发现序列化、进程启动方式、优雅关闭、队列大小每一个环节都能给你上一课。下面把这些内容完整拆开代码可以直接抄建议结合你自己的业务进程模型改一改再上生产。1. 为什么多进程下日志必须“异步”三个真实事故1.1 日志串行写入导致业务线程卡死绝大多数项目里日志使用的都是logging.FileHandler。这个handler在每次emit()的时候都会调用self.flush()也就是说每写一条日志都要和磁盘打一次交道。单机业务量小的时候看不出来一旦8个进程同时高频率打日志每个进程内部的业务线程都被磁盘IO拖住表现就是日志越多业务接口越慢。有人会反驳磁盘又不是SSD写一条几KB的日志不至于那么慢吧问题是物理磁盘在并发写入下的寻道开销会被放大而且Python的FileHandler在每次写入后同步flush这种频繁的fsync型操作在跨进程场景下就是灾难。我当时的压测数据是8个进程各自输出10万条日志同步方案整体耗时30秒以上业务线程单次日志调用最长等待超过900毫秒。这个等待还不能重试因为logger本身一般是不捕获也不吞掉异常的。1.2 FileHandler 的跨进程锁让人又爱又恨说到多进程写同一个文件很多人第一时间想到给日志加锁。问题是FileHandler内部用的锁是threading.RLock这个锁只能保证同一个进程内的线程安全对跨进程写文件完全无能为力。不同进程各自持有文件描述符写同一文件时靠的是操作系统底层的文件偏移量靠append模式不一定能避免半截消息。如果自己包一层multiprocessing.Lock呢能解决交错问题但会引入更严重的性能问题所有进程写日志都要抢同一把锁等锁的时间比写磁盘还长。我曾经试过这种方案结果是吞吐量暴跌日志系统的锁变成了整个业务的串行点。这个方案很快被我否决了。1.3 花式绕坑不如换模型在那之后我见过各种奇怪做法每个进程写独立日志文件之后再合并把日志发到本地UDP端口由另一个进程接收。前者排查问题时要在多个文件之间跳来跳去后者引入了额外网络依赖不够干净。后来我想通了问题根源是“日志产生”和“日志写入”耦合在同一个调用栈里。只要把这两个动作拆开多线程、多进程的写文件竞争就自然消失了。具体的拆法就是我后面要讲的所有业务进程只负责把LogRecord放进一个共享队列永远不直接接触文件句柄队列尽头安排一个独立消费者进程统一从队列取日志并落盘。这个模型日志产生端不再有任何磁盘IO业务线程自然就不会被日志拖死。2. 先复盘点基础QueueHandler、QueueListener 和 multiprocessing.Queue 是怎么配合的2.1 生产端和消费端的正统配合logging.handlers.QueueHandler和QueueListener是Python内置的一对搭档本身就是为了解决“日志生产和消费解耦”设计的。QueueHandler作为logger的handler会把LogRecord放进一个queueQueueListener在后台从同一个queue取Record然后转发给真正干活的handler。问题在于内置的QueueListener只是一个线程默认是单进程内使用的。多进程场景下你不能在每个业务进程里都开一个QueueListener线程——那样还是多个线程/进程写同一文件问题一点没解决。正确做法是把这个Listener放到一个独立进程中去跑或者干脆不用QueueListener自己写一个消费者循环。我更推荐后者原因后面会详细说。2.2 一套完整的最小链路先看最简单的异步日志链路。这个链路里主进程创建Queue启动一个消费者进程再派发若干worker进程每个worker进程的root logger挂上QueueHandler业务日志只需要正常调用logger.info即可。import logging import logging.handlers import multiprocessing import queue import time SENTINEL __SENTINEL__ def setup_child_logger(q: multiprocessing.Queue) - None: root logging.getLogger() root.setLevel(logging.DEBUG) for h in root.handlers[:]: root.removeHandler(h) h.close() handler logging.handlers.QueueHandler(q) root.addHandler(handler) root.propagate False def consumer_main(q: multiprocessing.Queue, log_file: str) - None: logger logging.getLogger(consumer) logger.setLevel(logging.DEBUG) logger.propagate False fh logging.FileHandler(log_file, encodingutf-8) fh.setFormatter(logging.Formatter(%(asctime)s %(processName)s %(levelname)s %(message)s)) logger.addHandler(fh) while True: try: record q.get(timeout0.5) except queue.Empty: continue if record SENTINEL: q.task_done() break logger.handle(record) q.task_done()如果只用内置的QueueListener消费者进程里大概长这样listener.start()之后进程必须保持运行退出时再调用listener.stop()。但QueueListener的start内部是启动一个非守护线程进程退出时如果线程没有自然结束行为并不好控制。尤其当你需要在主进程结束前确保所有日志都写完时手动循环能精确配合task_done()和queue.join()内置QueueListener没有这么方便的钩子。2.3 为什么不能直接把QueueListener放在业务进程里有一种偷懒做法在每一个worker进程内部都创建一个QueueListener。这种方案只要worker一多日志消费端变成多个写文件又回到了竞争状态而且不同进程可能都拿到同一批记录造成日志重复。异步方案的核心前提是消费者只有一个哪怕将来扩到多个消费者也必须想清楚接收方之间如何协调。所以生产者可以有很多个消费者路径一定要收敛。多进程异步日志的本质就是“生产端分散消费端集中”。3. 从零实现 MultiProcessAsyncLogger3.1 架构与进程拓扑我最终实现的MultiProcessAsyncLogger并不复杂进程拓扑如下主进程创建 multiprocessing.Queue启动消费者进程然后派发业务worker进程。worker进程调用setup_child_logger(queue)把root logger接上QueueHandler之后所有logging.getLogger(某名字)写出的日志都会进入共享队列。消费者进程从队列取出LogRecord交给带有RotatingFileHandler的logger执行真正的写盘。这个结构里业务进程再也不会直接打开日志文件日志文件的句柄只属于消费者进程。队列本身具备跨进程传递能力所以从Linux的fork到Windows的spawn都能正常工作前提是代码结构符合multiprocessing的基本要求。3.2 完整实现一个可落地的类下面这段代码我在Python 3.12下验证过核心逻辑也适用于3.7以上版本。为了控制复杂度这里只保留最关键的部分。import logging import logging.handlers import multiprocessing import queue import pickle from typing import Optional SENTINEL __SENTINEL__ class SafeQueueHandler(logging.handlers.QueueHandler): 解决跨进程pickle序列化问题见第4章详细说明。 def prepare(self, record: logging.LogRecord) - logging.LogRecord: record super().prepare(record) try: pickle.dumps(record) except Exception: for key, value in list(record.__dict__.items()): try: pickle.dumps(value) except Exception: setattr(record, key, str(value)) return record def setup_child_logger(queue: multiprocessing.Queue, level: int logging.DEBUG) - None: root logging.getLogger() root.setLevel(level) for handler in root.handlers[:]: root.removeHandler(handler) handler.close() handler SafeQueueHandler(queue) root.addHandler(handler) root.propagate False def consumer_main(queue: multiprocessing.Queue, log_file: str, level: int logging.DEBUG, max_bytes: int 100 * 1024 * 1024, backup_count: int 5) - None: logger logging.getLogger(async-consumer) logger.setLevel(level) logger.propagate False file_handler logging.handlers.RotatingFileHandler( log_file, maxBytesmax_bytes, backupCountbackup_count, encodingutf-8, ) formatter logging.Formatter( %(asctime)s | %(processName)s | %(threadName)s | %(levelname)s | %(name)s | %(message)s ) file_handler.setFormatter(formatter) logger.addHandler(file_handler) while True: try: record queue.get(timeout0.5) except queue.Empty: continue if record SENTINEL: queue.task_done() break logger.handle(record) queue.task_done() class MultiProcessAsyncLogger: def __init__(self, log_file: str, level: int logging.INFO, queue_size: int 10000, backup_count: int 5): self.log_file log_file self.level level self.queue: multiprocessing.Queue multiprocessing.Queue(maxsizequeue_size) self.consumer: Optional[multiprocessing.Process] None def start(self) - None: self.consumer multiprocessing.Process( targetconsumer_main, args(self.queue, self.log_file, self.level), nameasync-log-consumer, daemonTrue, ) self.consumer.start() def attach(self) - None: setup_child_logger(self.queue, logging.DEBUG) def shutdown(self) - None: self.queue.join() self.queue.put(SENTINEL) self.consumer.join(timeout10) if self.consumer.is_alive(): self.consumer.terminate()使用示例import multiprocessing import time def worker(job_id: int, q: multiprocessing.Queue) - None: setup_child_logger(q) logger logging.getLogger(fworker.{job_id}) for i in range(100): logger.info(job %s processing item %s, job_id, i) time.sleep(0.01) if __name__ __main__: async_logger MultiProcessAsyncLogger(/tmp/async_app.log, levellogging.INFO) async_logger.start() processes [] for j in range(4): p multiprocessing.Process(targetworker, args(j, async_logger.queue)) p.start() processes.append(p) for p in processes: p.join() async_logger.shutdown()这里的attach()方法比直接在外部调用setup_child_logger更方便管理worker进程启动后只需要async_logger.attach()一次root logger就接到队列上了。注意不要为了省事把整个async_logger对象作为参数传给Processspawn模式下Process只序列化args不会自动帮你序列化成员对象传queue足够。3.3 优雅关闭的顺序很关键很多人在自己的异步日志方案里遇到过“最后几条日志丢了”的问题十有八九是关闭顺序错了。我建议的顺序是for p in processes: p.join() async_logger.shutdown()先让所有业务worker退出确保不再有新日志产生然后queue.join()等待队列里所有的LogRecord被消费者处理完再发送SENTINEL哨兵消费者收到后从循环退出最后consumer.join()回收进程。如果不做queue.join()主进程直接给队列放一个SENTINEL消费者有可能在还有日志记录没处理完的情况下先碰到哨兵然后直接退出后面的日志全部丢失。task_done()和queue.join()这套机制就是为了避免这个问题才存在的。3.4 配置项可以再抽象一层上面的类参数只有日志路径、级别、队列大小、轮转体积和备份数量。实际生产里我建议再加两个参数消费者进程数量和格式化器模板。消费者进程数量后面会说默认用1格式化模板可以做成JSON格式方便后面接入日志采集平台只需要把consumer_main里的formatter换掉即可。如果项目里多个模块需要不同的日志级别可以再包一层dataclass配置。但核心原则不变所有配置最终都落到消费者进程里的handler上业务进程里只关心队列。4. 最容易翻车的序列化问题让日志记录能安全穿越队列4.1 QueueHandler 的 prepare 机制到底做了什么事multiprocessing.Queue在put数据时会对对象做pickleLogRecord本身不是一个天然的pickle对象所以QueueHandler在put之前会调用prepare()。这个方法默认做的事情是先把消息格式化一遍然后清空args和exc_info。也就是说异常堆栈对象不会原样穿过队列队列里携带的是一条已经格式化好的文本。这是很多人忽略但非常加分的机制。但是prepare()没有清理record.__dict__里的其他自定义属性。如果你用extra{trace_id: trace_id, user_id: user_id}这种写法给日志增加字段这些字段会原样放进record进而参与pickle。若字段值本身是能被pickle的对象比如字符串、数字、UUID没问题一旦放了无法pickle的对象整个队列写入就会失败。4.2 一个真实案例extra里放了不能被pickle的对象我遇到过一段业务代码日志里带了request对象想打请求ID。进程内运行没问题换成跨进程队列后在worker进程那侧调用logger.info时直接抛了TypeError: cannot pickle _io.BufferedReader object。原因是request里某个属性是文件流被QueueHandler放进了队列。更麻烦的是这个异常发生在业务代码的日志调用处如果不捕获会把原本可以正常跑完的业务逻辑打断。日志系统反而变成了业务故障源。4.3 解决方案SafeQueueHandler所以就有了SafeQueueHandler。它覆盖prepare()先让父类完成默认处理然后尝试对整个record做一次pickle。如果失败就逐个检查record.__dict__里的字段把不可pickle的值转换成字符串。这样至少保证日志不会因为一个带不出进程的属性而断掉。class SafeQueueHandler(logging.handlers.QueueHandler): def prepare(self, record: logging.LogRecord) - logging.LogRecord: record super().prepare(record) try: pickle.dumps(record) except Exception: for key, value in list(record.__dict__.items()): try: pickle.dumps(value) except Exception: setattr(record, key, str(value)) return record在实际业务中我建议除了SafeQueueHandler还要约定一条纪律日志的extra字段只允许放标量类型。序列化兜底是最后防线不应变成设计标准。5. 性能实测与参数调优5.1 同步与异步吞吐量对比我做了一组简单压测机器是4核云主机磁盘为普通云盘8个进程并发每个进程写10万条日志日志内容大约200字节。同步方案采用8个进程各自持有FileHandler异步方案就是上面的MultiProcessAsyncLogger。结果如下表方案业务日志调用最长耗时总体写盘完成时间日志是否完整同步FileHandler约900ms32秒文件交错、少量半截异步Queue 单消费者约2ms4秒完整有序按入队顺序不同环境的绝对数值会有差异但结论基本一致异步方式下业务线程的日志调用耗时下降了两个数量级。消费端多花的时间主要是在格式化器和RotatingFileHandler的写入上这部分时间被转移到了消费者进程不再干扰业务。5.2 队列大小要按“峰值内存”来算multiprocessing.Queue的maxsize不是缓存条数上限那么简单。每条LogRecord本身包含时间、进程名、线程名、消息文本、extra字段等如果消息里有大文本一条可能达到几KB甚至几十KB。假设队列大小设成10000排队日志的瞬时内存可能几十MB甚至几百MB。在Consumer写入速度跟上来的情况下队列会稳定在一个低水位日志峰值到来时队列开始积压。如果积压到maxsizeQueueHandler的put()会阻塞住调用方业务线程。从我的经验看队列大小设置为“一秒钟日志峰值条数 × 平均单条日志内存 × 2”比较合理。比如每秒峰值5000条单条约1KB那就是5000 × 1KB × 2 10MB对应的队列条目数。你可以先根据单条内存估算再用实际压测调整。5.3 单消费者还是多消费者顺序和吞吐的权衡默认我会坚持单消费者。因为日志文件只有一个单消费者顺序写是天然有序的RotatingFileHandler也不需要应对跨进程锁竞争。如果消费者进程成为瓶颈优先优化formatter和handler比如使用更简单的字符串格式化、避免正则而不是盲目增加消费者。只有一种情况我会考虑多消费者消费者进程还需要把日志转发给远程日志系统比如通过HTTP或者消息队列发送发送延迟不稳定单消费者会导致队列积压。这时可以让多个消费者进程各自从同一个queue里取记录但要接受跨进程写文件时的顺序错乱。如果日志系统的下游可以接受乱序再把消费者数量调大。5.4 Windows spawn 和 Linux fork 的差异Python 3.12 在Linux上默认fork在Windows和macOS上默认spawn。fork模式下子进程会继承父进程内存中的logger配置所以经常出现“子进程里日志重复打印”的情况spawn模式下子进程是全新的解释器不会继承任何父进程logger但你传给Process的queue可以作为参数正常传过去。所以代码里必须有一个独立于父进程的attach()或setup_child_logger(queue)在每个worker进程开头调用一次。不要试图在模块导入时配置logger因为spawn模式下子进程会重新导入模块如果你的日志配置写在模块顶层会在子进程里被重复执行。如果你在Windows上运行代码时遇到RuntimeError: An attempt has been made to start a new process before the current process has finished its bootstrapping phase十有八九是没把启动逻辑放进if __name__ __main__。这不是日志模块的问题而是multiprocessing的reuse模式要求。6. 生产环境避坑清单6.1 子进程日志重复打印重复打印问题最常见的来源不是QueueHandler而是子进程继承了父进程已经存在的handler。fork之后子进程的内存里本来就有父进程配置好的FileHandler你又调用了一次attach()于是同一份日志会经历两套handler。解决办法就是setup_child_logger开头清空root.handlers。我的代码里每次都先移除旧handler再添加QueueHandler就是防止这种情况。6.2 日志丢最后几条除了关闭顺序还有一个很容易被忽略的情况消费者进程被设置成daemonTrue。主进程退出时daemon进程不会优雅退出即使队列里还有日志也会被强制结束。所以在主进程结束前必须显式调用shutdown()。如果消费者进程不是daemon又可能造成主进程退出后进程仍然残留所以我选择daemon 显式shutdown两条路都走shutdown的join(timeout10)兜底超时直接terminate。6.3 消费者进程崩溃后队列塞满消费者进程如果意外崩溃最直接的后果是队列越积越多直到QueueHandler的put阻塞业务进程。我建议做一个简单的探活机制业务进程定期检查消费者是否alive如果发现死了先把队列里的数据落一个紧急本地文件再重启消费者进程。生产环境中消费者进程的日志文件如果磁盘满了也会导致写入失败所以日志目录的磁盘监控一定要做。6.4 与 RotatingFileHandler 结合时的注意点因为所有文件写入都集中在消费者进程RotatingFileHandler在这个模型里是安全的。它自己会按照maxBytes和backupCount切分文件不需要额外加锁。如果你的日志需要同时输出到多个文件或多种格式就多挂几个handler不要搞出多个消费者写同一个文件。6.5 在 Celery、任务队列等框架里怎么用只要框架允许你在worker进程启动时执行一个初始化函数就能用这套方案。以Celery为例在worker_process_init信号里调用async_logger.attach()然后在每个task里正常使用logging即可。任务进程退出时不需要每个worker都调用shutdown只有主进程负责在Celery worker主循环结束时调用shutdown()。我在实际项目里使用的最后版本就是把这个组件包成一个很小的三方库内部用logger的name区分模块所有日志统一经过队列再由消费者进程格式化输出。上线至今最大的感受是日志系统再也不是业务卡顿的制造者排查问题的时候也更愿意打开日志去看了。如果你正准备给多进程项目写日志系统记住三件事第一队列是所有进程共享的那根神经不要轻易换成普通线程queue第二关闭顺序永远遵循“停生产者—等队列清空—发哨兵—收消费者”第三给自定义extra字段设计兜底序列化策略否则某一天它一定会用一个奇怪的对象把业务打断。这些坑我都替你踩过一遍代码可以直接拿去做基础版本剩下的调优就看你的日志量级和应用场景了。
返回列表