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(f"remote 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_workers=12, thread_name_prefix="user-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])这段代码能解决大部分“批量处理慢”的问题。
留意, 我这边没将 记成 100、200。线程并非数量越多便越好。于接口运行迟缓之际, 增添线程确实能够把等待的时间累积起来 , 然而当线程数量增多之后 , 调度 、连接数量 、下游实施限流举措都会随之出现。
在线上的时候, 我通常是先从像8、12、16这样的数着手去试, 并非是随意地拍脑袋想出一个100来。
若任务之间存在共享数据的情况, 那就千万别随意去改动全局变量, 这个坑极为隐蔽不说, 在开发环境运行 20 条数据时不会有问题, 但上线后运行 20 万条, 数量偶尔会少几条, 且从日志中还看不出来。
比如下面这种写法,看着没毛病,其实不稳:
total = 0 defadd_count: global total total += 1对于total += 1这个行为而言, 它并非是那种不可拆分的动作, 在其执行进程当中, 存在着这样一种可能性, 即会有别的线程插足进来。
要么用锁:
import threading from concurrent.futures import ThreadPoolExecutor classCounter: def__init__(self): self.value = 0 self._lock = threading.Lock defincr(self, step=1): with self._lock: self.value += step defhandle_one_line(line, counter): if"ERROR"in 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_workers=6) 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(maxsize=1000) stop_flag = object defread_log_file(file_path): with open(file_path, "r", encoding="utf-8") as f: for line in f: if"orderId="in 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(f"worker={worker_no}, order_id={order_id}") finally: task_queue.task_done if __name__ == "__main__": reader = threading.Thread( target=read_log_file, args=("app.log",), name="log-reader" ) workers = [ threading.Thread(target=parse_worker, args=(i,), name=f"log-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, 那么线程池基本上就是首选的解决办法。首先要把控住并发数, 接着还要补齐异常、超时以及退出信号, 这样的代码才敢放置到生产脚本当中去运行。
