文章
多线程处理数据,主线程控制写数据
一、为什么需要这个模式
写并发程序时,一个很常见的需求是:
- 数据很多,单线程处理太慢,想用多线程加速。
- 处理结果需要写文件、写数据库。
- 但多线程直接写,会出现内容交错、顺序混乱、文件损坏、计数错误。
一个自然的解法是:让工作线程只负责计算,把结果交回主线程,由主线程统一写。
这就是“多线程处理数据,主线程控制写数据”。
它的核心不是靠锁,而是靠职责分离:并行的是处理,串行的是写。
二、先理解两种写法
Python 里实现这个模式,常见两种风格。
写法一:线程 + 队列
主线程把任务放进任务队列,工作线程从队列取任务、处理、把结果放进结果队列,主线程再从结果队列取结果写盘。
import threading
import queue
task_queue = queue.Queue()
result_queue = queue.Queue()
def worker():
while True:
try:
data = task_queue.get(timeout=1)
except queue.Empty:
break
result = process(data)
result_queue.put(result)
task_queue.task_done()特点:控制精细,适合流式处理,但代码量多,要自己管理线程生命周期。
写法二:线程池 + as_completed
用 concurrent.futures.ThreadPoolExecutor 提交任务,主线程用 as_completed 按完成顺序取结果,取一个写一个。
import concurrent.futures
with concurrent.futures.ThreadPoolExecutor(max_workers=3) as executor:
future_to_data = {
executor.submit(process, data): data
for data in data_list
}
for future in concurrent.futures.as_completed(future_to_data):
result = future.result()
f.write(result + "\n")特点:代码简洁,线程池自动管理线程,主线程流式消费结果。本文重点讲这种。
三、完整案例
场景:10 条数据,用 3 个线程并行处理,主线程统一写文件。
import concurrent.futures
import time
import random
data_list = [f"data_{i}" for i in range(1, 11)]
def process(data):
# 模拟耗时处理
time.sleep(random.uniform(0.1, 0.5))
return f"{data} -> 已处理"
def main():
written_count = 0
with open("output.txt", "w", encoding="utf-8") as f:
with concurrent.futures.ThreadPoolExecutor(max_workers=3) as executor:
# 提交任务,submit 立刻返回 Future
future_to_data = {
executor.submit(process, data): data
for data in data_list
}
# 谁先完成,谁先被主线程拿到
for future in concurrent.futures.as_completed(future_to_data):
data = future_to_data[future]
try:
result = future.result()
except Exception as e:
result = f"{data} -> 处理失败:{e}"
# 写操作只在主线程,且必须在 for 内部
f.write(result + "\n")
f.flush()
written_count += 1
print(f"[主线程] 写入第 {written_count} 条: {result}")
print(f"[主线程] 全部完成,共写入 {written_count} 条")
if __name__ == "__main__":
main()运行后你会看到:写入顺序不是 data_1 到 data_10,而是谁先处理完谁先写。但写入动作全部由主线程串行完成,文件内容不会交错。
四、几个关键点
1. 为什么不需要锁
这个模式里,f、written_count、result 都只在主线程访问。工作线程执行 process(data),只 return 结果,不碰任何共享资源。
锁的作用是保护共享资源。这里没有共享写,自然不需要锁。用结构设计消除竞争,比加锁更好。
ThreadPoolExecutor 和 as_completed 内部其实用了锁和条件变量,只是被封装了,你不需要手动管。
2. as_completed 到底怎么工作
executor.submit() 立刻返回一个 Future,不阻塞。as_completed(futures) 是迭代器,每次要“下一个元素”时,会阻塞等待,直到任意一个未完成的 Future 完成,就 yield 出来。
所以它等的不是某个特定任务,而是“下一个完成的”。快的先被处理,慢的不会堵住已经完成的结果。
3. 顺序问题
as_completed 按完成顺序返回,不是提交顺序。如果要求输出顺序和输入一致,要么收集完再排序,要么按提交顺序取:
for future in futures: # 按提交顺序
result = future.result() # 这里才阻塞等待
f.write(result + "\n")代价是会被最慢的任务拖住。
5. 异常处理
某个任务抛异常,future.result() 会重新抛出。不 try/except 的话,主线程会中断,后面的结果写不进去。所以循环里要单独捕获。
五、什么时候需要锁
只有多个线程真的并发写同一个资源时,才需要锁。比如让工作线程直接写文件:
lock = threading.Lock()
def process_and_write(data, f):
result = process(data)
with lock:
f.write(result + "\n")但这是不推荐的做法:
- 锁让写操作变成串行,拖慢并发。
- 写一半抛异常可能损坏文件。
- 顺序不可控。
更好的做法还是“工作线程只算,主线程只写”,或者用 queue.Queue(内部自带锁)传递结果。
六、对比表
| 对比项 | 多线程加锁写 | 主线程统一写 |
|---|---|---|
| 谁写 | 多个工作线程 | 只有主线程 |
| 是否需要锁 | 需要 | 不需要 |
| 顺序 | 不可控 | 可控 |
| 异常安全 | 差,可能写坏文件 | 好 |
| 并发效率 | 写被锁串行化 | 处理并行,写串行 |
| 代码复杂度 | 高 | 低 |
七、总结
- “多线程处理数据,主线程控制写数据”的核心是职责分离:工作线程只算,主线程只写。
- 写操作集中在单线程,天然串行,不需要锁。
as_completed按完成顺序返回结果,快的先处理,不被慢任务阻塞。- 写操作必须在
for循环体内,否则只会写最后一次结果。 - 锁只在多线程并发写同一资源时才需要,但更好的做法是用设计避免共享。
这个模式的价值不在于用了多少并发技巧,而在于把并行和串行放在正确的位置:并行处理,串行写入。
八、进阶实战:进度、重试、批量写入
真实场景里,光有 as_completed 还不够。你通常还需要:
- 进度显示:10000 条跑多久了,还剩多少。
- 失败重试:某条调用超时或报错,不能直接丢。
- 批量写入:每条都
flush会拖慢主线程,攒一批再写。 - 不让慢任务拖死整体:给单个任务设超时。
- 异常隔离:一条失败不影响其他条。
下面是一个完整加强版。
完整代码
import concurrent.futures
import time
import random
import json
import os
# ============================================================
# 模拟:一批待处理样本
# ============================================================
samples = [{"id": i, "text": f"sample_{i}"} for i in range(1, 101)]
# ============================================================
# 模拟:调用大模型质检,有一定概率失败
# ============================================================
def call_llm(sample):
# 模拟网络耗时
time.sleep(random.uniform(0.05, 0.3))
# 模拟 10% 概率失败
if random.random() < 0.1:
raise RuntimeError("模拟调用失败")
return {
"id": sample["id"],
"violations": [],
"overall_quality": "good",
}
# ============================================================
# 带重试的处理函数
# 只负责处理,不写文件
# ============================================================
def process_with_retry(sample, max_retries=3):
last_error = None
for attempt in range(1, max_retries + 1):
try:
result = call_llm(sample)
result["attempt"] = attempt
return result
except Exception as e:
last_error = e
# 退避,避免连续重试打爆下游
time.sleep(0.2 * attempt)
# 重试都失败,返回一个失败记录,不抛异常
return {
"id": sample["id"],
"error": str(last_error),
"overall_quality": "failed",
}
def main():
output_path = "results.jsonl"
batch_size = 20 # 每攒 20 条写一次
max_workers = 8 # 并发线程数
total = len(samples)
buffer = [] # 主线程的写缓冲区
written = 0 # 已写盘条数
failed = 0 # 处理失败条数
start_time = time.time()
with open(output_path, "w", encoding="utf-8") as f:
with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor:
# 提交所有任务
future_to_sample = {
executor.submit(process_with_retry, s): s
for s in samples
}
# 按完成顺序消费结果
for future in concurrent.futures.as_completed(future_to_sample):
sample = future_to_sample[future]
try:
# 这里加超时,防止某个 future 卡死
result = future.result(timeout=1)
except Exception as e:
result = {
"id": sample["id"],
"error": f"future 异常: {e}",
"overall_quality": "failed",
}
if result.get("overall_quality") == "failed":
failed += 1
# 放进缓冲区
buffer.append(result)
# 缓冲区满了,主线程统一写盘
if len(buffer) >= batch_size:
for item in buffer:
f.write(json.dumps(item, ensure_ascii=False) + "\n")
f.flush()
written += len(buffer)
buffer.clear()
# 进度显示
elapsed = time.time() - start_time
speed = written / elapsed if elapsed > 0 else 0
print(
f"[进度] 已写 {written}/{total} "
f"失败 {failed} "
f"速度 {speed:.1f} 条/秒 "
f"耗时 {elapsed:.1f}s"
)
# 循环结束后,把缓冲区剩余结果写掉
if buffer:
for item in buffer:
f.write(json.dumps(item, ensure_ascii=False) + "\n")
f.flush()
written += len(buffer)
buffer.clear()
elapsed = time.time() - start_time
print(f"[完成] 共写 {written} 条,失败 {failed} 条,耗时 {elapsed:.1f}s")
print(f"[完成] 结果文件:{os.path.abspath(output_path)}")
if __name__ == "__main__":
main()这个版本解决了什么
| 需求 | 实现方式 |
|---|---|
| 进度显示 | 每写一批打印已写/总数/失败数/速度 |
| 失败重试 | process_with_retry 最多重试 3 次,带退避 |
| 失败不丢 | 重试失败后返回 failed 记录,照样写盘 |
| 批量写入 | 攒够 batch_size 再写,减少 flush 次数 |
| 不让慢任务拖死 | future.result(timeout=1) 给单个结果设上限 |
| 异常隔离 | 每个 future 单独 try/except |
| 主线程写 | 所有 f.write 都在主线程循环内 |
运行输出示例
[进度] 已写 20/100 失败 2 速度 45.3 条/秒 耗时 0.4s
[进度] 已写 40/100 失败 5 速度 52.1 条/秒 耗时 0.8s
[进度] 已写 60/100 失败 7 速度 55.0 条/秒 耗时 1.1s
[进度] 已写 80/100 失败 9 速度 56.8 条/秒 耗时 1.4s
[进度] 已写 100/100 失败 11 速度 58.2 条/秒 耗时 1.7s
[完成] 共写 100 条,失败 11 条,耗时 1.7s
[完成] 结果文件:/path/to/results.jsonl几个设计细节
1. 重试放在工作线程里,不放在主线程
主线程只负责消费结果。如果主线程做重试,会拖慢消费速度,还可能让其他已完成结果堵在队列里。重试属于“处理”,应该在工作线程里做。
2. 失败也写盘
失败记录不能丢。后续可以单独捞出来分析、补跑。写盘时带上 overall_quality: "failed",方便过滤。
3. 缓冲区在主线程
buffer 只在主线程被读写,工作线程不碰,所以不需要锁。这是“主线程控制写”的延伸。
4. timeout 的作用
future.result(timeout=1) 只影响“主线程愿意等这个结果多久”,不会杀掉正在跑的线程。如果某个任务真的卡死,主线程会抛 TimeoutError,记录失败后继续处理其他结果。卡死的线程会在后台继续,但不阻塞整体。
5. 批量写入的权衡
batch_size太小 → flush 频繁,写盘慢。batch_size太大 → 内存占用高,中途崩溃丢的数据多。- 一般 20~100 比较合适,按你的结果大小调整。