서론: 파이썬 멀티프로세싱의 이해
파이썬에서 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: 프로세스 IDp.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('모든 작업 완료')
JoinableQueue는 task_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 코어 수를 넘지 않도록 관리
- 가능하면 공유 데이터보다 메시지 전달 방식 사용
- 대규모 작업에는 프로세스 풀 활용