Python多线程提速秘籍:告别缓慢的批量处理
接口速度并非那种明显迟缓的情况, 单个请求所耗费的时间是 200 毫秒, 然而在批量运行 5000 个的时候, 竟然需要十几分钟之久。这样的代码我晓得好多, 将其打开瞧一瞧, 基本上都是一个 for 循环从起始处一直怼到末尾哒:for user_id in user_ids: profile load_user_profile(user_id) write_snapshot(profile)这处所在, 我于第一眼之际不会产生怀疑, 慢下来, 也断不会率先着手去对函数内部予以优化。只要其中存在网络请求、文件读写以及数据库查询, 大概率便是线程未被加以运用。做多线程最常见有两种写法 . 和 。在平常的业务代码以内, 我更为倾向于运用后面这一个。线程池具备可控性, 不容易在情绪冲动之下创建出数千个线程, 继而将机器弄得风扇飞速运转。像是一个用于批量拉取用户状态的脚本, 该脚本对接的那个接口时不时会出现抖动现象, 绝不能够因为其中一次失败, 就将整批任务给毁掉:import time import random from concurrent.futures import ThreadPoolExecutor, as_completed defquery_user_status(user_id): begin time.time # 这里模拟一次远程接口调用 time.sleep(random.uniform(0.05, 0.3)) if user_id % 17 0: raise RuntimeError(fremote api timeout, user_id{user_id}) cost_ms int((time.time - begin) * 1000) return { user_id: user_id, status: ACTIVE, cost_ms: cost_ms } defbatch_query(user_ids): ok_rows bad_rows with ThreadPoolExecutor(max_workers12, thread_name_prefixuser-sync) as pool: future_map { pool.submit(query_user_status, user_id): user_id for user_id in user_ids } for future in as_completed(future_map): user_id future_map[future] try: row future.result ok_rows.append(row) except Exception as e: bad_rows.append((user_id, str(e))) return ok_rows, bad_rows if __name__ __main__: users list(range(1, 101)) ok, bad batch_query(users) print(success:, len(ok)) print(failed:, bad[:5])这段代码能解决大部分“批量处理慢”的问题。留意, 我在这儿没将其写成一百、二百, 那个情形不会出现。线程并非数量越多就越好, 当接口运行显得缓慢之际, 增添线程的确能够把等待的时间堆积起来, 然而线程数量一旦增多之后, 调度、连接数、处于下游时受到的流量限制就会相应接踵而至。在线上的时候呢, 我通常首先是从像数字8、12、16这样的数开始着手去尝试的, 而并非不去思考就随意地写下一个100。要是任务之间旨在共享数据, 那就千万别随意去更改全局变量。此坑极为隐蔽, 在开发环境下运行 20 条数据不成问题 , 可一上线运行 20 万条 , 数量偶尔有所减少几条 , 而日志里还看不出来差异状况。比如下面这种写法看着没毛病其实不稳total 0 defadd_count: global total total 1整体加上一, 并非是一个不能够拆解开来的动作, 在这个过程当中, 有可能会被其他的线程穿插进来。要么用锁import threading from concurrent.futures import ThreadPoolExecutor classCounter: def__init__(self): self.value 0 self._lock threading.Lock defincr(self, step1): with self._lock: self.value step defhandle_one_line(line, counter): ifERRORin line: counter.incr if __name__ __main__: lines [ INFO order created, ERROR payment timeout, WARN retry later, ERROR inventory locked, ] * 1000 counter Counter with ThreadPoolExecutor(max_workers6) as pool: for line in lines: pool.submit(handle_one_line, line, counter) print(counter.value)但要提醒的是, 锁可千万别随意添加。一旦锁被加大, 多线程便又会退化为仅单线程运行。若能够使得每个线程都各自计算属于自己的结果, 最终再进行汇总, 那就千万别去共享同一个变量了。还有一种写法, 更为常见, 那便是运用 queue.Queue 来充当生产者消费者。我在处理日志, 以及导文件, 并补数据之际, 常常这般去写。有一个线程专门负责读, 存在几个线程专门负责处理, 最终再进行统一的落盘或者写库。import queue import threading import time task_queue queue.Queue(maxsize1000) stop_flag object defread_log_file(file_path): with open(file_path, r, encodingutf-8) as f: for line in f: iforderIdin line: task_queue.put(line.strip) for _ in range(4): task_queue.put(stop_flag) defparse_worker(worker_no): whileTrue: line task_queue.get try: if line is stop_flag: return # 模拟解析日志里的订单号 order_id line.split(orderId)[-1].split[0] time.sleep(0.02) print(fworker{worker_no}, order_id{order_id}) finally: task_queue.task_done if __name__ __main__: reader threading.Thread( targetread_log_file, args(app.log,), namelog-reader ) workers [ threading.Thread(targetparse_worker, args(i,), nameflog-parser-{i}) for i in range(4) ] reader.start for t in workers: t.start reader.join task_queue.join for t in workers: t.join这段代码有两个细节。有一种情况是 Queue(1000) , 我对那种无限制的队列是持不倾向的态度呢。读取文件的速度表现得过于迅速, 然而进行处理的速度却显得迟缓不已, 如此一来内存就会出现被逐步支撑起来的状况。添置一个有关大小的限制条件, 通过这种方式起码能够使得处于读取阶段的那部分等待后续展开涉及处理线程才行。另外一个情况是, 线程没办法凭借猜测来结束, 特别是消费者线程。要是没有确切的退出信号, 那么就极易在get这个地方卡住, 脚本看上去没有报出错误然而就是不会结束。还有个绕不开的问题 多线程到底能不能提升性能得看任务类型。若是针对请求接口, 以及读写文件, 还有查数据库, 再者是扫日志的这种情况而言, 多线程很大程度上通常是具备有效性的, 这是由于线程于相当多的时间当中皆处于等待IO状态下的, 此期间CPU并没太多实际投入工作量去执行任务呢。设有这样一些情况, 像是进行压缩图片之作, 于其中实施跑复杂计算的行为, 去解析大 JSON , 又或者是操作大量加密解密的相关事宜在这般属于 CPU 密集任务的状况之下, 对于多线程所能够呈现出的效果, 就不要怀有太大的期望了。因为这里面存在着 GIL , 众多的线程并不能等同于多个 CPU 核能够同时进行迅猛的运行。在面临这种情况的时候, 我通常会选择进行更换之举, 亦或是把那些繁重的任务交付给 C 扩展、NumPy、还有外部服务去进行处理。还有一点线程里的异常不会像主线程那样直接把你喊醒。有着这样的益处就在这儿, 它会将异常给 вновь抛掷出来, 如此一来, 你便能知晓究竟是哪一条数据出现了问题, 而并非是线程悄然无声地灭亡, 主流程还打印出一句“执行完成”。多线程不是为了把代码写得高级。它所解决的, 是那种极端具体的问题, 有大量的任务都处于等待状态, 然而主线程, 却呆呆地站在那里, 一个个地去排队。只要能够确认瓶颈所在之处是在IO方面, 那么线程池基本上便是那解决问题的第一把利刃。先是把并发数给控制住, 接着, 再将异常、超时以及退出信号等方面予以补齐, 这样的代码, 才敢被放置到生产脚本里去运行。

相关新闻

最新新闻

日新闻

周新闻

月新闻