Python의 다중 처리 및 다중 스레딩
1) 디스크 (~ 95 % 시간)에서 매우 큰 파일을 읽고 2) 비교적 작은 출력 (~ 5 % 시간)을 처리하고 제공하는 파이썬 프로그램이 있습니다. 이 프로그램은 TeraBytes의 파일에서 실행됩니다.
이제 Multi Processing 및 Multi Threading을 활용하여이 프로그램을 최적화하려고합니다. 내가 실행중인 플랫폼은 가상 머신에 4 개의 프로세서가있는 가상 머신입니다.
4 개의 프로세스 (프로세서와 동일)를 실행하는 스케줄러 프로세스를 계획하고 각 프로세스에는 대부분의 부분이 I / O이므로 일부 스레드가 있어야합니다. 각 스레드는 1 개의 파일을 처리하고 결과를 메인 스레드에보고하고 IPC를 통해 스케줄러 프로세스에 다시보고합니다. 스케줄러는 이러한 항목을 대기열에 넣고 결국 순서대로 디스크에 기록 할 수 있습니다.
그렇다면 그러한 시나리오를 위해 생성 할 프로세스와 스레드의 수를 어떻게 결정합니까? 무엇이 최선의 조합인지 알아내는 수학적 방법이 있습니까?
감사합니다
답변
나는 그것을 당신이하는 일의 역으로 배열 할 것이라고 생각합니다. 즉, 결과 생성을 담당하는 특정 크기의 스레드 풀을 만듭니다. 이 풀에 제출 된 작업은 CPU 바운드 작업 부분을 제출하기 위해 작업자 스레드에서 사용할 수있는 프로세서 풀의 인수로 전달됩니다. 즉, 스레드 풀 작업자는 기본적으로 모든 디스크 관련 작업을 수행하고 CPU 집약적 인 작업을 프로세서 풀로 넘겨줍니다.
프로세서 풀의 크기는 환경에있는 프로세서 수 여야합니다. 스레드 풀에 정확한 크기를 지정하는 것은 어렵습니다. 수익 감소의 법칙이 적용되기 전에 처리 할 수있는 동시 디스크 작업 수에 따라 다릅니다. 또한 메모리에 따라 다릅니다. 풀이 클수록, 특히 처리를 위해 전체 파일을 메모리로 읽어야하는 경우에 사용되는 메모리 리소스가 많아집니다. 따라서이 값을 실험해야 할 수도 있습니다. 아래 코드는 이러한 아이디어를 설명합니다. 스레드 풀에서 얻는 것은 작은 프로세서 풀을 사용했을 때 얻을 수있는 것보다 더 큰 I / O 작업이 겹치는 것입니다.
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
from functools import partial
import os
def cpu_bound_function(arg1, arg2):
...
return some_result
def io_bound_function(process_pool_executor, file_name):
with open(file_name, 'r') as f:
# Do disk related operations:
. . . # code omitted
# Now we have to do a CPU-intensive operation:
future = process_pool_executor.submit(cpu_bound_function, arg1, arg2)
result = future.result() # get result
return result
file_list = [file_1, file_2, file_n]
N_FILES = len(file_list)
MAX_THREADS = 50 # depends on your configuration on how well the I/O can be overlapped
N_THREADS = min(N_FILES, MAX_THREADS) # no point in creating more threds than required
N_PROCESSES = os.cpu_count() # use the number of processors you have
with ThreadPoolExecutor(N_THREADS) as thread_pool_executor:
with ProcessPoolExecutor(N_PROCESSES) as process_pool_executor:
results = thread_pool_executor.map(partial(io_bound_function, process_pool_executor), file_list)
중요 사항 :
훨씬 더 간단한 또 다른 접근 방식은 크기가 사용자가 보유한 CPU 프로세서 수 (예 : 25 ) 보다 큰 단일 프로세서 풀을 보유하는 것입니다. 작업자 프로세스는 I / O 및 CPU 작업을 모두 수행합니다. CPU보다 많은 프로세스가 있더라도 많은 프로세스가 CPU 집약적 인 작업이 실행될 수 있도록 I / O가 완료되기를 기다리는 대기 상태에 있습니다.
이 접근법의 단점은 N 개의 프로세스를 생성 할 때의 오버 헤드가 N 개의 스레드 + 적은 수의 프로세스를 생성 할 때의 오버 헤드보다 훨씬 크다는 것입니다. 그러나 풀에 제출 된 작업의 실행 시간이 점점 더 길어짐에 따라이 증가 된 오버 헤드는 전체 실행 시간에서 점점 더 작은 비율이됩니다. 따라서 작업이 사소하지 않은 경우 이것은 합리적으로 수행되는 단순화가 될 수 있습니다.
업데이트 : 두 가지 접근 방식의 벤치 마크
크기가 약 10,000KB 인 24 개의 파일을 처리하는 두 가지 접근 방식에 대해 몇 가지 벤치 마크를 수행했습니다 (실제로 이들은 각각 8 번 처리 된 3 개의 서로 다른 파일이므로 캐싱이 수행되었을 수 있습니다).
방법 1 (스레드 풀 + 프로세서 풀)
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
from functools import partial
import os
from math import sqrt
import timeit
def cpu_bound_function(b):
sum = 0.0
for x in b:
sum += sqrt(float(x))
return sum
def io_bound_function(process_pool_executor, file_name):
with open(file_name, 'rb') as f:
b = f.read()
future = process_pool_executor.submit(cpu_bound_function, b)
result = future.result() # get result
return result
def main():
file_list = ['/download/httpd-2.4.16-win32-VC14.zip'] * 8 + ['/download/curlmanager-1.0.6-x64.exe'] * 8 + ['/download/Element_v2.8.0_UserManual_RevA.pdf'] * 8
N_FILES = len(file_list)
MAX_THREADS = 50 # depends on your configuration on how well the I/O can be overlapped
N_THREADS = min(N_FILES, MAX_THREADS) # no point in creating more threds than required
N_PROCESSES = os.cpu_count() # use the number of processors you have
with ThreadPoolExecutor(N_THREADS) as thread_pool_executor:
with ProcessPoolExecutor(N_PROCESSES) as process_pool_executor:
results = list(thread_pool_executor.map(partial(io_bound_function, process_pool_executor), file_list))
print(results)
if __name__ == '__main__':
print(timeit.timeit(stmt='main()', number=1, globals=globals()))
방법 2 (프로세서 풀 전용)
from concurrent.futures import ProcessPoolExecutor
from math import sqrt
import timeit
def cpu_bound_function(b):
sum = 0.0
for x in b:
sum += sqrt(float(x))
return sum
def io_bound_function(file_name):
with open(file_name, 'rb') as f:
b = f.read()
result = cpu_bound_function(b)
return result
def main():
file_list = ['/download/httpd-2.4.16-win32-VC14.zip'] * 8 + ['/download/curlmanager-1.0.6-x64.exe'] * 8 + ['/download/Element_v2.8.0_UserManual_RevA.pdf'] * 8
N_FILES = len(file_list)
MAX_PROCESSES = 50 # depends on your configuration on how well the I/O can be overlapped
N_PROCESSES = min(N_FILES, MAX_PROCESSES) # no point in creating more threds than required
with ProcessPoolExecutor(N_PROCESSES) as process_pool_executor:
results = list(process_pool_executor.map(io_bound_function, file_list))
print(results)
if __name__ == '__main__':
print(timeit.timeit(stmt='main()', number=1, globals=globals()))
결과 :
(나는 8 개의 코어가 있습니다)
스레드 풀 + 프로세서 풀 : 13.5 초 프로세서 풀 단독 : 13.3 초
결론 : 먼저 모든 것에 프로세서 풀을 사용하는 것보다 더 간단한 방법을 시도해 보겠습니다. 이제 까다로운 부분은 생성 할 최대 프로세스 수를 결정하는 것입니다. 이는 원래 질문의 일부였으며 CPU 집약적 인 계산 만했을 때 간단한 답을 얻었습니다. 읽고있는 파일의 수가 너무 많지 않다면 요점은 문제입니다. 파일 당 하나의 프로세스를 가질 수 있습니다. 그러나 수백 개의 파일이있는 경우 풀에 수백 개의 프로세스가있는 것을 원하지 않을 것입니다 (생성 할 수있는 프로세스 수에 대한 상한선도 있으며 이러한 불쾌한 메모리 제약도 있습니다). 정확한 번호를 줄 수있는 방법은 없습니다. 많은 수의 파일이있는 경우 작은 풀 크기로 시작하여 더 이상 이점이 없을 때까지 계속 증가합니다 (물론 이러한 테스트에 대해 일부 최대 수보다 많은 파일을 처리하고 싶지 않을 것입니다. 실제 실행을위한 좋은 풀 크기를 결정하기 만하면 영원히 실행됩니다).
병렬 처리의 경우 : 나는 이 질문을 보았고 수락 된 대답을 인용했습니다.
실제로는 최적의 스레드 수를 찾기가 어려울 수 있으며 그 수조차도 프로그램을 실행할 때마다 달라질 수 있습니다. 따라서 이론적으로 최적의 스레드 수는 컴퓨터에있는 코어 수입니다. 코어가 "하이퍼 스레드"인 경우 (Intel이 부르는대로) 각 코어에서 2 개의 스레드를 실행할 수 있습니다. 그런 다음이 경우 최적의 스레드 수는 컴퓨터의 코어 수의 두 배입니다.
다중 처리의 경우 : 누군가 여기 에서 비슷한 질문 을했고 수락 된 답변은 다음과 같습니다.
모든 스레드 / 프로세스가 실제로 CPU 바운드 인 경우 CPU가 코어를보고하는만큼 많은 프로세스를 실행해야합니다. HyperThreading으로 인해 각 물리적 CPU 코어는 여러 가상 코어를 표시 할 수 있습니다.
multiprocessing.cpu_count가상 코어 수를 가져 오려면 호출하십시오 .
스레드 1 개 중 p 만 CPU 바운드 인 경우 p를 곱하여 해당 수를 조정할 수 있습니다. 예를 들어 프로세스의 절반이 CPU 바운드 (p = 0.5)이고 각각 4 개의 코어와 2x HyperThreading이있는 2 개의 CPU가있는 경우 0.5 * 2 * 4 * 2 = 8 개 프로세스를 시작해야합니다.
여기서 핵심은 어떤 머신을 사용하고 있는지 이해하는 것입니다. 그로부터 코드 실행을 분할하기 위해 거의 최적의 스레드 / 프로세스 수를 선택할 수 있습니다. 그리고 스크립트를 실행할 때마다 조금씩 달라지기 때문에 거의 최적 이라고 말했습니다 . 따라서 수학적 관점에서이 최적의 숫자를 예측하는 것은 어려울 것입니다.
특정 상황에서 머신에 4 개의 코어가있는 경우 최대 4 개의 스레드 만 생성 한 다음 분할하는 것이 좋습니다.
- 메인 스레드에 1.
- 3 파일 읽기 및 처리.
여러 프로세스를 사용하여 IO 성능 속도를 높이는 것은 좋은 생각이 아닐 수 있습니다. 이 기능 과 그 아래 의 샘플 코드 를 확인하여 도움이되는지 확인하십시오.
한 가지 아이디어는 스레드가 파일을 읽기만하고 (잘 이해했다면 파일이 하나뿐 임) 독립적 인 부분 (예 : 행)을 메시지와 함께 큐에 넣는 것입니다.
메시지는 4 개의 스레드로 처리 할 수 있습니다. 이러한 방식으로 프로세서 간의로드를 최적화 할 수 있습니다.
강력한 I / O 바인딩 프로세스 (설명하는 내용과 같이)에서는 반드시 멀티 스레딩이나 멀티 프로세싱이 필요하지 않습니다. OS에서 고급 I / O 프리미티브를 사용할 수도 있습니다.
예를 들어 Linux에서 적절한 크기의 가변 버퍼와 함께 커널에 읽기 요청을 제출하고 버퍼가 채워지면 알림을받을 수 있습니다. 이것은 순수 파이썬 바인딩을 작성한 AIO API를 사용하여 수행 할 수 있습니다 : python-libaio ( libaio on pypi)) 또는 CFFI 파이썬 바인딩 이있는 것으로 보이는 최신 io_uring API ( liburing on pypy) (나는 io_uring도이 파이썬 바인딩도 사용하지 않았습니다).
이렇게하면 레벨에서 병렬 처리의 복잡성이 제거되고 OS / 사용자 영역 컨텍스트 전환의 수를 줄일 수 있으며 (cpu 시간을 더욱 단축) OS가 수행하려는 작업에 대해 더 많이 알 수 있으므로 스케줄링 기회를 얻을 수 있습니다. IO가 더 효율적입니다 (가상화 된 환경에서 직접 시도하지는 않았지만 데이터 사본 수를 줄이면 놀라지 않을 것입니다).
물론 단점은 프로그램이 실행중인 OS에 더 밀접하게 연결되어 다른 프로그램에서 실행하려면 더 많은 노력이 필요하다는 것입니다.