当前位置: 首页 > news >正文

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, 那么线程池基本上就是首选的解决办法。首先要把控住并发数, 接着还要补齐异常、超时以及退出信号, 这样的代码才敢放置到生产脚本当中去运行。

http://www.cnnetsun.cn/news/4258693.html

相关文章:

  • C++函数模板实战:从距离计算到泛型编程核心原理
  • 浏览器鼓机音序器进阶:Web Audio时钟调度与架构拆解
  • 做弱电工程,这些线材一定要认识
  • 基于Django与Python的适老化健康预警系统:架构设计与工程实践
  • FANUC上位机开发实战:C#连接PMC与MES回传设计
  • 从数学建模到数据挖掘实战:古代玻璃成分分析全流程解析
  • 不会Python?AI帮你写脚本,自动化办公(保姆级教程)
  • chatgpt赋能python:Python可以跨平台吗?
  • 基于YOLOv11的柑橘果柄识别:从数据集构建到模型部署的完整实践
  • 工业级智能决策系统:DSAC+双层MLP落地实践
  • 3 个独立开发者,用 AI 给自己做了融资 FA、求职诊断和效率工具
  • 基于样本平均近似与机器学习的血管机器人订购策略建模与Matlab实现
  • 从集合到范畴:图解范畴论核心概念与编程实践
  • 蓝桥杯国赛真题“123”解析:从数学规律到二分查找的算法优化实践
  • Python模拟退火算法求解整数规划:从原理到实战调优
  • 原码、反码、补码与位运算(与/或/异或/取反)
  • Python数学建模实战:数据拟合、优化与蒙特卡洛模拟核心技巧
  • Agentic RAG工作流:轻量级智能体问答系统实战
  • 线性规划建模与求解:从数学建模到MATLAB/Python实战
  • C++模板类与STL实战:构建泛型数据管理器的工程化指南
  • 气动系统电磁阀选型
  • Signal拟推免手机号注册:一次性付费背后的账号体系设计与反滥用权衡
  • 蓝桥杯国赛“扩散”题解:从BFS模拟到曼哈顿距离的算法优化
  • VOC格式路面缺陷数据集的工程化解析与实战指南
  • AI助理技术拆解:用RAG打造企业知识库实战
  • 字符串周期模式匹配:贪心算法与分组统计实战解析
  • FPGA驱动VGA显示:从时序原理到工程实践全解析
  • ASP.NET返利购物商城系统:架构设计与佣金计算引擎实现
  • Parallels Desktop 27图形与AI性能提升全解析
  • Zero-Mem:零Token消耗的LLM Agent记忆管理新方案