Вопросы о многопоточности для чайников в python

Недавно наткнулся на статью о параллельных вычислениях. Раньше всегда все делал через циклы в одном потоке. Решил опробовать на своем коде. Пример кода:

import queue
import threading
data=[1000,500,2000,500,5000,1205,12412,124512,...]
data= pd.DataFrame(data, columns=["value"])
period=1
def mov_avg(df):
    global period
    period=period+1
    y=df.apply(lambda x: x.value, axis=1).rolling(period).sum() / period
while period<=20:
    worker1 = threading.Thread(target=mov_avg(df=data))
    worker2 = threading.Thread(target=mov_avg(df=data))
worker1.start()
worker2.start()

И возникли вопросы:

  1. Почему при сравнении времени работы worker1,worker2(работают вместе) и worker1(работет один), worker1 не отстает от работы их вместе(бывает и опережает). Как это исправить?
  2. Как создать генератор потоков чтобы можно было запустить максимальное кол-во процессов? Чем это лимитировано?
  3. Как передавать значение из одного потока в другой?
  • Пример 1: генератор создал n кол-во потоков для подсчета скользящих средних. Какой-то случайный поток освободился от просчета и чтобы не ждать завершения работы других потоков, он получает следующее непросчитанное значение и приступает к подсчету.(асинхронный)
  • Пример 2: генератор создал n кол-во потоков для подсчета скользящих средних. Какой-то случайный поток освободился от просчета и он ЖДЕТ завершения других потоков, после чего он получает следующее непросчитанное значение и приступает к подсчету.(синхронный)
  1. В своем коде я использовал глобальную переменную для передачи значения между потоками. Каковы различия и преимущества моего способа и способа из вопроса №3?

Ответы (1 шт):

Автор решения: MaxU

Попробуйте так:

from multiprocessing import Pool

def parallelize_dataframe(df, func, n_cores=4):
    df_split = np.array_split(df, n_cores)
    pool = Pool(n_cores)
    df = pd.concat(pool.map(func, df_split))
    pool.close()
    pool.join()
    return df
→ Ссылка