์นดํ…Œ๊ณ ๋ฆฌ ์—†์Œ

Python asyncio์™€ multiprocessing ์กฐํ•ฉ ์‹œ ๋ฐœ์ƒํ•˜๋Š” IPC ์ง๋ ฌํ™” ๋ณ‘๋ชฉ ๋ฐ ์ด๋ฒคํŠธ ๋ฃจํ”„ ์ง€์—ฐ ํ•ด๊ฒฐ ๊ฐ€์ด๋“œ

๊ฒŒ์ž„๊ต์ˆ˜ 2026. 8. 30. 22:01
๋ฐ˜์‘ํ˜•

๐ŸŒ English Abstract

This technical deep-dive addresses severe event loop latency and starvation issues in Python real-time telemetry ingestion pipelines combining asyncio and multiprocessing. It examines how IPC serialization (pickle) synchronously blocks the asynchronous event loop thread during high-throughput metric handling. Furthermore, it presents concrete troubleshooting steps, architectural optimizations using shared memory, zero-copy techniques, fast binary serialization formats, and thread-executor offloading to achieve sub-millisecond telemetry ingestion latency.

๋Œ€๊ทœ๋ชจ ์ˆ˜์ง‘ ํŒŒ์ดํ”„๋ผ์ธ(Ingestion Pipeline)์„ ๊ตฌ์ถ•ํ•  ๋•Œ ํŒŒ์ด์ฌ(Python)์€ ๋›ฐ์–ด๋‚œ ์ƒํƒœ๊ณ„์™€ ์ƒ์‚ฐ์„ฑ ๋•๋ถ„์— ์ž์ฃผ ์„ ํƒ๋˜๋Š” ์–ธ์–ด์ž…๋‹ˆ๋‹ค. ํŠนํžˆ ๋Œ€๋Ÿ‰์˜ ๋„คํŠธ์›Œํฌ I/O ๋™์‹œ์„ฑ์„ ์ฒ˜๋ฆฌํ•˜๊ธฐ ์œ„ํ•ด asyncio๋ฅผ ํ™œ์šฉํ•˜๊ณ , CPython์˜ GIL(Global Interpreter Lock) ํ•œ๊ณ„๋ฅผ ๊ทน๋ณตํ•˜์—ฌ CPU ์—ฐ์‚ฐ(๋ฐ์ดํ„ฐ ๊ฒ€์ฆ, ํŒฉ ํŒŒ์‹ฑ ๋“ฑ)์„ ๋ถ„์‚ฐํ•˜๊ธฐ ์œ„ํ•ด multiprocessing ๋ชจ๋“ˆ์„ ํ˜ผ์šฉํ•˜๋Š” ์•„ํ‚คํ…์ฒ˜๊ฐ€ ์ž์ฃผ ์‚ฌ์šฉ๋ฉ๋‹ˆ๋‹ค.

๊ทธ๋Ÿฌ๋‚˜ ์ดˆ๋‹น ์ˆ˜๋งŒ ๊ฑด ์ด์ƒ์˜ ํ…”๋ ˆ๋ฉ”ํŠธ๋ฆฌ(Telemetry) ๋ฐ์ดํ„ฐ๊ฐ€ ์œ ์ž…๋˜๋Š” ์‹ค์‹œ๊ฐ„ ํŒŒ์ดํ”„๋ผ์ธ์„ ์šด์˜ํ•˜๋‹ค ๋ณด๋ฉด, ๋„คํŠธ์›Œํฌ ํƒ€์ž„์•„์›ƒ์ด ๋ฐœ์ƒํ•˜๊ฑฐ๋‚˜ asyncio ์ด๋ฒคํŠธ ๋ฃจํ”„(Event Loop)์˜ ์ง€์—ฐ ์‹œ๊ฐ„(Latency)์ด ๊ธฐํ•˜๊ธ‰์ˆ˜์ ์œผ๋กœ ์ฆ๊ฐ€ํ•˜๋Š” ํ˜„์ƒ์„ ๋ชฉ๊ฒฉํ•˜๊ฒŒ ๋ฉ๋‹ˆ๋‹ค. ์ด ๋ฌธ์ œ์˜ ๊ทผ๋ณธ ์›์ธ์€ ํ”„๋กœ์„ธ์Šค ๊ฐ„ ํ†ต์‹ (IPC, Inter-Process Communication) ๊ณผ์ •์—์„œ ๋ฐœ์ƒํ•˜๋Š” ๋ฐ์ดํ„ฐ ์ง๋ ฌํ™”(Serialization) ๋ณ‘๋ชฉ์— ์žˆ์Šต๋‹ˆ๋‹ค. ๋ณธ ๊ธ€์—์„œ๋Š” ์ด ํŠธ๋Ÿฌ๋ธ”์ŠˆํŒ… ๊ณผ์ •๊ณผ ์ด๋ฅผ ๊ทน๋ณตํ•˜๊ธฐ ์œ„ํ•œ ์•„ํ‚คํ…์ฒ˜ ์ตœ์ ํ™” ๋ฐฉ์•ˆ์„ ์ƒ์„ธํžˆ ๋‹ค๋ฃน๋‹ˆ๋‹ค.

 



 

1. ๋ฌธ์ œ์˜ ๋ณธ์งˆ: asyncio ์ด๋ฒคํŠธ ๋ฃจํ”„์™€ multiprocessing IPC์˜ ์ถฉ๋Œ

asyncio๋Š” ๋‹จ์ผ ์Šค๋ ˆ๋“œ ๊ธฐ๋ฐ˜์˜ ๋น„๋™๊ธฐ I/O ์ด๋ฒคํŠธ ๋ฃจํ”„(Event Loop)๋ฅผ ๊ธฐ๋ฐ˜์œผ๋กœ ์ž‘๋™ํ•ฉ๋‹ˆ๋‹ค. ํƒœ์Šคํฌ(Task)๊ฐ€ I/O ์ž‘์—…์„ ๊ธฐ๋‹ค๋ฆฌ๋Š” ๋™์•ˆ await ํ‚ค์›Œ๋“œ๋ฅผ ํ†ตํ•ด ์ œ์–ด๊ถŒ์„ ์ด๋ฒคํŠธ ๋ฃจํ”„์— ์–‘๋„(Yield)ํ•จ์œผ๋กœ์จ ํš๊ธฐ์ ์ธ ๋™์‹œ์„ฑ์„ ํ™•๋ณดํ•ฉ๋‹ˆ๋‹ค. ๋ฐ˜๋ฉด ํŒŒ์ด์ฌ์˜ multiprocessing ๋ชจ๋“ˆ์€ ํ”„๋กœ์„ธ์Šค ๊ฐ„์— ๋ฐ์ดํ„ฐ๋ฅผ ์ฃผ๊ณ ๋ฐ›๊ธฐ ์œ„ํ•ด ๋‚ด๋ถ€์ ์œผ๋กœ pickle ๋ผ์ด๋ธŒ๋Ÿฌ๋ฆฌ๋ฅผ ์‚ฌ์šฉํ•ด ๊ฐ์ฒด๋ฅผ ๋ฐ”์ด๋„ˆ๋ฆฌ๋กœ ์ง๋ ฌํ™”ํ•˜๊ณ , Pipe๋‚˜ Queue๋ฅผ ํ†ตํ•ด ์ „์†กํ•ฉ๋‹ˆ๋‹ค.

์ด๋ฒคํŠธ ๋ฃจํ”„ ์ฐจ๋‹จ(Event Loop Blocking)์˜ ๋ฉ”์ปค๋‹ˆ์ฆ˜

์‹ค์‹œ๊ฐ„ ํ…”๋ ˆ๋ฉ”ํŠธ๋ฆฌ ์ˆ˜์ง‘๊ธฐ์—์„œ ๋ฉ”์ธ ํ”„๋กœ์„ธ์Šค์˜ asyncio ๋ฃจํ”„๊ฐ€ ๋„คํŠธ์›Œํฌ ์†Œ์ผ“์œผ๋กœ๋ถ€ํ„ฐ ๋ณต์žกํ•œ JSON ๋˜๋Š” ๋”•์…”๋„ˆ๋ฆฌ ํ˜•ํƒœ์˜ ๋ฐ์ดํ„ฐ๋ฅผ ์ˆ˜์‹ ํ•œ ํ›„, ์ด๋ฅผ ์ž‘์—…์ž ํ”„๋กœ์„ธ์Šค(Worker Process)์— ์ „๋‹ฌํ•˜๊ธฐ ์œ„ํ•ด multiprocessing.Queue.put()์„ ํ˜ธ์ถœํ•  ๋•Œ ๋ฌธ์ œ๊ฐ€ ๋ฐœ์ƒํ•ฉ๋‹ˆ๋‹ค.

  • ๋™๊ธฐ์‹ ์ง๋ ฌํ™” ์—ฐ์‚ฐ: multiprocessing.Queue.put() ๋‚ด๋ถ€์—์„œ ์‹คํ–‰๋˜๋Š” pickle.dumps() ๊ณผ์ •์€ ์™„์ „ํ•œ CPU ๋ฐ”์šด๋“œ(CPU-bound) ์ž‘์—…์ด๋ฉฐ ๋™๊ธฐ(Blocking) ๋ฐฉ์‹์œผ๋กœ ๋™์ž‘ํ•ฉ๋‹ˆ๋‹ค.
  • ์ด๋ฒคํŠธ ๋ฃจํ”„ ๋ฉˆ์ถค(Starvation): ์ง๋ ฌํ™”ํ•  ๋ฐ์ดํ„ฐ์˜ ํฌ๊ธฐ๊ฐ€ ํฌ๊ฑฐ๋‚˜ ๊ฐœ์ˆ˜๊ฐ€ ๋งŽ์•„์ง€๋ฉด pickle.dumps() ์—ฐ์‚ฐ์ด ์™„๋ฃŒ๋  ๋•Œ๊นŒ์ง€ ์ด๋ฒคํŠธ ๋ฃจํ”„ ์Šค๋ ˆ๋“œ ์ „์ฒด๊ฐ€ ๋ธ”๋กœํ‚น๋ฉ๋‹ˆ๋‹ค.
  • ๋„คํŠธ์›Œํฌ ์ง€์—ฐ ํŒŒ๊ธ‰: ์ด๋ฒคํŠธ ๋ฃจํ”„๊ฐ€ ๋ฉˆ์ถ”๋ฉด ์ƒˆ๋กœ์šด ์†Œ์ผ“ ์ ‘์† ์ฒ˜๋ฆฌ, ํ•‘/ํ(Ping/Pong) ํ—ฌ์Šค์ฒดํฌ, ์ฝ”๋ฃจํ‹ด ์Šค์ผ€์ค„๋ง์ด ๋ชจ๋‘ ์ค‘๋‹จ๋˜์–ด ํด๋ผ์ด์–ธํŠธ ์ธก์—์„œ ์ปค๋„ฅ์…˜ ํƒ€์ž„์•„์›ƒ(Connection Timeout)์ด ๋ฐœ์ƒํ•ฉ๋‹ˆ๋‹ค.

 

2. ํ˜„์ƒ ์ง„๋‹จ ๋ฐ ํŠธ๋Ÿฌ๋ธ”์ŠˆํŒ… ๋‹จ๊ณ„๋ณ„ ์ ‘๊ทผ

์‹œ์Šคํ…œ ์„ฑ๋Šฅ์ด ์ €ํ•˜๋˜์—ˆ์„ ๋•Œ ๋ณ‘๋ชฉ์˜ ์›์ธ์ด ๋„คํŠธ์›Œํฌ์ธ์ง€, CPU ์—ฐ์‚ฐ์ธ์ง€, ์•„๋‹ˆ๋ฉด IPC ์ง๋ ฌํ™”์ธ์ง€๋ฅผ ์ •ํ™•ํžˆ ํŒ๋ณ„ํ•˜๋Š” ๊ตฌ์ฒด์ ์ธ ์ง„๋‹จ ์ ˆ์ฐจ์ž…๋‹ˆ๋‹ค.

๋‹จ๊ณ„ 1: asyncio ์ด๋ฒคํŠธ ๋ฃจํ”„ ์ง€์—ฐ ๋ชจ๋‹ˆํ„ฐ๋ง

๊ฐ€์žฅ ๋จผ์ € ์ด๋ฒคํŠธ ๋ฃจํ”„๊ฐ€ ํŠน์ • ํƒœ์Šคํฌ์— ์˜ํ•ด ์–ผ๋งˆ๋‚˜ ์˜ค๋žซ๋™์•ˆ ๋ฉˆ์ถฐ ์žˆ๋Š”์ง€๋ฅผ ์ธก์ •ํ•ด์•ผ ํ•ฉ๋‹ˆ๋‹ค. ํŒŒ์ด์ฌ ๊ธฐ๋ณธ ์ œ๊ณต ๋””๋ฒ„๊ทธ ๋ชจ๋“œ๋‚˜ ์ปค์Šคํ…€ ๋ฃจํ”„ ๋ชจ๋‹ˆํ„ฐ๋ฅผ ํ™œ์šฉํ•  ์ˆ˜ ์žˆ์Šต๋‹ˆ๋‹ค.

import asyncio
import time
import logging

# asyncio ๋””๋ฒ„๊ทธ ๋ชจ๋“œ ํ™œ์„ฑํ™” (์‹คํ–‰์— 100ms ์ด์ƒ ๊ฑธ๋ฆฌ๋Š” ๋ธ”๋กœํ‚น ํƒœ์Šคํฌ ๋กœ๊น…)
loop = asyncio.get_event_loop()
loop.set_debug(True)
loop.slow_callback_duration = 0.1  # 100ms ๊ธฐ์ค€

async def monitor_event_loop_drift():
    """์ด๋ฒคํŠธ ๋ฃจํ”„์˜ ์‹ค์ œ ์ง€์—ฐ ์‹œ๊ฐ„์„ ์ง€์†์ ์œผ๋กœ ์ธก์ •ํ•˜๋Š” ์ฝ”๋ฃจํ‹ด"""
    while True:
        before = time.monotonic()
        await asyncio.sleep(1)
        after = time.monotonic()
        drift = (after - before) - 1.0
        if drift > 0.05:  # 50ms ์ด์ƒ์˜ ์ง€์—ฐ์ด ๋ฐœ์ƒํ•œ ๊ฒฝ์šฐ ๊ฒฝ๊ณ 
            logging.warning(f"[Event Loop Lag] Drift: {drift * 1000:.2f}ms")

๋‹จ๊ณ„ 2: IPC ์ง๋ ฌํ™” ํ”„๋กœํŒŒ์ผ๋ง (Pickle Overhead Quantification)

cProfile ๋˜๋Š” yappi ํ”„๋กœํŒŒ์ผ๋Ÿฌ๋ฅผ ์‚ฌ์šฉํ•˜์—ฌ ์‹คํ–‰ ์‹œ๊ฐ„์„ ๋ถ„์„ํ•˜๋ฉด, multiprocessing/queues.py ๋‚ด๋ถ€์˜ _feed ๋ฉ”์„œ๋“œ์™€ pickle.dumps์—์„œ ์ „์ฒด CPU ์‹œ๊ฐ„์˜ 40~60% ์ด์ƒ์„ ์†Œ๋น„ํ•˜๋Š” ์—ฃ์ง€ ์ผ€์ด์Šค๋ฅผ ๋ฐœ๊ฒฌํ•  ์ˆ˜ ์žˆ์Šต๋‹ˆ๋‹ค.

์ˆ˜์ง‘ ๋ฐ์ดํ„ฐ ์œ ํ˜• ๊ฑด๋‹น ํ‰๊ท  ๋ฐ์ดํ„ฐ ํฌ๊ธฐ ์ดˆ๋‹น ์ฒ˜๋ฆฌ ๊ฑด์ˆ˜ (TPS) Pickle ์ง๋ ฌํ™” ์‹œ ์ด๋ฒคํŠธ ๋ฃจํ”„ ์ง€์—ฐ์‹œ๊ฐ„
๋‹จ์ˆœ ๋ฉ”ํŠธ๋ฆญ (Integer/Float) 1 KB ์ดํ•˜ 50,000 TPS 12ms ~ 25ms (์–‘ํ˜ธ)
์ค‘์ฒฉ ๊ตฌ์กฐ ํ…”๋ ˆ๋ฉ”ํŠธ๋ฆฌ (Dict/List) 50 KB 10,000 TPS 180ms ~ 450ms (์‹ฌ๊ฐํ•œ ๋ณ‘๋ชฉ)
๊ณ ๋ฐ€๋„ ์„ผ์„œ ๋กœ๊ทธ (String/Array) 500 KB 2,000 TPS 1,200ms ์ด์ƒ (์ด๋ฒคํŠธ ๋ฃจํ”„ ๋งˆ๋น„)

 



 

3. IPC ๋ณ‘๋ชฉ ํ•ด๊ฒฐ์„ ์œ„ํ•œ 4๊ฐ€์ง€ ์•„ํ‚คํ…์ฒ˜ ์ตœ์ ํ™” ์ „๋žต

์ง๋ ฌํ™” ๋ณ‘๋ชฉ์„ ๊ทน๋ณตํ•˜๊ณ  ์‹ค์‹œ๊ฐ„ ์ˆ˜์ง‘ ์‹œ์Šคํ…œ์˜ ์ฒ˜๋ฆฌ๋Ÿ‰(Throughput)์„ ๊ทน๋Œ€ํ™”ํ•˜๊ธฐ ์œ„ํ•œ ๋‹จ๊ณ„๋ณ„ ์ ์šฉ ๊ธฐ์ˆ ์ž…๋‹ˆ๋‹ค.

์ „๋žต 1: loop.run_in_executor๋ฅผ ํ†ตํ•œ ์ง๋ ฌํ™” ์ž‘์—… ๋ณ„๋„ ์Šค๋ ˆ๋“œ ์œ„์ž„

๊ฐ€์žฅ ์†์‰ฝ๊ฒŒ ์ ์šฉ ๊ฐ€๋Šฅํ•œ ์ž„์‹œ ๋ฐฉํŽธ์€ multiprocessing.Queue.put()๊ณผ ๊ฐ™์€ blocking ์ž‘์—…์„ ์ด๋ฒคํŠธ ๋ฃจํ”„์˜ ๋ฉ”์ธ ์Šค๋ ˆ๋“œ๊ฐ€ ์•„๋‹Œ ThreadPoolExecutor๋กœ ์˜คํ”„๋กœ๋”ฉ(Offloading)ํ•˜๋Š” ๊ฒƒ์ž…๋‹ˆ๋‹ค.

import asyncio
from concurrent.futures import ThreadPoolExecutor
import multiprocessing as mp

executor = ThreadPoolExecutor(max_workers=8)
mp_queue = mp.Queue()

async def send_telemetry_async(payload):
    loop = asyncio.get_running_loop()
    # Queue.put(payload) ์ž‘์—…์„ ThreadPool์—์„œ ์‹คํ–‰ํ•˜์—ฌ ์ด๋ฒคํŠธ ๋ฃจํ”„ ์ฐจ๋‹จ์„ ๋ฐฉ์ง€ํ•จ
    await loop.run_in_executor(executor, mp_queue.put, payload)

์ฃผ์˜์ : ์Šค๋ ˆ๋“œ ์ƒ์„ฑ ๋น„์šฉ๊ณผ ์ปจํ…์ŠคํŠธ ์Šค์œ„์นญ(Context Switching) ์˜ค๋ฒ„ํ—ค๋“œ๊ฐ€ ๋ฐœ์ƒํ•˜๋ฏ€๋กœ ๋ฐ์ดํ„ฐ๋Ÿ‰์ด ๋งค์šฐ ๋งŽ์„ ๋•Œ๋Š” ๊ทผ๋ณธ์ ์ธ ํ•ด๊ฒฐ์ฑ…์ด ๋˜์ง€ ๋ชปํ•ฉ๋‹ˆ๋‹ค.

์ „๋žต 2: ๊ณ ์„ฑ๋Šฅ ๋ฐ”์ด๋„ˆ๋ฆฌ ์ง๋ ฌํ™” ํฌ๋งท ๋„์ž… (MessagePack, Orjson, Protobuf)

ํŒŒ์ด์ฌ์˜ ํ‘œ์ค€ pickle์€ ์œ ์—ฐํ•˜์ง€๋งŒ ์†๋„๊ฐ€ ๋ฌด๊ฒ์Šต๋‹ˆ๋‹ค. C ์–ธ์–ด๋กœ ์ž‘์„ฑ๋œ ๋ฐ”์ธ๋”ฉ(C-extension)์„ ์‚ฌ์šฉํ•˜๋Š” Orjson์ด๋‚˜ MessagePack, Protocol Buffers(Protobuf)๋กœ ์ง๋ ฌํ™” ํฌ๋งท์„ ๊ต์ฒดํ•˜๋ฉด CPU ์ ์œ ์œจ์„ ํฌ๊ฒŒ ๋‚ฎ์ถœ ์ˆ˜ ์žˆ์Šต๋‹ˆ๋‹ค.

  • Orjson: ํŒŒ์ด์ฌ ์ตœ์†์˜ JSON ๋ผ์ด๋ธŒ๋Ÿฌ๋ฆฌ๋กœ, Rust ๊ธฐ๋ฐ˜์œผ๋กœ ์ž‘์„ฑ๋˜์–ด ๋น ๋ฅธ bytes ๋ณ€ํ™˜์ด ๊ฐ€๋Šฅํ•ฉ๋‹ˆ๋‹ค.
  • MessagePack: JSON ๊ตฌ์กฐ๋ฅผ ์œ ์ง€ํ•˜๋ฉด์„œ ๋ฐ”์ด๋„ˆ๋ฆฌ ํ˜•ํƒœ๋กœ ์ถ•์†Œํ•˜์—ฌ ์ง๋ ฌํ™” ์†๋„์™€ ์ „์†ก ํฌ๊ธฐ๋ฅผ ๋™์‹œ์— ์ค„์ž…๋‹ˆ๋‹ค.

์ „๋žต 3: ๊ณต์œ  ๋ฉ”๋ชจ๋ฆฌ(SharedMemory) ๊ธฐ๋ฐ˜ ์ œ๋กœ์นดํ”ผ(Zero-Copy) ์•„ํ‚คํ…์ฒ˜ ๋„์ž…

๋Œ€์šฉ๋Ÿ‰ ํ…”๋ ˆ๋ฉ”ํŠธ๋ฆฌ ์ˆ˜์ง‘์˜ ๊ฐ€์žฅ ์šฐ์ˆ˜ํ•œ ํ•ด๊ฒฐ์ฑ…์€ ๋ฐ์ดํ„ฐ๋ฅผ IPC ํŒŒ์ดํ”„๋‚˜ ํ๋กœ ์ง์ ‘ ๋ณต์‚ฌ ์ „์†กํ•˜์ง€ ์•Š๊ณ , Python 3.8+ ๊ณต์œ  ๋ฉ”๋ชจ๋ฆฌ(multiprocessing.shared_memory) ๊ณต๊ฐ„์— ์“ด ๋’ค ์ž‘์—…์ž ํ”„๋กœ์„ธ์Šค์— ๋ฉ”๋ชจ๋ฆฌ ์ฃผ์†Œ ๋ฐ ์ด๊ฒฉ๊ฑฐ๋ฆฌ(Offset/Size)๋งŒ ์ „๋‹ฌํ•˜๋Š” ์ œ๋กœ์นดํ”ผ ๊ธฐ๋ฒ•์ž…๋‹ˆ๋‹ค.

from multiprocessing import shared_memory
import numpy as np

# 1. ๋ฉ”์ธ ํ”„๋กœ์„ธ์Šค: ๊ณต์œ  ๋ฉ”๋ชจ๋ฆฌ ํ• ๋‹น ๋ฐ ๋ฐ์ดํ„ฐ ๊ธฐ๋ก
shm = shared_memory.SharedMemory(create=True, size=1024 * 1024) # 1MB ๊ณต๊ฐ„
buffer = shm.buf

# ์ง๋ ฌํ™”๋œ ๋ฐ”์ด๋„ˆ๋ฆฌ ๋ฐ์ดํ„ฐ๋ฅผ ๊ณต์œ  ๋ฉ”๋ชจ๋ฆฌ์— ์ง์ ‘ ๋ณต์‚ฌ
raw_bytes = b"Telemetry High Velocity Packet Data..."
buffer[:len(raw_bytes)] = raw_bytes

# 2. Worker ํ”„๋กœ์„ธ์Šค์—๋Š” ๋ฉ”ํƒ€๋ฐ์ดํ„ฐ(๋ฉ”๋ชจ๋ฆฌ ์ด๋ฆ„๊ณผ ๊ธธ์ด)๋งŒ Queue๋กœ ์ „๋‹ฌ (์ง๋ ฌํ™” ๋น„์šฉ ์ตœ์†Œํ™”)
ipc_metadata = {"shm_name": shm.name, "data_size": len(raw_bytes)}
# Queue์—๋Š” ์•„์ฃผ ์ž‘๊ณ  ๊ฒฝ๋Ÿ‰ํ™”๋œ Dict๋งŒ ์ „๋‹ฌ๋˜๋ฏ€๋กœ ์ด๋ฒคํŠธ ๋ฃจํ”„ ์ง€์—ฐ์ด 1ms ์ดํ•˜๋กœ ์œ ์ง€๋จ

์ „๋žต 4: ๋งˆ์ดํฌ๋กœ ๋ฐฐ์น˜(Micro-batching) ๋ฐ๋ง ๋ฒ„ํผ(Ring Buffer) ํŒจํ„ด

๋‹จ์ผ ๋ฐ์ดํ„ฐ ๋‹จ์œ„๋กœ IPC๋ฅผ ์ˆ˜ํ–‰ํ•˜๋Š” ๋Œ€์‹ , asyncio ์ธก์—์„œ ๋ฐ์ดํ„ฐ๋ฅผ ์†Œํ˜• ๋ฉ”๋ชจ๋ฆฌ ๋ฒ„ํผ์— ๋ชจ์€ ํ›„ ์ผ์ • ์‹œ๊ฐ„Interval(์˜ˆ: 10ms)์ด๋‚˜ ๊ฐœ์ˆ˜Threshold(์˜ˆ: 500๊ฑด)์— ๋„๋‹ฌํ–ˆ์„ ๋•Œ ๋‹จ 1ํšŒ์˜ IPC ํ˜ธ์ถœ๋กœ ์ „์†กํ•˜๋Š” **๋ฐฐ์น˜ ๋ฌถ์Œ ์ฒ˜๋ฆฌ(Micro-Batching)**๋ฅผ ์ ์šฉํ•ฉ๋‹ˆ๋‹ค. ์ด๋กœ ์ธํ•ด ์ดˆ๋‹น ๋ฐœ์ƒํ•˜๋Š” ์ง๋ ฌํ™” ํšŸ์ˆ˜ ์ž์ฒด๊ฐ€ 1/500 ์ดํ•˜๋กœ ๊ฐ์†Œํ•ฉ๋‹ˆ๋‹ค.

 

4. ์‹ค๋ฌด ์šด์˜ ์‹œ ์—ฃ์ง€ ์ผ€์ด์Šค ๋ฐ ๋ถ€์ž‘์šฉ(Pitfalls) ๊ด€๋ฆฌ

์œ„ ์ตœ์ ํ™” ์ „๋žต์„ ํ˜„์—… ์šด์˜ ํ™˜๊ฒฝ์— ๋„์ž…ํ•  ๋•Œ ๋ฐ˜๋“œ์‹œ ๋Œ€๋น„ํ•ด์•ผ ํ•˜๋Š” ๋ฌธ์ œ์ ๋“ค๊ณผ ๊ทธ ๋Œ€์ฒ˜ ๋ฐฉ์•ˆ์ž…๋‹ˆ๋‹ค.

1) SharedMemory ๋ฉ”๋ชจ๋ฆฌ ๋ˆ„์ˆ˜(Memory Leak)์™€ ํ”„๋กœ์„ธ์Šค ์ถฉ๋Œ

SharedMemory๋Š” OS ์ฐจ์›์˜ RAM ์ž์›์„ ์ง์ ‘ ์ ์œ ํ•ฉ๋‹ˆ๋‹ค. ์ž‘์—…์„ ๋งˆ์นœ ์ž‘์—…์ž ํ”„๋กœ์„ธ์Šค์—์„œ shm.close()๋ฅผ ์ˆ˜ํ–‰ํ•˜๊ณ , ์ตœ์ข…์ ์œผ๋กœ ๋ฆฌ์†Œ์Šค๋ฅผ ํ•ด์ œํ•˜๋Š” ํ”„๋กœ์„ธ์Šค์—์„œ shm.unlink()๋ฅผ ๋ช…์‹œ์ ์œผ๋กœ ํ˜ธ์ถœํ•˜์ง€ ์•Š์œผ๋ฉด, ์˜ˆ์™ธ(Exception) ๋ฐœ์ƒ ์‹œ ๋น„๋””์˜ค ๋ฉ”๋ชจ๋ฆฌ๋‚˜ ์‹œ์Šคํ…œ RAM์ด ๋ฐ”๋‹ฅ๋‚˜๋Š” ๊ณ ์งˆ์ ์ธ ๋ˆ„์ˆ˜ ํ˜„์ƒ์ด ์ผ์–ด๋‚ฉ๋‹ˆ๋‹ค.

ํ•ด๊ฒฐ์ฑ…: try...finally ๋ธ”๋ก์„ ๋ณด์žฅํ•˜๊ฑฐ๋‚˜ contextlib.contextmanager๋ฅผ ํ™œ์šฉํ•ด ๋ฆฌ์†Œ์Šค ํ•ด์ œ๋ฅผ ๊ฐ•์ œํ•ด์•ผ ํ•ฉ๋‹ˆ๋‹ค.

2) ๋ฐฐ์„ธํ”„๋ ˆ์…”(Backpressure) ๋ฏธ๋น„๋กœ ์ธํ•œ Queue Full ๋ฐ๋“œ๋ฝ

์ž‘์—…์ž ํ”„๋กœ์„ธ์Šค์˜ CPU ์ฒ˜๋ฆฌ ์†๋„๋ณด๋‹ค asyncio์˜ ๋ฐ์ดํ„ฐ ์ˆ˜์ง‘ ์†๋„๊ฐ€ ์›”๋“ฑํžˆ ๋น ๋ฅด๋ฉด multiprocessing.Queue๊ฐ€ ๋ฌดํ•œ์ • ํŒฝ์ฐฝํ•˜์—ฌ OOM(Out of Memory)์ด ๋ฐœ์ƒํ•ฉ๋‹ˆ๋‹ค. ๋ฐ˜๋Œ€๋กœ Queue(maxsize=1000)๋กœ ํฌ๊ธฐ๋ฅผ ์ œํ•œํ•ด ๋‘๋ฉด Queue.put()์ด ๋ฌดํ•œ ๋ธ”๋กœํ‚น๋˜์–ด ์ด๋ฒคํŠธ ๋ฃจํ”„๊ฐ€ ์™„์ „ํžˆ ๋ฉˆ์ถ”๊ฒŒ ๋ฉ๋‹ˆ๋‹ค.

ํ•ด๊ฒฐ์ฑ…: non-blocking ๋ฐฉ์‹์ธ Queue.put_nowait()์„ ์‹คํ–‰ํ•˜๊ณ , queue.Full ์˜ˆ์™ธ ๋ฐœ์ƒ ์‹œ ์ˆ˜์ง‘ ์ธก์—์„œ ์š”์ฒญ์„ ๊ฑฐ์ ˆ(Drop)ํ•˜๊ฑฐ๋‚˜ ํด๋ผ์ด์–ธํŠธ์—๊ฒŒ Throttling ์‘๋‹ต์„ ๋ณด๋‚ด๋Š” **๋ฐฐํ”„๋ ˆ์…” ์ œ์–ด ํ๋ฆ„(Backpressure Control)**์„ ๊ตฌํ˜„ํ•ด์•ผ ํ•ฉ๋‹ˆ๋‹ค.

 



 

5. ๊ฒฐ๋ก  ๋ฐ ์š”์•ฝ

Python์˜ asyncio์™€ multiprocessing์„ ๊ณ ์„ฑ๋Šฅ ํ…”๋ ˆ๋ฉ”ํŠธ๋ฆฌ ํŒŒ์ดํ”„๋ผ์ธ์— ํ˜ผ์šฉํ•  ๋•Œ ๋ฐœ์ƒํ•˜๋Š” ์ด๋ฒคํŠธ ๋ฃจํ”„ ์ง€์—ฐ์˜ ํ•ต์‹ฌ ๋ฒ”์ธ์€ ๋™๊ธฐ์‹ IPC ์ง๋ ฌํ™”(Pickle) ์—ฐ์‚ฐ์ž…๋‹ˆ๋‹ค.

์•ˆ์ •์ ์ธ ์‹ค์‹œ๊ฐ„ ํŒŒ์ดํ”„๋ผ์ธ ๊ตฌ์ถ•์„ ์œ„ํ•ด ๋‹ค์Œ 3๊ฐ€์ง€ ํ•ต์‹ฌ ์š”์†Œ๋ฅผ ๋ฐ˜๋“œ์‹œ ๊ธฐ์–ตํ•˜๊ณ  ์ ์šฉํ•ด์•ผ ํ•ฉ๋‹ˆ๋‹ค:

  1. ์ง๋ ฌํ™” ๋น„์šฉ ์ธก์ •: ์ด๋ฒคํŠธ ๋ฃจํ”„ ๋“œ๋ฆฌํ”„ํŠธ ๋ชจ๋‹ˆํ„ฐ๋ง์„ ํ†ตํ•ด pickle์— ์˜ํ•œ ์ง€์—ฐ์„ ์ˆ˜์น˜ํ™”ํ•ฉ๋‹ˆ๋‹ค.
  2. ์ œ๋กœ์นดํ”ผ ๋ฐ ๋ฐ”์ด๋„ˆ๋ฆฌ ์ง๋ ฌํ™” ํ™œ์šฉ: ๋Œ€์šฉ๋Ÿ‰ ๋ฐ์ดํ„ฐ๋Š” SharedMemory ๋ฐ Orjson/MessagePack์„ ๊ฒฐํ•ฉํ•˜์—ฌ ์ง๋ ฌํ™” ๋ถ€ํ•˜๋ฅผ ๊ทน์ ์œผ๋กœ ์ค„์ž…๋‹ˆ๋‹ค.
  3. ๋ฐฐ์น˜ ์ „์†ก ๋ฐ ๋ฐฑํ”„๋ ˆ์…” ์„ค๊ณ„: Queue.put์˜ ๋™๊ธฐ ์ฐจ๋‹จ์„ ๋ฐฉ์ง€ํ•˜๊ธฐ ์œ„ํ•ด micro-batching ๊ธฐ๋ฒ•๊ณผ non-blocking ๋ฐฑํ”„๋ ˆ์…” ์ •์ฑ…์„ ๊ตฌํ˜„ํ•ฉ๋‹ˆ๋‹ค.

์ด๋Ÿฌํ•œ ์•„ํ‚คํ…์ฒ˜ ๊ฐœ์„ ์„ ์ ์šฉํ•˜๋ฉด CPython ํ™˜๊ฒฝ์—์„œ๋„ ์ด๋ฒคํŠธ ๋ฃจํ”„ ์ง€์—ฐ ์‹œ๊ฐ„์„ ๋ฐ€๋ฆฌ์ดˆ(ms) ์ดํ•˜ ์ˆ˜์ค€์œผ๋กœ ์•ˆ์ •์ ์œผ๋กœ ์œ ์ง€ํ•˜๋ฉด์„œ ์ดˆ๋‹น ์ˆ˜๋งŒ ๊ฑด์˜ ํ…”๋ ˆ๋ฉ”ํŠธ๋ฆฌ ๋ฐ์ดํ„ฐ๋ฅผ ๋ณ‘๋ ฌ ์ฒ˜๋ฆฌํ•  ์ˆ˜ ์žˆ์Šต๋‹ˆ๋‹ค.

๋ฐ˜์‘ํ˜•