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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 06:45:00