如何修改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
相关产品推荐
相关产品推荐

