Python 멀티프로세싱 프로그래밍 완벽 가이드

서론: 파이썬 멀티프로세싱의 이해

파이썬에서 GIL(Global Interpreter Lock)때문에 스레드를 사용해도 실제 다중 코어를 활용하지 못합니다. 따라서 CPU 자원을 최대한 활용하려면 멀티프로세싱이 필수적입니다. multiprocessing 모듈은 서브프로세스를 생성하고 관리하는 강력한 기능을 제공합니다.

1. multiprocessing 모듈概述

이 모듈은 다음과 같은 기능을 지원합니다:

  • 서브프로세스 생성 및 관리
  • 프로세스 간 통신(IPC)
  • 데이터 공유 메커니즘
  • 동기화 프리미티브(Lock, Semaphore, Event 등)

중요: 스레드와 달리 프로세스는 메모리를 공유하지 않습니다. 한 프로세스에서 데이터를 변경해도 다른 프로세스에는 영향을 주지 않습니다.

2. Process 클래스 상세

클래스 생성자

Process([group, target, name, args, kwargs])

주요 파라미터

  • target: 실행할 함수
  • args: 타겟 함수에 전달할 위치 인자(튜플)
  • kwargs: 타겟 함수에 전달할 키워드 인자(딕셔너리)
  • name: 프로세스 이름

핵심 메서드

  • p.start(): 프로세스 시작
  • p.run(): 프로세스 실행 시 호출되는 메서드
  • p.terminate(): 프로세스 강제 종료
  • p.is_alive(): 프로세스 실행 상태 확인
  • p.join([timeout]): 프로세스 대기

주요 속성

  • p.daemon: 데몬 프로세스 설정
  • p.pid: 프로세스 ID
  • p.name: 프로세스 이름

3. 프로세스 생성 방법

Windows에서는 반드시 if __name__ == '__main__': 블록 내에 코드를 작성해야 합니다.

방법 1: 함수 기반 생성

import time
import random
from multiprocessing import Process

def work(name):
    print(f'{name} 작업 시작')
    time.sleep(random.randint(1, 5))
    print(f'{name} 작업 완료')

if __name__ == '__main__':
    tasks = ['task1', 'task2', 'task3', 'task4']
    process_list = []
    
    for task in tasks:
        p = Process(target=work, args=(task,))
        p.start()
        process_list.append(p)
    
    for p in process_list:
        p.join()
    
    print('모든 프로세스 종료')

방법 2: 클래스 기반 생성

import time
import random
from multiprocessing import Process

class Worker(Process):
    def __init__(self, name):
        super().__init__()
        self.name = name
    
    def run(self):
        print(f'{self.name} 실행 중')
        time.sleep(random.randint(1, 3))
        print(f'{self.name} 종료')

if __name__ == '__main__':
    workers = [Worker(f'worker_{i}') for i in range(4)]
    
    for w in workers:
        w.start()
    
    print('메인 프로세스继续')

메모리 격리 확인

from multiprocessing import Process

counter = 100

def modify_value():
    global counter
    counter = 0
    print(f'서브프로세스: {counter}')

if __name__ == '__main__':
    p = Process(target=modify_value)
    p.start()
    p.join()
    print(f'메인 프로세스: {counter}')

각 프로세스는 독립된 메모리 공간을 가지므로 변수가 서로 영향을 주지 않습니다.

4. 소켓 서버 멀티프로세싱实例

서버 코드

from socket import *
from multiprocessing import Process

server_socket = socket(AF_INET, SOCK_STREAM)
server_socket.setsockopt(SOL_SOCKET, SO_REUSEADDR, 1)
server_socket.bind(('127.0.0.1', 9999))
server_socket.listen(5)

def handle_client(conn, addr):
    print(f'클라이언트 연결: {addr}')
    while True:
        try:
            data = conn.recv(1024)
            if not data:
                break
            conn.send(data.upper())
        except Exception:
            break
    conn.close()

if __name__ == '__main__':
    while True:
        conn, addr = server_socket.accept()
        p = Process(target=handle_client, args=(conn, addr))
        p.start()

이 방식의 한계: 클라이언트 수만큼 프로세스가 생성되어 시스템 부하가 발생합니다. 이를 해결하려면 프로세스プーリング이 필요합니다.

5. join 메서드 활용

from multiprocessing import Process
import time
import random

class MyTask(Process):
    def __init__(self, name):
        self.name = name
        super().__init__()
    
    def run(self):
        print(f'{self.name} 시작')
        time.sleep(random.randint(1, 3))
        print(f'{self.name} 완료')

if __name__ == '__main__':
    tasks = [MyTask(f'task_{i}') for i in range(4)]
    
    for t in tasks:
        t.start()
    
    # 메인 프로세스가 각 프로세스 완료를 대기
    for t in tasks:
        t.join()
    
    print('모든 작업 종료')

join()은 메인 스레드가 대기하는 것이지 프로세스가 순차 실행되는 것이 아닙니다. 모든 프로세스는 시작과 동시에 병렬로 실행됩니다.

6. 프로세스 종료 및 상태 확인

from multiprocessing import Process
import time
import random

class Worker(Process):
    def __init__(self, name):
        super().__init__()
        self.name = name
    
    def run(self):
        print(f'{self.name} 실행')
        time.sleep(random.randint(1, 5))

if __name__ == '__main__':
    w = Worker('my_worker')
    w.start()
    
    print(f'실행 상태: {w.is_alive()}')
    w.terminate()
    print(f'종료 후 상태: {w.is_alive()}')
    print(f'PID: {w.pid}')

7. 좀비 프로세스와 고아 프로세스

좀비 프로세스

자식 프로세스가 종료しても親프로세스가 상태 정보를 회수하지 않으면 좀비 프로세스가 됩니다. 부모가 wait() 또는 join()을 호출하지 않으면 발생합니다.

고아 프로세스

부모 프로세스가 자식보다 먼저 종료되면 자식 프로세스는 init 프로세스(PID 1)에 의해收养됩니다. 고아 프로세스는危害가 없습니다.

좀비 프로세스 방지

from multiprocessing import Process
import time
import os

def child_task():
    print(f'자식进程的 PID: {os.getpid()}')

if __name__ == '__main__':
    p = Process(target=child_task)
    p.start()
    
    print(f'부모 PID: {os.getpid()}')
    time.sleep(1)  # 자식 프로세스 종료 대기
    
    # join으로 자원 회수
    p.join()
    print('좀비 프로세스 방지 완료')

8. 데몬 프로세스

데몬 프로세스는 메인 프로세스 종료 시 자동으로 종료됩니다.

from multiprocessing import Process
import time
import random

class Service(Process):
    def __init__(self, name):
        super().__init__()
        self.name = name
    
    def run(self):
        print(f'{self.name} 서비스 시작')
        time.sleep(random.randint(1, 3))
        print(f'{self.name} 서비스 종료')

if __name__ == '__main__':
    s = Service('background')
    s.daemon = True  # 반드시 start() 전에 설정
    s.start()
    
    print('메인 프로세스 종료')

주의사항

  • 데몬 프로세스는 자식 프로세스를 생성할 수 없습니다
  • 메인 프로세스 종료 시 데몬 프로세스도 즉시 종료됩니다

9. 프로세스 동기화 - Lock

여러 프로세스가同一 리소스에 접근하면 경합이 발생합니다. 이를 방지하려면 Lock이 필요합니다.

터미널 출력 경합 해결

from multiprocessing import Process, Lock
import os
import time

def print_task(task_id, lock):
    lock.acquire()
    try:
        print(f'프로세스 {os.getpid()}가 {task_id} 처리')
        time.sleep(2)
        print(f'프로세스 {task_id} 완료')
    finally:
        lock.release()

if __name__ == '__main__':
    lock = Lock()
    processes = []
    
    for i in range(5):
        p = Process(target=print_task, args=(i, lock))
        p.start()
        processes.append(p)
    
    for p in processes:
        p.join()

파일 접근 동시성 해결

from multiprocessing import Process, Lock
import time
import json

# 초기 데이터
initial_data = {"count": 10}
with open('inventory.txt', 'w') as f:
    json.dump(initial_data, f)

def read_stock():
    with open('inventory.txt', 'r') as f:
        return json.load(f)

def purchase(lock):
    with lock:
        stock = read_stock()
        time.sleep(0.1)  # 네트워크 지연 시뮬레이션
        
        if stock['count'] > 0:
            stock['count'] -= 1
            time.sleep(0.2)  # 쓰기 지연 시뮬레이션
            with open('inventory.txt', 'w') as f:
                json.dump(stock, f)
            print('구매 성공!')
        else:
            print('재고 없음')

if __name__ == '__main__':
    lock = Lock()
    processes = [Process(target=purchase, args=(lock,)) for _ in range(20)]
    
    for p in processes:
        p.start()
    
    for p in processes:
        p.join()

결론: Lock을 사용하면 데이터 일관성을 보장할 수 있지만 성능이 저하됩니다. 따라서 가능하면 공유 데이터 대신 메시지 전달 방식(Queue, Pipe)을 사용하는 것이 좋습니다.

10. 큐 (Queue) - 권장 방법

프로세스 간 통신에는 Lock보다 큐를 사용하는 것이 좋습니다. 큐는 내부적으로 Pipe와 Lock을 조합하여 구현됩니다.

기본 사용법

from multiprocessing import Queue
import time

q = Queue(3)  # 최대 크기 3

q.put('item1')
q.put('item2')
q.put('item3')

print(f'포화 상태: {q.full()}')
print(f'항목 수: {q.qsize()}')

print(q.get())
print(q.get())
print(q.get())

print(f'빈 상태: {q.empty()}')

生产者-소비자 모델

from multiprocessing import Process, Queue
import time
import random
import os

def producer(q, name):
    for i in range(5):
        time.sleep(random.uniform(0.5, 1.5))
        item = f'{name}_item_{i}'
        q.put(item)
        print(f'[{os.getpid()}] 생산: {item}')
    q.put(None)  # 종료 신호

def consumer(q):
    while True:
        item = q.get()
        if item is None:
            break
        time.sleep(random.uniform(1, 2))
        print(f'[{os.getpid()}] 소비: {item}')

if __name__ == '__main__':
    q = Queue()
    
    p1 = Process(target=producer, args=(q, '공장A'))
    p2 = Process(target=producer, args=(q, '공장B'))
    c1 = Process(target=consumer, args=(q,))
    c2 = Process(target=consumer, args=(q,))
    
    p1.start()
    p2.start()
    c1.start()
    c2.start()
    
    p1.join()
    p2.join()
    
    # 종료 신호 전송 (소비자 수만큼)
    q.put(None)
    q.put(None)
    
    print('메인 프로세스 종료')

JoinableQueue 사용

from multiprocessing import Process, JoinableQueue
import time
import random
import os

def producer(q, name):
    for i in range(3):
        time.sleep(random.uniform(0.5, 1.5))
        item = f'{name}_{i}'
        q.put(item)
        print(f'[{os.getpid()}] 생산: {item}')
    q.join()  # 모든 항목 처리 대기

def consumer(q):
    while True:
        item = q.get()
        time.sleep(random.uniform(1, 2))
        print(f'[{os.getpid()}] 소비: {item}')
        q.task_done()  # 처리 완료 신호

if __name__ == '__main__':
    q = JoinableQueue()
    
    producers = [Process(target=producer, args=(q, f'공장{i}')) for i in range(3)]
    consumers = [Process(target=consumer, args=(q,)) for _ in range(2)]
    
    # 소비자를 데몬으로 설정
    for c in consumers:
        c.daemon = True
    
    for p in producers:
        p.start()
    
    for c in consumers:
        c.start()
    
    for p in producers:
        p.join()
    
    print('모든 작업 완료')

JoinableQueuetask_done()join() 메서드를 제공하여 더 elegant한生产者-소비자 모델을 구현할 수 있습니다.

11. 파이프 (Pipe)

파이프는 두 프로세스 간 양방향 통신을 제공합니다.

from multiprocessing import Process, Pipe
import time

def server(pipe):
    send_pipe, recv_pipe = pipe
    recv_pipe.close()  # 사용하지 않는쪽关闭
    
    while True:
        try:
            message = recv_pipe.recv()
            if message == 'exit':
                break
            print(f'수신: {message}')
            send_pipe.send(f'응답: {message}')
        except EOFError:
            break
    
    send_pipe.close()
    print('서버 종료')

if __name__ == '__main__':
    server_conn, client_conn = Pipe()
    
    p = Process(target=server, args=((server_conn, client_conn),))
    p.start()
    
    server_conn.close()
    
    for msg in ['안녕하세요', '점심먹었어요', 'exit']:
        client_conn.send(msg)
        if msg != 'exit':
            print(f'클라이언트 수신: {client_conn.recv()}')
    
    client_conn.close()
    p.join()
    print('메인 프로세스 종료')

주의: 파이프의 한쪽 끝을 사용하지 않으면 반드시 닫아야 합니다. 그렇지 않으면 recv()에서 무한 대기할 수 있습니다.

12. 공유 데이터 (Manager)

프로세스 간 데이터를 공유하려면 Manager를 사용합니다.

from multiprocessing import Manager, Process, Lock

def update_counter(shared_dict, lock):
    with lock:
        shared_dict['counter'] -= 1

if __name__ == '__main__':
    manager = Manager()
    lock = Lock()
    shared_data = manager.dict({'counter': 100})
    
    processes = [Process(target=update_counter, args=(shared_data, lock)) for _ in range(50)]
    
    for p in processes:
        p.start()
    
    for p in processes:
        p.join()
    
    print(f'최종 값: {shared_data["counter"]}')

권장 사항: 가능한 한 공유 데이터 방식보다는 큐나 파이프를 사용하는 것이 좋습니다. Lock 관리가 복잡해지고 성능 저하가 발생할 수 있습니다.

13. 세마포어 (Semaphore)

세마포어는 동시에 자원을 사용할 수 있는 프로세스 수를 제한합니다.

from multiprocessing import Process, Semaphore
import time
import random

def use_resource(sem, user_id):
    sem.acquire()
    print(f'[{user_id}] 자원 사용 중')
    time.sleep(random.uniform(0.5, 2))
    print(f'[{user_id}] 자원 반납')
    sem.release()

if __name__ == '__main__':
    sem = Semaphore(3)  # 최대 3개 프로세스同期使用
    
    processes = [Process(target=use_resource, args=(sem, i)) for i in range(10)]
    
    for p in processes:
        p.start()
    
    for p in processes:
        p.join()
    
    print('모든 프로세스 완료')

14. 이벤트 (Event)

이벤트는 프로세스 간 신호 전달을 위해 사용됩니다.

from multiprocessing import Process, Event
import time
import random

def traffic_car(event, car_id):
    while not event.is_set():
        print(f'[차 {car_id}] 신호 대기 중')
        event.wait(timeout=1)
    
    print(f'[차 {car_id}] 출발!')
    time.sleep(random.uniform(1, 2))
    print(f'[차 {car_id}] 도착')

def traffic_light(event, interval):
    while True:
        time.sleep(interval)
        if event.is_set():
            event.clear()
            print('신호: 적색')
        else:
            event.set()
            print('신호: 녹색')

if __name__ == '__main__':
    event = Event()
    event.set()  # 초기값: 녹색
    
    # 신호등 프로세스
    light = Process(target=traffic_light, args=(event, 3))
    light.daemon = True
    light.start()
    
    # 차량 프로세스
    cars = [Process(target=traffic_car, args=(event, i)) for i in range(5)]
    
    for c in cars:
        c.start()
    
    for c in cars:
        c.join()
    
    print('모든 차량 도착')

15. 프로세스 풀 (Process Pool)

프로세스 풀은 생성할 프로세스 수를 제한하고 재사용합니다. 이는 시스템 리소스를 효율적으로 관리합니다.

Pool 클래스

Pool([processes, initializer, initargs])
  • processes: 작업자 프로세스 수 (기본값: CPU 코어 수)
  • initializer: 각 프로세스 시작 시 실행할 함수

주요 메서드

  • p.apply(func, args): 동기 실행
  • p.apply_async(func, args, callback): 비동기 실행
  • p.map(func, iterable): 맵 적용
  • p.close(): 새 작업 제한
  • p.join(): 모든 작업 완료 대기

동기 실행 (apply)

from multiprocessing import Pool
import os
import time

def compute_square(n):
    print(f'[{os.getpid()}] 제곱 계산: {n}')
    time.sleep(1)
    return n ** 2

if __name__ == '__main__':
    with Pool(3) as pool:
        results = []
        for i in range(10):
            result = pool.apply(compute_square, args=(i,))
            results.append(result)
        
        print(f'결과: {results}')

비동기 실행 (apply_async)

from multiprocessing import Pool
import os
import time

def calculate(n):
    print(f'[{os.getpid()}] 계산 중: {n}')
    time.sleep(2)
    return n * 2

if __name__ == '__main__':
    with Pool(4) as pool:
        results = [pool.apply_async(calculate, args=(i,)) for i in range(8)]
        
        pool.close()
        pool.join()
        
        for r in results:
            print(r.get())  # 결과 수집

콜백 함수 사용

from multiprocessing import Pool
import requests
import os

def fetch_url(url):
    print(f'[{os.getpid()}] 가져오는 중: {url}')
    response = requests.get(url)
    return {'url': url, 'status': response.status_code}

def save_result(result):
    print(f'[{os.getpid()}] 저장: {result}')
    with open('results.txt', 'a') as f:
        f.write(f"{result['url']}: {result['status']}\n")

if __name__ == '__main__':
    urls = [
        'https://python.org',
        'https://github.com',
        'https://stackoverflow.com'
    ]
    
    with Pool(2) as pool:
        for url in urls:
            pool.apply_async(fetch_url, args=(url,), callback=save_result)
        
        pool.close()
        pool.join()
    
    print('모든 작업 완료')

Socket 서버에 프로세스 풀 적용

from socket import *
from multiprocessing import Pool
import os

server = socket(AF_INET, SOCK_STREAM)
server.setsockopt(SOL_SOCKET, SO_REUSEADDR, 1)
server.bind(('127.0.0.1', 8888))
server.listen(5)

def handle_request(conn_info):
    conn, addr = conn_info
    print(f'[{os.getpid()}] 클라이언트 연결: {addr}')
    
    while True:
        try:
            data = conn.recv(1024)
            if not data:
                break
            conn.send(data.upper())
        except Exception:
            break
    
    conn.close()

if __name__ == '__main__':
    pool = Pool(3)  # 3개 프로세스로 제한
    
    while True:
        conn, addr = server.accept()
        pool.apply_async(handle_request, args=((conn, addr),))

이 구현에서는 최대 3개의 프로세스가 동시에 클라이언트를 처리하므로 시스템 리소스를 효율적으로 사용할 수 있습니다.

map 메서드 사용

from multiprocessing import Pool
import time

def heavy_computation(n):
    time.sleep(1)
    return n ** 2

if __name__ == '__main__':
    with Pool(4) as pool:
        results = pool.map(heavy_computation, range(10))
        print(results)

마무리

멀티프로세싱은 CPU 집약적인 작업을 병렬화하는 데 필수적입니다. 주요 고려사항:

  • CPU 바운드 작업에는 멀티프로세싱 사용
  • I/O 바운드 작업에는 비동기 atau 스레딩 고려
  • 프로세스 수가 CPU 코어 수를 넘지 않도록 관리
  • 가능하면 공유 데이터보다 메시지 전달 방식 사용
  • 대규모 작업에는 프로세스 풀 활용

태그: python multiprocessing Process concurrency parallel-programming

7월 23일 06:34에 게시됨