如何用Python多线程实现主从线程聚合求和随机整数数组
线程改造:拆分从属线程实现计算结果汇总
需求说明
我需要实现一个主线程(master thread),用于聚合并汇总多个数量可变的从属线程(slave thread)的计算结果。每个从属线程需生成包含1000个随机整数的数组并计算其和。
现有代码情况
基于数组逻辑的可行实现
from random import randint from multiprocessing import * from queue import Queue q = Queue() def array(n): sumOfSlaves = 0 randomList = [0 for i in range(100)] sizeOfArray = 1000 // n if 1000 % n != 0: for i in range((1000 % n)): randomList = [(randint(0, 1000)) for i in range(int(1000 / n) + 1)] print("Objects in array: ", randomList) print("Size of array: ", len(randomList)) print("Sum of objects in array: ", sum(randomList)) print(sum(randomList)) sumOfSlaves = sumOfSlaves + sum(randomList) print("Aggregate of slave threads so far: ", sumOfSlaves) for i in range(n - int(1000 % n)): randomList = [(randint(0, 1000)) for i in range(int(1000 / n))] print("Objects in array: ", randomList) print("Size of array: ", len(randomList)) print("Sum of objects in array: ", sum(randomList)) print(sum(randomList)) sumOfSlaves = sumOfSlaves + sum(randomList) print("Aggregate of slave threads so far: ", sumOfSlaves) return () else: for i in range(n): randomList = [(randint(0, 1000)) for i in range(sizeOfArray)] print("Objects in array: ", randomList) print("Size of array: ", len(randomList)) print("Sum of objects in array: ", sum(randomList)) sumOfSlaves = sumOfSlaves + sum(randomList) print("Aggregate of slave threads so far: ", sumOfSlaves) return () if __name__ == '__main__': o = int(input("please enter how many slave processor you need: ")) array(o)
当前线程版本代码(未拆分从属线程)
from random import randint import logging import threading import time def master_thread(n): randomList = [0 for i in range(100)] sizeOfArray = 1000 // n sumOfSlaves = 0 if 1000 % n != 0: for i in range((1000 % n)): randomList = [(randint(0, 1000)) for i in range(int(1000 / n) + 1)] print("Objects in array: ", randomList) print("Size of array: ", len(randomList)) print("Sum of objects in array: ", sum(randomList)) print(sum(randomList)) sumOfSlaves = sumOfSlaves + sum(randomList) for i in range(n - int(1000 % n)): randomList = [(randint(0, 1000)) for i in range(int(1000 / n))] print("Objects in array: ", randomList) print("Size of array: ", len(randomList)) print("Sum of objects in array: ", sum(randomList)) print(sum(randomList)) sumOfSlaves = sumOfSlaves + sum(randomList) return () else: for i in range(n): randomList = [(randint(0, 1000)) for i in range(sizeOfArray)] print("Objects in array: ", randomList) print("Size of array: ", len(randomList)) print("Sum of objects in array: ", sum(randomList)) sumOfSlaves = sumOfSlaves + sum(randomList) print("Aggregate of slave threads so far: ", sumOfSlaves) return () if __name__ == "__main__": o = int(input("please enter how many slave processor you need: ")) x = threading.Thread(target=master_thread, args=(o,)) logging.info("Main : before running thread") x.start()
改造方案
核心思路是将生成随机数组+计算和的逻辑抽离为独立的从属线程函数,主线程负责创建线程、等待线程完成、汇总结果。以下提供两种实现方式:
方式一:基础线程+队列实现
from random import randint import threading from queue import Queue # 从属线程函数:负责生成指定长度的随机数组并计算和,将结果存入队列 def slave_thread(array_length, result_queue): random_list = [randint(0, 1000) for _ in range(array_length)] array_sum = sum(random_list) print(f"从属线程生成数组长度: {len(random_list)}") print(f"该数组的和: {array_sum}") result_queue.put(array_sum) # 主线程函数:负责创建从属线程、等待执行、汇总结果 def master_thread(n): result_queue = Queue() threads = [] total_sum = 0 # 计算每个从属线程的数组长度 base_length = 1000 // n remainder = 1000 % n # 创建处理余数的线程(每个多处理1个元素) for _ in range(remainder): t = threading.Thread(target=slave_thread, args=(base_length + 1, result_queue)) threads.append(t) t.start() # 创建处理基础长度的线程 for _ in range(n - remainder): t = threading.Thread(target=slave_thread, args=(base_length, result_queue)) threads.append(t) t.start() # 等待所有从属线程执行完成 for t in threads: t.join() # 从队列中取出所有结果并汇总 while not result_queue.empty(): total_sum += result_queue.get() print(f"\n所有从属线程计算结果汇总: {total_sum}") if __name__ == "__main__": o = int(input("请输入从属线程的数量: ")) master_thread(o)
方式二:使用ThreadPoolExecutor简化实现
concurrent.futures.ThreadPoolExecutor可以自动管理线程池,无需手动处理线程创建和队列:
from random import randint from concurrent.futures import ThreadPoolExecutor # 从属任务函数:生成数组并返回和 def slave_task(array_length): random_list = [randint(0, 1000) for _ in range(array_length)] array_sum = sum(random_list) print(f"从属任务生成数组长度: {len(random_list)}") print(f"该数组的和: {array_sum}") return array_sum # 主线程任务:分配任务、收集结果、汇总 def master_task(n): base_length = 1000 // n remainder = 1000 % n # 构建所有任务的数组长度列表 task_lengths = [base_length + 1] * remainder + [base_length] * (n - remainder) total_sum = 0 # 创建线程池并执行任务 with ThreadPoolExecutor(max_workers=n) as executor: # 提交所有任务并迭代获取结果 results = executor.map(slave_task, task_lengths) for res in results: total_sum += res print(f"\n所有从属任务计算结果汇总: {total_sum}") if __name__ == "__main__": o = int(input("请输入从属线程的数量: ")) master_task(o)
改造说明
- 抽离独立的从属线程/任务函数,每个线程只负责自己的计算逻辑,实现真正的多线程并行。
- 使用线程安全的队列(方式一)或线程池的结果收集机制(方式二)传递计算结果,避免线程间数据竞争。
- 主线程统一等待所有从属线程完成后再汇总结果,确保数据完整性。
内容的提问来源于stack exchange,提问作者QRebourneQ
相关产品推荐
相关产品推荐

