You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

多线程处理大数据分块异常:重复循环首批次数据未正常执行

多线程处理分块数据时重复循环首批次、无法处理后续批次的问题

问题描述

我会持续从服务器接收20K或50K条记录的大数据,将其拆分为5K条的分块,计划通过多线程并行调用函数处理这些分块,同时还要处理后续批次的分块数据。但程序持续加载,未按预期处理记录,而是不断重复循环首批次数据,无法处理后续分块。

尝试的代码

def divide_chunks(l, n):
    small_msgs = []
    for i in range(0, len(l), n):
        small_msgs.append(l[i:i + n])
    return small_msgs

def process_data(data, i):
    #process data for chunks
    try:
        # data processing here according to my requirement
        # it may take 20-25 seconds of process that is why am planning for parallel processing
    except exceptions.BadRequest as exc:
        print(json.dumps({'error': str(exc)}))
    return True

#msgs are nothing but bulk data recieving from server continuously am appending to msgs
chunk_msgs = divide_chunks(msgs, 5000)

#clearing msgs to append next data after chunking previous data
msgs.clear()

for n in range(0, len(chunk_msgs)):
    threading.Thread(target=process_data, args=(chunk_msgs[n],n)).start()

问题分析

  1. 缺乏持续处理逻辑:原代码仅执行了一次分块和线程启动操作,没有循环监听新进入msgs的后续批次数据,导致后续接收的数据无法被处理;若外层存在重复执行的循环,又会因msgs未接收完新数据就被分块/清空,出现“重复循环首批次”的假象。
  2. 线程安全问题:如果有其他线程持续向msgs追加数据,原代码中的divide_chunks、msgs.clear()操作无锁保护,会导致数据读取与修改冲突,引发数据错乱或重复处理。
  3. 无并发控制:直接为每个分块启动新线程,数据量较大时会创建大量线程,耗尽系统资源,导致程序陷入加载状态,无法及时处理后续批次。

解决方案

1. 核心改进方向

  • 用线程池控制并发数,避免线程泛滥;
  • 用线程安全队列存储原始数据,解决多线程操作冲突;
  • 增加循环监听逻辑,自动处理后续批次数据。

修正后的示例代码

import threading
import queue
from concurrent.futures import ThreadPoolExecutor
import json
# 请确保导入自定义异常模块
# import your_exceptions_module as exceptions

# 线程安全队列,存储接收的原始数据
data_queue = queue.Queue()
# 线程池,根据CPU核心数或业务需求调整并发数
executor = ThreadPoolExecutor(max_workers=4)

def divide_chunks(l, n):
    small_msgs = []
    for i in range(0, len(l), n):
        small_msgs.append(l[i:i + n])
    return small_msgs

def process_data(data, chunk_id):
    try:
        # 替换为你的实际数据处理逻辑
        print(f"开始处理分块 {chunk_id},数据量:{len(data)}")
        # 模拟处理耗时(可删除)
        # time.sleep(20)
        print(f"分块 {chunk_id} 处理完成")
    except exceptions.BadRequest as exc:
        print(json.dumps({'error': f"分块 {chunk_id} 处理失败: {str(exc)}"}))
    return True

def data_receiver():
    # 模拟持续从服务器接收数据的线程,实际场景替换为真实接收逻辑
    batch_count = 1
    while True:
        # 模拟一批20K数据
        fake_data = [f"record_{batch_count}_{i}" for i in range(20000)]
        for item in fake_data:
            data_queue.put(item)
        print(f"已接收第 {batch_count} 批次数据")
        batch_count += 1
        # 模拟接收间隔(可调整或删除)
        # time.sleep(60)

def data_processor():
    chunk_size = 5000
    while True:
        # 积累到指定数量后分块,不足时超时处理剩余数据
        current_batch = []
        while len(current_batch) < chunk_size:
            try:
                item = data_queue.get(timeout=30)
                current_batch.append(item)
            except queue.Empty:
                if current_batch:
                    break
                continue
        
        if not current_batch:
            continue
        
        # 拆分数据块并提交到线程池
        chunks = divide_chunks(current_batch, chunk_size)
        for idx, chunk in enumerate(chunks):
            executor.submit(process_data, chunk, f"{threading.get_ident()}_{idx}")

# 启动数据接收线程(守护线程随主进程退出)
receiver_thread = threading.Thread(target=data_receiver, daemon=True)
receiver_thread.start()

# 启动数据处理线程
processor_thread = threading.Thread(target=data_processor, daemon=True)
processor_thread.start()

# 保持主进程运行
try:
    while True:
        threading.Event().wait()
except KeyboardInterrupt:
    print("程序终止")
    executor.shutdown()

关键优化点

  • queue.Queue保证数据接收与处理的线程安全,避免数据冲突;
  • ThreadPoolExecutor控制并发数,防止系统资源耗尽;
  • 分离接收与处理线程,职责清晰,持续循环监听新数据;
  • 处理超时逻辑,避免因数据不足导致程序阻塞。

内容的提问来源于stack exchange,提问作者ditil

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.04 19:50:29