全栈开发

多线程处理数据,主线程控制写数据

一、为什么需要这个模式

写并发程序时,一个很常见的需求是:

  • 数据很多,单线程处理太慢,想用多线程加速。
  • 处理结果需要写文件、写数据库。
  • 但多线程直接写,会出现内容交错、顺序混乱、文件损坏、计数错误。

一个自然的解法是:让工作线程只负责计算,把结果交回主线程,由主线程统一写。

这就是“多线程处理数据,主线程控制写数据”。

它的核心不是靠锁,而是靠职责分离:并行的是处理,串行的是写。

二、先理解两种写法

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_1data_10,而是谁先处理完谁先写。但写入动作全部由主线程串行完成,文件内容不会交错。

四、几个关键点

1. 为什么不需要锁

这个模式里,fwritten_countresult 都只在主线程访问。工作线程执行 process(data),只 return 结果,不碰任何共享资源。

锁的作用是保护共享资源。这里没有共享写,自然不需要锁。用结构设计消除竞争,比加锁更好。

ThreadPoolExecutoras_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 比较合适,按你的结果大小调整。