How to keep the original order of input when using ThreadPoolExecutor?

Viewed 2304
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()

If I set max_works=1, it will print out from 1 to 10 in order; if I set max_workers = 10, the order could be random:

5
3
10
7
1
8
6
2
4
9

How to keep the original order of the input when using ThreadPoolExecutor to process a list of items as in this example?

6 Answers

You can use the map method of the ThreadExecutor:

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()

Updated answer based on comments requirements:

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()

This mixes two incompatible ideas!

When you use a thread/process/whatever pool, the work will be done in an arbitrary order (largely the result of unrelated system load). Some work may happen at the same time as other work (which is normally the benefit of such a system; parallelization). However, unless you go out of your way to order the results, they will be in whatever order the pool did the work in.

Rather than attempting to "order" the results, consider mapping them back to some collection, like a dictionary, so you can read them back by-key (which may have some order to it).

Just remove as_completed from your code!

executor.submit will return Future object, and you append it to list in order.

But as_completed will generate Future objects in completed order.

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(threads):
            result = task.result()
            #print(result)
            all_results.append(result)

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

process()

Per @marlon's request, here is a variation of the @rorra solution that uses itertools.repeat to reduce some of the parameter passing complexity.

Example:

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()

Output:

3
4
5
6
7
8
9
10
11
12

This can be one of the ways to get the desired result

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()

I know that this question was posted a some time ago, but I had a similar problem and I could sovle it creating a dict using the result of the method id() as key, and the chunk processed for the future function (in my case if one fails I need to try process this chunk one more time). Ex:

control_dict = {}
future = executor.submit(my_func, chunk))
control_dict[id(future)] = chunk

You could use the same strategy to "find" the order.

Related