【问题标题】:How to keep the original order of input when using ThreadPoolExecutor?使用ThreadPoolExecutor时如何保持原来的输入顺序?
【发布时间】:2021-04-21 04:27:13
【问题描述】:
from concurrent.futures import ThreadPoolExecutor, as_completed

def add_one(number, n):
    return number + 1 + n

def process():
    all_numbers = []
    for i in range(0, 10):
        all_numbers.append(i)

    threads = []
    all_results = []
    with ThreadPoolExecutor(max_workers=10) as executor:
        for number in all_numbers:
            threads.append(executor.submit(add_one, number))

        for index, task in enumerate(as_completed(threads)):
            result = task.result()
            #print(result)
            all_results.append(result)

    for index, result in enumerate(all_results):
        print(result)

process()

如果我设置 max_works=1,它会从 1 到 10 依次打印出来;如果我设置 max_workers = 10,则顺序可能是随机的:

5
3
10
7
1
8
6
2
4
9

如本例中使用 ThreadPoolExecutor 处理项目列表时如何保持输入的原始顺序?

【问题讨论】:

    标签: python threadpoolexecutor


    【解决方案1】:

    可以使用ThreadExecutor的map方法:

    from concurrent.futures import ThreadPoolExecutor, as_completed
    
    def add_one(number):
        return number + 1
    
    def process():
        all_numbers = []
        for i in range(0, 10):
            all_numbers.append(i)
    
        all_results = []
        with ThreadPoolExecutor(max_workers=10) as executor:
            for i in executor.map(add_one, all_numbers):
                print(i)
                all_results.append(i)
    
        for index, result in enumerate(all_results):
            print(result)
    
    process()
    

    根据 cmets 要求更新答案:

    from concurrent.futures import ThreadPoolExecutor, as_completed
    
    def add_one(args):
        return args[0] + 1 + args[1]
    
    def process():
        all_numbers = []
        for i in range(0, 10):
            all_numbers.append([i, 2])
    
        all_results = []
        with ThreadPoolExecutor(max_workers=10) as executor:
            for i in executor.map(add_one, all_numbers):
                print(i)
                all_results.append(i)
    
        for index, result in enumerate(all_results):
            print(result)
    
    process()
    

    【讨论】:

    • 我的实际 add_one() 多了一个参数,'n'。在这种情况下,如何将附加参数与 all_numbers 一起传递给 executor.map() 函数?我编辑了 add_one() 函数。
    • @marlon 我根据你的 cmets 更新了答案
    • 使用来自itertoolsrepeat 可以避免一些额外的参数args 的复杂性。它允许add_one(number, n) 使用map.executor(add_one, all_numbers, itertools.repeat(2))
    • @rhurwitz,repeat(2) 是否意味着函数 add_one 有两个参数?
    • @rhurwitz 你能写一个完整的答案吗?在你的 cmets 中,应该是 map.executor() 还是 executor.map()?
    【解决方案2】:

    这混合了两个不相容的想法!

    当您使用线程/进程/任何池时,工作将以任意顺序完成(主要是不相关的系统负载的结果)。一些工作可能与其他工作同时发生(这通常是这种系统的好处;并行化)。但是,除非您不遗余力地对结果进行排序,否则它们将按照池执行工作的任何顺序进行。

    与其尝试对结果“排序”,不如考虑将它们映射回某个集合,例如字典,这样您就可以按键读回它们(可能有某种顺序)。

    【讨论】:

      【解决方案3】:

      根据@marlon 的要求,这里是@rorra 解决方案的一个变体,它使用itertools.repeat 来降低一些参数传递的复杂性。

      例子:

      from concurrent.futures import ThreadPoolExecutor, as_completed
      import time
      import itertools
      
      def add_one(number, n):
          return number + 1 + n
      
      def process():
          all_numbers = list(range(0, 10))
      
          with ThreadPoolExecutor(max_workers=10) as executor:
              
              for result in executor.map(add_one, all_numbers, itertools.repeat(2)):
                  print(result)
      
      process()
      

      输出:

      3
      4
      5
      6
      7
      8
      9
      10
      11
      12
      

      【讨论】:

      【解决方案4】:

      这可能是获得所需结果的方法之一

      from concurrent.futures import ThreadPoolExecutor, as_completed
      
      
      def add_one(number, index):
          return number + 1, index
      
      
      def process():
          all_numbers = []
          for i in range(0, 10):
              all_numbers.append(i)
      
          threads = []
          all_results = []
          with ThreadPoolExecutor(max_workers=10) as executor:
              for index, number in enumerate(all_numbers):
                  threads.append(executor.submit(add_one, number, index))
              for task in as_completed(threads):
                  result, index = task.result()
                  all_results.append([result, index])
              all_results = sorted(all_results, key=lambda x: x[-1])
      
          for index, result in enumerate(all_results):
              print(result[0])
      
      
      process()
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2020-05-30
        • 1970-01-01
        • 1970-01-01
        • 2016-11-19
        • 1970-01-01
        • 1970-01-01
        • 2014-06-18
        • 2021-07-09
        相关资源
        最近更新 更多