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

如何修改Python生产者消费者队列实现批量存取适配机器学习模型

批量处理队列数据适配机器学习模型需求

要实现Producer批量存入、Consumer批量获取数据以适配模型一次性处理2000条的需求,只需调整队列操作逻辑,将单条数据改为批次数据即可。以下是修改后的完整代码:

import time
from threading import Thread
from queue import Queue

def read_data():
    torque_data_queryset = Torque.get_torque_data()
    return list(torque_data_queryset)

def producer(queue, batch_size=2000):
    print('Producer: Running')
    data = read_data()
    # 将数据按指定批次大小拆分
    for i in range(0, len(data), batch_size):
        batch = data[i:i+batch_size]
        time.sleep(1)  # 模拟数据准备延迟,可根据实际情况调整或移除
        queue.put(batch)
        print(f'> Producer added batch of {len(batch)} items')
    
    queue.put(None)  # 发送结束信号
    print('Producer: Done')

def consumer(queue):
    print('Consumer: Running')
    while True:
        batch = queue.get()
        if batch is None:
            break
        
        time.sleep(1)  # 模拟模型处理延迟,可替换为实际模型调用
        # 这里调用你的机器学习模型处理批次数据
        # model.predict(batch)
        print(f'> Consumer processed batch of {len(batch)} items')
    
    print('Consumer: Done')

def thread_function():
    queue = Queue()
    
    consumer_thread = Thread(target=consumer, args=(queue,))
    consumer_thread.start()
    
    producer_thread = Thread(target=producer, args=(queue,))
    producer_thread.start()
    
    producer_thread.join()
    consumer_thread.join()

def thread_result(request):
    thread_function()
    return render(request, 'async_processing_result.html', {'message': 'Threads completed successfully!'})

关键修改说明

  • Producer端:通过切片将全量数据拆分为batch_size(默认2000)条一组的批次,批量放入队列,避免逐条操作的冗余开销,同时精准匹配模型的批量处理要求。
  • Consumer端:每次从队列取出完整批次,直接传入模型处理(代码中已预留模型调用位置),无需逐条处理数据。
  • 结束信号:保留None作为数据传输完成的标识,确保Consumer能正确终止循环。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 17:23:20