Issue 01중국 AI
AC POST
중국 AI 목록
掘金2026년 9월 18일 09:10중국어 → 한국어

AI Agent 도구 호출, DAG 병렬 오케스트레이션 실무

완전 동시 실행의 경쟁 조건과 직렬 처리의 지연을 피하려고 DAG 의존성 그래프로 도구 호출을 스케줄링하며, 스케줄러 구현과 LangGraph·Temporal 통합, 핵심 경로 분석을 담았다.

중국어 원문을 AI로 번역했습니다. 고유명사와 수치는 원문 표기를 우선하며, 중요한 판단에는 아래 출처 원문을 함께 확인하세요.

나를 깊이 각인시킨 한 번의 프로덕션 장애

작년 우리 팀이 유지보수하던 AI 리서치 리포트 생성 Agent는 리포트 한 부를 생성할 때마다 7개의 도구를 호출해야 했다. 시장 데이터 가져오기, 경쟁사 정보 조회, 뉴스 요약 검색, 재무 데이터 획득, 기술 지표 분석, 업종 벤치마크 가져오기, 사용자 과거 선호도 읽기였다.

출시 초기에는 이 7개의 도구 호출이 순차 실행되었다. 평균적으로 리포트 한 부에 22초를 기다려야 했다. 사용자는 "너무 느리다"고 불평했고, 프로덕트 매니저는 P95 지연을 뚫어지게 보며 발을 동동 굴렀다. 우리는 첫 번째 버전의 최적화를 했다. 7개 도구 전부를 asyncio.gather로 동시에 병렬 실행하도록 바꾼 것이다. 지연은 22초에서 6초로 줄었다. 그러나 곧바로 새로운 유형의 장애가 나타났다. 바로 간헐적 데이터 불일치였다.

문제는 두 도구 사이에 암묵적 의존이 있다는 점이었다. get_financial_data는 get_company_profile이 반환한 ticker symbol을 사용해야 했다. 이 두 도구가 동시에 시작되면, get_financial_data가 가끔 get_company_profile보다 먼저 완료되어 이전 캐시의 ticker를 사용했고, 최종적으로 잘못된 회사의 재무 데이터를 받아왔다.

이 bug는 테스트 환경에서는 한 번도 재현되지 않았다. 테스트할 때는 매번 ticker가 동일했기 때문이다. 프로덕션에서 ticker가 바뀔 때 두 도구의 타이밍이 불확정해지면서 bug가 드러났다.

이것이 바로 도구 호출 DAG 오케스트레이션 문제의 전형적인 출발점이다. 당신이 "병렬로 처리할 수 있다"고 생각한 것이 사실은 숨겨진 의존 관계를 가지고 있고, 당신이 "직렬이 안전하다"고 생각한 것이 사실은 겹쳐서 실행할 수 있었던 대량의 시간을 헛되이 버리고 있는 것이다.

왜 순차 호출과 완전 병렬 모두 충분하지 않은가

LLM 도구 호출의 실행 모드는 대략 세 가지로 나뉜다.

순차 실행(Sequential)

Tool1 → Tool2 → Tool3 → Tool4

장점: 구현이 간단하고 의존 문제가 없다. 단점: 동시성 기회를 완전히 포기한다. 만약 Tool1이 800ms, Tool2가 600ms, Tool3이 700ms 걸리고 세 도구 사이에 의존이 없다면, 당신은 무려 2100ms를 기다린 셈이다.

전량 병렬(Bulk Parallel)

Tool1 ─┐ Tool2 ─┤→ join → 다음 단계 Tool3 ─┘

장점: 총 대기 시간이 가장 짧다(이론상). 단점: 의존 관계가 있는 도구들 사이의 데이터 흐름을 깨뜨려 경쟁 조건(race condition)을 유발한다.

DAG 스케줄링(Dependency-Aware Parallel)

Tool1 ──────────────────┐ Tool2 ─→ Tool4 ─────────┤→ join → 다음 단계 Tool3 ─→ Tool5 ─→ Tool6 ─┘

의존 토폴로지 순서에 따라, 의존이 없는 도구는 병렬로 돌리고 의존이 있는 도구는 상류가 완료된 뒤 시작한다. 이것이 "가장 빠른 실행"과 "의존의 정확성"을 동시에 만족하는 유일한 방안이다.

현실의 Agent 도구 호출 그래프는 거의 대부분 DAG이며, 선형 시퀀스도 아니고 완전 병렬 팬아웃도 아니다. 그러나 대부분의 엔지니어링 구현은 가장 손쉬운 양 끝 중 하나를 골랐고, 중간에 있는 올바른 선택지를 고르지 않았다.

DAG 도구 스케줄링의 핵심 개념

1. 의존 선언과 자동 추론

도구 사이에는 두 종류의 의존이 있다.

명시적 의존(Explicit Dependency): Tool B의 입력이 Tool A의 출력 필드를 직접 참조한다. 이는 도구 schema의 $ref나 파라미터 바인딩에서 정적으로 추론할 수 있다.

도구 호출 계획(LLM 생성) tool_plan = [ { "id" : "t1" , "tool" : "get_company_profile" , "args" : { "company_name" : "{{user_query.company}}" }, "deps" : [] }, { "id" : "t2" , "tool" : "get_financial_data" , "args" : { "ticker" : "{{t1.result.ticker}}" }, # t1에 의존 "deps" : [ "t1" ] }, { "id" : "t3" , "tool" : "get_news" , "args" : { "query" : "{{user_query.company}}" }, "deps" : [] # t1에 의존하지 않으므로 동시 실행 가능 }, { "id" : "t4" , "tool" : "analyze_sentiment" , "args" : { "news" : "{{t3.result.articles}}" , # t3에 의존 "financials" : "{{t2.result}}" # t2에 의존 }, "deps" : [ "t2" , "t3" ] # t2와 t3가 모두 완료되어야 함 } ]

암묵적 의존(Implicit Dependency): 두 도구가 동일한 외부 상태(예: 데이터베이스, 파일)에 쓰기를 하므로 일관성을 보장하기 위해 순서를 정해야 한다. 이런 유형의 의존은 파라미터 참조에서 추론할 수 없으며, 도구 schema에 명시적으로 표기해야 한다.

@tool( writes=[ "user_profile.preferences" ], reads=[ "user_profile.history" ] ) def update_user_preferences ( user_id: str , new_prefs: dict ) -> dict : ...

만약 두 도구가 동일한 writes 필드를 선언했다면, 스케줄러는 이들을 강제로 직렬화해야 한다.

2. 토폴로지 정렬과 임계 경로

의존 그래프가 있으면 Kahn 알고리즘으로 토폴로지 정렬을 수행하여 합법적인 실행 계층(layer)을 얻는다.

from collections import defaultdict, deque def compute_execution_layers ( tool_plan: list [ dict ] ) -> list [ list [ str ]]: """도구 계획을 의존 관계에 따라 실행 계층으로 나누고, 같은 계층 내에서는 병렬 실행 가능하게 한다.""" in_degree = {t[ "id" ]: 0 for t in tool_plan} graph = defaultdict( list ) for tool in tool_plan: for dep in tool.get( "deps" , []): graph[dep].append(tool[ "id" ]) in_degree[tool[ "id" ]] += 1 # Kahn 알고리즘 queue = deque([tid for tid, deg in in_degree.items() if deg == 0 ]) layers = [] while queue: layer = [] for _ in range ( len (queue)): node = queue.popleft() layer.append(node) for neighbor in graph[node]: in_degree[neighbor] -= 1 if in_degree[neighbor] == 0 : queue.append(neighbor) layers.append(layer) if sum ( len (l) for l in layers) != len (tool_plan): raise ValueError( "도구 의존 그래프에 순환 의존이 존재합니다!" ) return layers # 예시 출력 # Layer 0: ["t1", "t3"] <- 의존 없음, 동시 실행 # Layer 1: ["t2"] <- t1 완료에 의존 # Layer 2: ["t4"] <- t2와 t3 완료에 의존

임계 경로(Critical Path)는 소스에서 싱크까지의 가장 긴 소요 시간 경로이며, 이론상 최단 완료 시간을 결정한다. 프로덕션에서는 임계 경로 위의 도구에 최적화 여지가 있는지 모니터링해야 한다. 임계 경로가 아닌 위의 도구를 최적화하는 것은 총 지연에 아무런 기여도 하지 않는다.

def compute_critical_path ( tool_plan: list [ dict ], estimated_durations: dict [ str , float ] ) -> tuple [ list [ str ], float ]: """임계 경로의 도구 시퀀스와 총 소요 시간을 반환한다.""" # 동적 계획법: est[id] = 최조 완료 시간 est = {} pred = {} # 전임자 추적 tool_by_id = {t[ "id" ]: t for t in tool_plan} def earliest_finish ( tid: str ) -> float : if tid in est: return est[tid] deps = tool_by_id[tid].get( "deps" , []) if not deps: est[tid] = estimated_durations.get(tid, 1.0 ) pred[tid] = None return est[tid] max_dep_finish = 0 max_dep_id = None for dep in deps: dep_finish = earliest_finish(dep) if dep_finish > max_dep_finish: max_dep_finish = dep_finish max_dep_id = dep est[tid] = max_dep_finish + estimated_durations.get(tid, 1.0 ) pred[tid] = max_dep_id return est[tid] for tool in tool_plan: earliest_finish(tool[ "id" ]) # 가장 늦게 완료된 노드에서 역추적 end_node = max (est, key=est.get) path = [] node = end_node while node is not None : path.append(node) node = pred.get(node) return list ( reversed (path)), est[end_node]

프로덕션 스케줄러의 완전한 구현

이론은 이해했으니, 실제로 프로덕션에 사용할 수 있는 스케줄러를 살펴보자.

import asyncio import time import traceback from dataclasses import dataclass, field from typing import Any, Callable, Awaitable from enum import Enum

class ToolStatus(Enum): PENDING = "pending" RUNNING = "running" SUCCESS = "success" FAILED = "failed" SKIPPED = "skipped" # 상위 실패로 인한 건너뜀

@dataclass class ToolResult: tool_id: str status: ToolStatus result: Any = None error: Exception | None = None started_at: float = 0.0 finished_at: float = 0.0

@property def duration_ms(self) -> float: return (self.finished_at - self.started_at) * 1000

class DAGToolScheduler: """ 프로덕션급 DAG 도구 스케줄러. 지원: 계층 병렬 실행, 인자 템플릿 바인딩, 실패 전파 정책, 타임아웃, 관측. """

def __init__( self, failure_policy: str = "stop_on_required_failure", global_timeout_s: float = 30.0, on_tool_start: Callable | None = None, on_tool_done: Callable | None = None, ): self.failure_policy = failure_policy self.global_timeout_s = global_timeout_s self.on_tool_start = on_tool_start self.on_tool_done = on_tool_done

def _resolve_args(self, args: dict, context: dict) -> dict: """{{tool_id.field.path}} 자리표시자를 실제 값으로 치환한다.""" import re

def resolve_value(val): if not isinstance(val, str): return val

pattern = r'\{\{([^}]+)\}\}' matches = re.findall(pattern, val) if not matches: return val

값 전체가 하나의 자리표시자라면 원래 타입을 반환한다 (str로 강제 변환하지 않음) if val.strip() == f"{{{{ {matches[0]} }}}}": path = matches[0].strip() return self._get_nested(context, path.split("."))

문자열 삽입(interpolation) def replacer(m): path = m.group(1).strip() value = self._get_nested(context, path.split(".")) return str(value) if value is not None else m.group(0)

return re.sub(pattern, replacer, val)

def resolve_dict(d): return {k: resolve_value(v) if not isinstance(v, dict) else resolve_dict(v) for k, v in d.items()}

return resolve_dict(args)

def _get_nested(self, obj: Any, path: list[str]) -> Any: for key in path: if obj is None: return None if isinstance(obj, dict): obj = obj.get(key) else: obj = getattr(obj, key, None) return obj

def _should_skip(self, tool: dict, results: dict[str, ToolResult]) -> bool: """상위 실패 때문에 이 도구를 건너뛰어야 하는지 판단한다.""" if self.failure_policy == "continue_all": return False

for dep_id in tool.get("deps", []): dep_result = results.get(dep_id) if dep_result is None: return True # 상위가 아직 끝나지 않음 (발생해선 안 됨) if dep_result.status in (ToolStatus.FAILED, ToolStatus.SKIPPED): # 의존성이 required로 표시되었는지 확인 (기본값 required) if tool.get("dep_required", {}).get(dep_id, True): return True return False

async def _execute_tool( self, tool: dict, executor: Callable[[str, dict], Awaitable[Any]], context: dict, results: dict[str, ToolResult], ) -> ToolResult: tool_id = tool["id"] tool_result = ToolResult(tool_id=tool_id, status=ToolStatus.RUNNING, started_at=time.time())

if self._should_skip(tool, results): tool_result.status = ToolStatus.SKIPPED tool_result.finished_at = time.time() return tool_result

if self.on_tool_start: self.on_tool_start(tool_id, tool["tool"])

try: resolved_args = self._resolve_args(tool.get("args", {}), context) timeout = tool.get("timeout_s", self.global_timeout_s) result = await asyncio.wait_for( executor(tool["tool"], resolved_args), timeout=timeout, ) tool_result.result = result tool_result.status = ToolStatus.SUCCESS except asyncio.TimeoutError: tool_result.error = TimeoutError(f"도구 {tool_id} 초과(> {timeout}s)") tool_result.status = ToolStatus.FAILED except Exception as e: tool_result.error = e tool_result.status = ToolStatus.FAILED finally: tool_result.finished_at = time.time() if self.on_tool_done: self.on_tool_done(tool_id, tool_result)

return tool_result

async def run( self, tool_plan: list[dict], executor: Callable[[str, dict], Awaitable[Any]], initial_context: dict | None = None, ) -> dict[str, ToolResult]: """ DAG 토폴로지에 따라 도구 계획을 실행한다. 각 도구의 ToolResult를 반환하며, 호출자는 여기서 result 또는 error를 읽을 수 있다. """ results: dict[str, ToolResult] = {} context = {"user_query": initial_context or {}} layers = compute_execution_layers(tool_plan)

async def run_layer(layer: list[str]): tool_by_id = {t["id"]: t for t in tool_plan} tasks = [ asyncio.create_task( self._execute_tool(tool_by_id[tid], executor, context, results) ) for tid in layer ] layer_results = await asyncio.gather(*tasks, return_exceptions=False) for tool_result in layer_results: results[tool_result.tool_id] = tool_result # 성공한 결과를 context에 주입하여 이후 계층의 인자 바인딩에 사용 if tool_result.status == ToolStatus.SUCCESS: context[tool_result.tool_id] = {"result": tool_result.result}

try: await asyncio.wait_for( asyncio.gather(*[run_layer(layer) for layer in [layers[0]]]), timeout=self.global_timeout_s, ) # 계층별로 실행한다 (모든 계층을 gather할 수 없다 — 계층 간에 순서 의존성이 있다) for layer in layers: await run_layer(layer) except asyncio.TimeoutError: # 전역 타임아웃: 완료되지 않은 모든 도구를 FAILED로 표시 for tool in tool_plan: if tool["id"] not in results: results[tool["id"]] = ToolResult( tool_id=tool["id"], status=ToolStatus.FAILED, error=TimeoutError("전역 타임아웃"), started_at=time.time(), finished_at=time.time(), )

return results

시스템을 폭파시키는 세 가지 프로덕션 함정

함정 1: asyncio.gather를 DAG 스케줄링으로 착각하기

잘못된 예시: "모두 의존성이 없다"고 가정한 경우 results = await asyncio.gather( call_tool( "get_company_profile" , { "name" : company}), call_tool( "get_financial_data" , { "ticker" : ???}), # ticker는 어디서 오는가? call_tool( "get_news" , { "query" : company}), )

LLM이 생성한 도구 호출 목록에 파라미터 참조 관계가 있을 때는, 먼저 의존성 그래프를 분석한 다음 어떤 것을 병렬 실행할 수 있는지 결정해야 한다. 무작정 gather 하면 안 된다.

수정: 도구 호출 계획 생성 단계에서 LLM이 deps 필드(또는 파라미터 `{{ref}}` 형식)를 명시적으로 출력하도록 요구하고, 스케줄러가 실행 전에 정적 분석을 수행한다.

함정 2: 실패 전파 전략의 비대칭

도구가 5개 있고, 그중 t2가 실패했다고 가정하자. 두 가지 극단적 전략 모두 문제가 있다:

- 즉시 전면 중단(fail-fast all): t3, t4, t5(t2와 무관)까지 취소된다. 최종 보고서에는 무관한 데이터조차 남지 않는다.

- 전부 계속(continue all): t4는 t2의 출력에 의존하는데, t2가 실패한 후 t4가 `None`을 받아 `NullPointerError`나 데이터 오류가 발생하고, 오히려 더 찾기 어려운 2차 장애를 만들어낸다.

올바른 전략은 의존성 인식 실패 전파다:

failure_policy_matrix = { "stop_on_required_failure" : True , # 의존성 실패 → 다운스트림 건너뛰기(기본값) "continue_all" : False , # 전파하지 않고, 다운스트림은 None으로 계속 "stop_all" : True , # 어떤 실패든 → 모든 작업 취소 } # 도구 수준의 required/optional 표기 tool_plan = [ { "id" : "t2" , "tool" : "get_financial_data" , ... }, { "id" : "t4" , "tool" : "analyze_sentiment" , "deps" : [ "t2" , "t3"], "dep_required" : { "t2" : True , # t2 실패 → t4 건너뛰기 "t3" : False , # t3 실패 → t4는 여전히 실행(t2 결과만 사용) } } ]

함정 3: 타임아웃 예산이 하향 전달되지 않음

전역 timeout을 30초로 설정했지만, 각 계층 실행 시 이미 사용한 시간을 차감하는 것을 아무도 기억하지 못한다. 그 결과:

- Layer 0이 25초 실행됨

- Layer 1도 각 도구에 30초 timeout을 할당함

- 실제 요청 체인은 55초에 타임아웃되고, 업스트림의 deadline은 이미 훨씬 지나버렸다

올바른 방법은 남은 deadline 전달을 사용하는 것이다:

import time async def run_with_deadline ( tool_plan: list [ dict ], executor, deadline: float , # absolute timestamp ): layers = compute_execution_layers(tool_plan) results = {} context = {} for layer in layers: remaining = deadline - time.time() if remaining <= 0 : # deadline 초과, 모든 미실행 도구를 SKIPPED로 표시 for t in tool_plan: if t[ "id" ] not in results: results[t[ "id" ]] = ToolResult( tool_id=t[ "id" ], status=ToolStatus.SKIPPED, error=TimeoutError( "Deadline exceeded before execution" ), started_at=time.time(), finished_at=time.time(), ) break # 각 도구는 최대 remaining / layer_parallelism 시간만 사용 per_tool_timeout = min (remaining * 0.8 , 10.0 ) tasks = [ run_tool_with_timeout(t, executor, context, per_tool_timeout) for t in layer_tools(layer, tool_plan) ] layer_results = await asyncio.gather(*tasks) for r in layer_results: results[r.tool_id] = r if r.status == ToolStatus.SUCCESS: context[r.tool_id] = { "result" : r.result} return results

관측성: 스케줄링 과정을 실제로 디버깅 가능하게 만들기

도구 호출 그래프가 한 번 돌고 나면, 당신은 알아야 한다: 어떤 도구가 병렬로 실행되었는가? 어느 것이 크리티컬 패스인가? 어디서 가장 많은 시간이 소요되었는가?

OpenTelemetry로 각 계층과 각 도구에 Span을 붙인다:

from opentelemetry import trace tracer = trace.get_tracer( "tool-dag-scheduler" ) async def _execute_tool_traced ( self, tool, executor, context, results ): with tracer.start_as_current_span( f"tool. {tool[ 'tool' ]} " , attributes={ "tool.id" : tool[ "id" ], "tool.name" : tool[ "tool" ], "tool.deps" : "," .join(tool.get( "deps" , [])), } ) as span: result = await self._execute_tool(tool, executor, context, results) span.set_attribute( "tool.status" , result.status.value) span.set_attribute( "tool.duration_ms" , result.duration_ms) if result.error: span.record_exception(result.error) return result

출력된 Trace는 Jaeger에서 중첩 Span으로 표시되고, 같은 계층의 도구는 같은 시간 구간에 나타나며, 의존 관계는 Span의 선후 순서에 드러난다. Trace에서 당신은 곧바로 알 수 있다:

- 크리티컬 패스가 실제로 어느 것인지(가장 긴 연속 Span 체인)

- 어떤 도구가 의존성 대기 때문에 시간을 헛되이 잃었는지

- 병렬 효과가 얼마나 되는지(계층 내 가장 넓은 Span 너비 vs 직렬이라고 가정했을 때의 총 너비)

하나의 실제 Trace 비교 수치(우리 내부 측정):

시나리오 | 7개 도구 총 지연 | 병렬도 전체 직렬 | 18.4초 | 1.0x 전체 병렬(오류) | 6.1초(+ 데이터 오류) | 3.0x DAG 스케줄링 | 7.3초(오류 없음) | 2.5x

DAG 스케줄링은 정확성을 훼손하지 않는 전제하에 "잘못된 전체 병렬"보다 1.2초만 느릴 뿐이지만, 이전의 경쟁 상태 장애를 완전히 제거했다.

주류 프레임워크와의 통합

LangGraph

LangGraph는 fan-out/fan-in 패턴을 네이티브로 지원하지만, 도구 수준의 의존성은 스스로 모델링해야 한다:

from langgraph.graph import StateGraph, END from typing import TypedDict, Annotated import operator class AgentState ( TypedDict ): company_name: str company_profile: dict | None financial_data: dict | None news_articles: list | None sentiment_analysis: dict | None def fetch_company_profile ( state: AgentState ) -> AgentState: # t1: 의존성 없음 result = get_company_profile(state[ "company_name" ]) return { "company_profile" : result} def fetch_financial_data ( state: AgentState ) -> AgentState: # t2: t1에 의존(state를 통해 전달) ticker = state[ "company_profile" ][ "ticker" ] result = get_financial_data(ticker) return { "financial_data" : result} def fetch_news ( state: AgentState ) -> AgentState: # t3: t1에 의존하지 않으므로 t2와 동시 실행 가능 result = get_news(state[ "company_name" ]) return { "news_articles" : result} def analyze_sentiment ( state: AgentState ) -> AgentState: # t4: t2와 t3에 의존 result = analyze(state[ "financial_data" ], state[ "news_articles" ]) return { "sentiment_analysis" : result} # 그래프 구축 builder = StateGraph(AgentState) builder.add_node( "fetch_profile" , fetch_company_profile) builder.add_node( "fetch_financials" , fetch_financial_data) builder.add_node( "fetch_news" , fetch_news) builder.add_node( "analyze" , analyze_sentiment) # 토폴로지 엣지 builder.set_entry_point( "fetch_profile" ) builder.add_edge( "fetch_profile" , "fetch_financials" ) builder.add_edge( "fetch_profile" , "fetch_news" ) # fan-out # LangGraph는 fetch_financials와 fetch_news가 모두 완료된 후에야 analyze로 진입한다 builder.add_edge( "fetch_financials" , "analyze" ) builder.add_edge( "fetch_news" , "analyze" ) # fan-in builder.add_edge( "analyze" , END) graph = builder. compile ()

LangGraph의 super-step 메커니즘은 같은 배치로 스케줄된 노드들을 자동으로 식별한다(fetch_financials와 fetch_news는 모두 fetch_profile에만 의존하므로, fetch_profile이 완료된 후 같은 super-step 안에서 동시에 실행된다). 주의: 이것은 엄격한 asyncio 동시성이 아니다. 같은 super-step 안의 노드들은 서로 다른 스레드에서 실행되지만 여전히 Python GIL의 제약을 받으므로 CPU 집약적인 작업에는 도움이 되지 않는다.

Temporal Workflow(프로덕션 중량급 작업)

영속화와 재시작 간 복구가 필요한 도구 DAG에는 Temporal이 더 적합한 선택이다:

import asyncio from temporalio import workflow, activity from temporalio.common import RetryPolicy @activity.defn async def get_company_profile ( company_name: str ) -> dict : ... @activity.defn async def get_financial_data ( ticker: str ) -> dict : ... @activity.defn async def get_news ( query: str ) -> list : ... @workflow.defn class ResearchReportWorkflow : @workflow.run async def run ( self, company_name: str ) -> dict : retry = RetryPolicy(maximum_attempts= 3 , backoff_coefficient= 2.0 ) # t1: 串行 profile = await workflow.execute_activity( get_company_profile, company_name, retry_policy=retry, start_to_close_timeout=timedelta(seconds= 10 ), ) # t2 和 t3: 并发 financial_task = asyncio.ensure_future( workflow.execute_activity( get_financial_data, profile[ "ticker" ], retry_policy=retry, start_to_close_timeout=timedelta(seconds= 15 ), ) ) news_task = asyncio.ensure_future( workflow.execute_activity( get_news, company_name, retry_policy=retry, start_to_close_timeout=timedelta(seconds= 8 ), ) ) financials, news = await asyncio.gather(financial_task, news_task) # t4: 等 t2、t3 完成 sentiment = await workflow.execute_activity( analyze_sentiment, args=[financials, news], retry_policy=retry, start_to_close_timeout=timedelta(seconds= 10 ), ) return { "profile" : profile, "financials" : financials, "news" : news, "sentiment" : sentiment}

Temporal의 장점은 워크플로 상태가 Event History에 영속화되어, 프로세스가 크래시하더라도 재시작 후 마지막 checkpoint부터 이어서 진행하며 이미 완료된 Activity는 다시 실행되지 않는다는 것이다. 10초를 초과하고 외부 API 호출이 얽힌 도구 DAG에는 이것이 순수 asyncio보다 더 신뢰할 수 있는 선택이다.

언제 어떤 방식을 쓸까: 의사 결정 매트릭스

시나리오 권장 방식 이유 도구 5개 이하, 강한 의존 없음, 단일 요청 10초 미만 asyncio + 수작업 의존성 검사 경량, 충분함 도구 5~20개, 복잡한 의존 그래프, 단일 요청 10~60초 본문의 DAGToolScheduler 범용, 확장 가능 도구 20개 초과, 또는 요청 간 영속화 필요 Temporal / Prefect 영속화, 재시도, 가시성 프레임워크가 이미 LangGraph LangGraph super-step 네이티브 지원, 새 의존성 도입 없음 도구 호출 그래프가 런타임에 동적으로 변함 DAGToolScheduler + 동적 계획 재생성 정적 그래프로는 적응 불가

소결

도구 호출 병렬화의 핵심은 "모든 도구를 gather에 던져넣는 것"이 아니라, 먼저 의존 관계를 명확히 파악한 뒤 의존 관계가 없는 도구들을 동시에 실행하게 하는 것이다.

구체적인 실행 제안:

- LLM 도구 계획 생성 단계에서부터 deps 필드를 요구하고, 사후에 추론하지 말 것.

- Kahn 위상 정렬로 계획을 계층화하고, 같은 계층은 동시 실행, 계층 간은 순차 실행.

- 실패 전파는 required/optional 의존을 구분하여, 도구 하나의 실패가 전체 계획을 붕괴시키거나 2차 오류를 낳지 않도록 할 것.

- 고정 timeout이 아니라 남은 deadline으로 각 계층의 시간을 배분할 것.

- 각 도구에 Span을 부여하여 Trace에서 크리티컬 패스를 보고, 데이터에 근거해 최적화할 것.

마지막으로 기억할 만한 숫자 하나: 내부 측정에서 올바른 DAG 스케줄링은 전면 순차 실행보다 2.5배 빨랐고, 동시에 "잘못된 전면 동시 실행"보다는 20%도 안 되게 느렸다. 이 20%의 대가는 가치가 있다——그 대가로 얻는 것은 제로 레이스 장애, 그리고 Trace에서 실행 과정을 실제로 이해할 수 있는 디버깅 가능성이다.

AINative 소프트웨어 엔지니어링

AI Infrastructure Engineer

144

32k

읽음

28

팔로워