Многопроцессный код(multiprocessing) python не работает как надо
У меня есть код:
import multiprocessing as mpl
import stocker
from time import time
import threading
def fill_d(data):
d, name, lock = data
val = stocker.predict.tomorrow(name) # val = [predicted price, error(%), date of the next business day]
try:
lock.acquire()
data = d[name][val[-1]]
data.append(val[0])
d[name][val[-1]] = data
except KeyError:
d[name] = {}
d[name][val[-1]] = [val[0]]
finally:
lock.release()
def runner(name, d, l):
t = time()
with mpl.Pool() as p:
p.map(fill_d, [(d, name, l) for _ in range(5)])
print(f'{round(time() - t, 2)}s')
stock = ['AAPL', 'AMZN']
if __name__ == '__main__':
manager = mpl.Manager()
d = manager.dict()
threads = []
for name in stock:
t = threading.Thread(target=runner, args=(name, d, manager.Lock()))
threads.append(t)
print(f'[+] Start thread {t.name}')
t.start()
'''Я думаю что каждый поток должен иметь свой lock, в противном случае потоки будут "ждать друг-друга"
'''
for i in threads:
i.join()
print(d)
Который должен выдавать такой результат:
{'AAPL': {'2020-11-03': [109.56, 109.13, 109.48, 109.2, 108.44]}, 'AMZN': {'2020-11-03': [3014.6, 3003.31, 3009.34, 3011.31, 2994.48]}}
А я получаю такой результат:
{'AAPL': {}, 'AMZN': {}}
Я убрал из кода multiprocessing и threading, и всё заработало.
import stocker
d, name = {}, 'AAPL'
def fill(n):
global d, name
for i in range(n):
val = stocker.predict.tomorrow(name) # val = [predicted price, error(%), date of the next business day]
try:
data = d[name][val[-1]]
data.append(val[0])
d[name][val[-1]] = data
except KeyError:
d[name] = {}
d[name][val[-1]] = [val[0]]
fill(5)
name = 'AMZN'
fill(5)
print(
d) # {'AAPL': {'2020-11-03': [109.56, 109.13, 109.48, 109.2, 108.44]}, 'AMZN': {'2020-11-03': [3014.6, 3003.31, 3009.34, 3011.31, 2994.48]}}
Чего же я хочу?
Я хочу чтобы два потока заполняли d некоторыми данными через multiprocessing. Или просто заставить первый код работать.
Спасибо.
Ответы (1 шт):
Автор решения: gil9red
→ Ссылка
Попробуйте просто из функции для map возвращать значение, тогда они, значения, вернутся как список.
Пример:
import multiprocessing as mpl
import stocker
from time import time
import threading
def fill_d(data):
name = data
val = stocker.predict.tomorrow(name) # val = [predicted price, error(%), date of the next business day]
# val = [266.07, 1.276, '2019-11-11']
return val
def runner(name, d):
t = time()
with mpl.Pool() as p:
d[name] = p.map(fill_d, [[name] for _ in range(5)])
print(f'{round(time() - t, 2)}s')
stock = ['AAPL', 'AMZN']
if __name__ == '__main__':
manager = mpl.Manager()
d = manager.dict()
threads = []
for name in stock:
t = threading.Thread(target=runner, args=(name, d))
threads.append(t)
print(f'[+] Start thread {t.name}')
t.start()
for i in threads:
i.join()
print(d)