Python. Помогите с многопоточностью или многопроцессорностью
Проблема в том, что мои потоки не работают параллельно, а запускаются, после завершения предыдущего (ну по крайней мере на это указывает print и не изменчивая скорость с потоками и без потоков) Я уже столько всего попробовал, но ничего не выходит, кол-во говнокода растёт, а толку от этого никакого... Мне кажется, что всё дело в списках, которые я передаю, но я так же пробовал создавать их копии, чтобы копии срезов были у каждого потока свои и они не обращались к списку из главного потока, но это так же не помогло, скорость всё равно медленная...
import os
from typing import Tuple, List
from threading import Thread
from csv import reader
import time
class Data:
WORKED_THREADS = 0
NUMBER = 0
COINCES_FOUND = 0
PROCESSED_LINES = 0
LAST_LINE_NUMBER = 0
THREAD_LIST: List[Thread] = []
NEW_EMAILS = []
def GetSearchedEmails(fileName: str='./email.txt') -> List[str]:
with open(fileName, 'r', encoding='utf-8') as f:
return [k.strip() for k in f.read().split('\n')]
def GetCsvLines(fileName: str='./source.csv') -> List[List[str]]:
print('Открытие csv-файла...')
with open(fileName, 'r', newline='') as f:
return [row for row in reader(f, delimiter=';')]
def EditEmails(searchedEmails: Tuple[str], emails: List[str]):
for searchedEmail in searchedEmails:
for _ in emails:
if searchedEmail in emails:
Data.COINCES_FOUND += 1
emails.remove(searchedEmail)
return emails
def CheckCurrentInterval(emailLines: List[List[str]], start: int, stop: int, searchedEmails: List[str], index: int):
searchedEmails = searchedEmails.copy()
newEmailLines = []
for line in emailLines[start:stop].copy():
emails = [k.strip() for k in line[index].split(',')]
newLine = line.copy()
newLine[index] = ', '.join(EditEmails(searchedEmails, emails))
newEmailLines.append(newLine)
# Data.PROCESSED_LINES += 1
# print(Data.PROCESSED_LINES)
# Data.WORKED_THREADS -= 1
# Data.NEW_EMAILS += newEmailLines
if __name__ == '__main__':
index = 11
emailLines = GetCsvLines()[:10000]
searchedEmails = GetSearchedEmails()
newEmailLines = []
threadCount = 100
csvLen = len(emailLines)
searchedEmailsLen = len(searchedEmails)
interval = csvLen // threadCount
argList = []
start = 0
stop = 0
from multiprocessing import Process
startT = time.time()
while stop < csvLen:
start = interval * Data.NUMBER
stop = interval * Data.NUMBER + interval
Data.NUMBER += 1
Data.THREAD_LIST.append(Process(target=CheckCurrentInterval, args=(emailLines, start, stop, searchedEmails, index)))
Data.THREAD_LIST[-1].start()
Data.WORKED_THREADS += 1
print(Data.WORKED_THREADS)
while True:
os.system('cls')
print(
f'Длина csv-файла: {csvLen}'
f'\nДлина txt-файла: {searchedEmailsLen}'
f'\nИнтервал: {interval}'
f'\nПоследний stop-индекс: {Data.LAST_LINE_NUMBER}'
f'\nЗапущенные потоки: {Data.WORKED_THREADS} / {threadCount}'
f'\nОбработано строк: {Data.PROCESSED_LINES} / {csvLen}'
f'\nСовпадений найдено: {Data.COINCES_FOUND}'
)
if Data.WORKED_THREADS <= 0 and Data.LAST_LINE_NUMBER >= csvLen:
print('Готово!')
break
time.sleep(1)
print(f'\nКонец: {time.time() - startT} сек.')
UPD.
Вот новый код, вроде работает получше, хз
import os
import time
from csv import reader, writer
from multiprocessing import Pool
from pprint import pprint
from threading import Thread
from typing import List, Tuple
class Data:
WORKED_THREADS = 0
NUMBER = 0
COINCES_FOUND = 0
PROCESSED_LINES = 0
LAST_LINE_NUMBER = 0
THREAD_LIST: List[Thread] = []
NEW_EMAILS = []
def GetSearchedEmails(fileName: str = './email.txt') -> List[str]:
with open(fileName, 'r', encoding='utf-8') as f:
return [k.strip() for k in f.read().split('\n')]
def GetCsvLines(fileName: str = './source.csv') -> List[List[str]]:
print('Открытие csv-файла...')
with open(fileName, 'r', newline='') as f:
return [row for row in reader(f, delimiter=';')]
def EditEmails(args: Tuple[Tuple[str], List[str]] = None, searchedEmails: Tuple[str] = None, emails: List[str] = None):
if args:
searchedEmails = args[0]
emails = args[1]
for searchedEmail in searchedEmails:
for _ in emails:
if searchedEmail in emails:
emails.remove(searchedEmail)
return emails
def CheckCurrentInterval(emailLines: List[List[str]], start: int, stop: int, searchedEmails: List[str], index: int):
searchedEmails = searchedEmails
newEmailLines = []
argList = []
for line in emailLines:
emails = [k.strip() for k in line[index].split(',')]
argList.append((searchedEmails, emails))
print(f'Кол-во аргументов: {len(argList)}')
with Pool(threadCount) as pool:
for res in pool.map(EditEmails, argList):
newLine = line.copy()
newLine[index] = ', '.join(res)
newEmailLines.append(newLine)
return newEmailLines
# Data.WORKED_THREADS -= 1
# Data.NEW_EMAILS += newEmailLines
if __name__ == '__main__':
index = 11
_emailLines = GetCsvLines()
searchedEmails = GetSearchedEmails()
newEmailLines = []
print(len(_emailLines))
additionalEmailLines = []
emailLines = []
for k in _emailLines:
if k[index].strip() != '':
emailLines.append(k)
else:
additionalEmailLines.append(k)
print(len(emailLines))
threadCount = 32
csvLen = len(emailLines)
searchedEmailsLen = len(searchedEmails)
interval = csvLen // threadCount
argList = []
start = 0
stop = 0
kf = 0
startT = time.time()
while stop < csvLen:
if Data.WORKED_THREADS < threadCount:
start = interval * kf
stop = interval * kf + interval
kf += 1
print(f'Пул {kf} запущен!')
newEmailLines += CheckCurrentInterval(emailLines[start:stop], start, stop, searchedEmails, index)
print(f'\nКонец: {time.time() - startT} сек.')
with open('new.csv', 'w', encoding='cp1251', newline='') as f:
writer(f, delimiter=';').writerows(newEmailLines + additionalEmailLines)