Многопоточная запись в файл в python
Есть огромный лист B и огромный словарь A. Для каждого элемента B нужно проверить наличие его в A и в случае наличия записать в файл. Я сделал это таким образом и получаю неверный результат работы.
from multiprocessing import Pool
def func(i):
if i in A:
ff.write(str(A[i]) + "\n")
if __name__ == '__main__':
A = {2 * x : 0 for x in range(5000000)}
B = [x for x in range(10000000)]
ff = open("task.del", "w")
pool = Pool(4)
pool.map(func, B)
pool.close()
pool.join()
В результате в файле неверное количество ответов (должно быть около 5 миллионов +- 1, а получается, к примеру на 4 процессах 4997120). Работает эта система, к тому же, еще и дольше чем однопоточная. Я понимаю, что какие-то накладки по скорости будут, но в данном случае, подозреваю, проблема в постоянном перехвате дескриптора файла для записи. Посоветуйте, как лучше.
Ответы (2 шт):
Могу только предложить всё же писать в файл в одном потоке, а данные для записи генерировать в несколько потоков. Тогда всё нормально работает. Когда питон создаёт потоки он копирует в каждый поток все переменные основного потока и, видимо, поэтому с дескриптором файла получается какая-то накладка, он так не умеет работать, чтобы один дескриптор (вернее его копия) работал из нескольких потоков сразу.
from multiprocessing import Pool
def func(i):
if i in A:
return A[i]
if __name__ == '__main__':
A = {2 * x : 0 for x in range(5000000)}
B = [x for x in range(10000000)]
with open("task.del", "w") as ff:
with Pool(4) as pool:
for r in pool.map(func, B):
if r != None:
ff.write(str(r) + "\n")
Что конкретно многопоточно?
Можно например, данные формировать в отдельном процессе, и выдавать их по мере формирования, а в Главном процессе получать их и записывать в файл.
Или можно многопоточно формировать сразу все данные, а потом их записывать однопоточно.
И не следует пытатся многопоточно записывать в сам файл.
t1 = time.monotonic() # 0 сек
A = {2 * x : 0 for x in range(5000000)}
B = [x for x in range(10000000)]
t2 = time.monotonic() # + 2.0 сек
s = A.keys() & set(B)
t3 = time.monotonic() # + 1.3 сек
s = '\n'.join(map(str, s))
t4 = time.monotonic() # + 1.7 сек
with open('task.del', 'w') as f:
f.write(s)
t5 = time.monotonic() # + 0.6 сек
print((t2 - t1), (t3 - t2), (t4 - t3), (t5 - t4))