Computer >> 컴퓨터 >  >> 프로그래밍 >> Python

Python 멀티프로세싱 완전 정복: 프로세스 동기화부터 풀링까지

파이썬에서 병렬 처리를 구현할 때 가장 많이 활용되는 것이 바로 multiprocessing 패키지입니다. 이 글에서는 프로세스 간 동기화 방법과 Pool 클래스를 이용한 작업 풀링 기법을 예제 코드와 함께 자세히 살펴보겠습니다.

프로세스 간 동기화(Synchronization)

multiprocessing은 API를 통해 프로세스 생성을 지원하는 패키지로, 로컬 환경과 원격 환경 모두에서 동시성(Concurrency) 처리에 사용할 수 있습니다. 이 모듈을 활용하면 프로그래머가 하나의 머신에서 여러 프로세서를 효율적으로 사용할 수 있으며, Windows와 UNIX 계열 운영체제에서 모두 동작합니다.

또한 이 패키지에는 threading 모듈과 동등한 모든 동기화 프리미티브(synchronization primitives)가 내장되어 있어, 여러 프로세스가 공유 자원에 접근하는 상황에서도 안정적인 코드 작성이 가능합니다.

동기화 예제 코드

from multiprocessing import Process, Lock

def my_function(x, y):
    x.acquire()
    print('hello world', y)
    x.release()

if __name__ == '__main__':
    lock = Lock()
    for num in range(10):
        Process(target=my_function, args=(lock, num)).start()

위 예제에서는 락(Lock) 인스턴스를 사용하여 특정 시점에 단 하나의 프로세스만 표준 출력을 사용할 수 있도록 보장합니다. 이처럼 락은 여러 프로세스가 동시에 출력하면서 발생할 수 있는 데이터 충돌과 출력 꼬임 문제를 방지해 줍니다.

풀링(Pooling)

작업 풀링에는 Pool 클래스를 사용합니다. Pool은 제출된 전체 작업을 처리할 워커 프로세스들의 집합을 생성하며, 다음과 같은 형태로 선언합니다.

class multiprocessing.Pool([processes[, initializer[, initargs[, maxtasksperchild]]]])

Pool 객체는 어떤 작업을 워커에 제출할지 관리하며, 타임아웃(timeout), 콜백(callback), 그리고 병렬 map 구현을 지원하는 비동기 결과(asynchronous result) 기능도 함께 제공합니다.

참고로 processes 인자가 None으로 지정되면 cpu_count() 값이 기본값으로 사용되며, initializer가 None이 아니라면 initializer(*initargs)가 초기화 시점에 호출됩니다.

주요 메서드

apply(func[, args[, kwds]])

내장 함수 apply()와 동일하게 동작하며, 결과가 준비될 때까지 블록킹(blocking)됩니다. 병렬로 작업을 수행하려면 apply_async() 메서드를 사용하는 것이 좋습니다.

apply_async(func[, args[, kwds[, callback]]])

결과 객체(result object)를 반환하는 비동기 버전입니다.

map(func, iterable[, chunksize])

내장 함수 map()과 유사하며 하나의 iterable 인수만 지원하고, 결과가 준비될 때까지 블록킹됩니다. 내부적으로는 iterable을 여러 개의 작은 청크(chunk)로 분할한 뒤, 각 청크를 개별 작업으로 프로세스 풀에 제출하여 병렬 처리합니다.

map_async(func, iterable[, chunksize[, callback]])

map()의 비동기 버전으로, 결과 객체를 반환합니다.

imap(func, iterable[, chunksize])

itertools.imap()과 동일하게 동작하며, 인수 크기는 map()에서 사용되는 것과 같습니다. 결과를 즉시 하나씩 순회할 수 있다는 장점이 있습니다.

imap_unordered(func, iterable[, chunksize])

imap()과 거의 동일하지만, 반환되는 이터레이터의 결과 순서는 입력 순서가 아닌 작업 완료 순서에 따라 결정됩니다.

close()

더 이상 새로운 작업을 받지 않도록 풀을 닫습니다. 이미 할당된 작업을 모두 마친 후 워커 프로세스가 종료됩니다.

terminate()

진행 중인 작업의 완료 여부와 관계없이 워커 프로세스를 즉시 중지시켜야 할 때 사용합니다.

join()

워커 프로세스의 종료를 기다립니다. 반드시 close() 또는 terminate()를 먼저 호출한 후에 사용해야 합니다.

AsyncResult 클래스

multiprocessing.pool.AsyncResultPool.apply_async()Pool.map_async()가 반환하는 객체로, 비동기 작업의 결과를 다루는 다양한 메서드를 제공합니다.

get([timeout])
결과가 도착하면 해당 값을 반환합니다. timeout이 지나면 TimeoutError가 발생합니다.

wait([timeout])
결과가 준비되거나 timeout 초가 경과할 때까지 대기합니다.

ready()
호출(작업)이 완료되었는지 여부를 불리언 값으로 반환합니다.

successful()
오류 없이 호출이 정상적으로 완료되었는지 여부를 반환합니다.

풀링 예제 코드

# -*- coding: utf-8 -*-
"""
Created on Sun Sep 30 12:17:58 2018
@author: Tutorials Point
"""
from multiprocessing import Pool
import time

def myfunction(m):
    return m * m

if __name__ == '__main__':
    my_pool = Pool(processes=4)  # 4개의 워커 프로세스 시작

    result = my_pool.apply_async(myfunction, (10,))  # f(10)을 단일 프로세스에서 비동기 평가
    print(result.get(timeout=1))

    print(my_pool.map(myfunction, range(10)))  # "[0, 1, 4, ..., 81]" 출력

    my_it = my_pool.imap(myfunction, range(10))
    print(my_it.next())              # "0" 출력
    print(my_it.next())              # "1" 출력
    print(my_it.next(timeout=1))     # 컴퓨터가 매우 느리지 않다면 "4" 출력

    result = my_pool.apply_async(time.sleep, (10,))
    print(result.get(timeout=1))     # multiprocessing.TimeoutError 발생

위 예제를 통해 Pool 객체 생성, 비동기 실행(apply_async), 병렬 매핑(map), 순차적 결과 순회(imap), 그리고 타임아웃 처리까지 멀티프로세싱 풀링의 핵심 흐름을 확인할 수 있습니다. CPU 집약적인 작업을 다룰 때 이러한 기법을 활용하면 처리 성능을 크게 향상시킬 수 있습니다.