LangGraph 런타임
LangGraph 런타임 (Pregel)
LangGraph에는 애플리케이션의 실행을 관리하는 런타임이 있어요. 그 주인공이 바로 Pregel이에요.
StateGraph를 컴파일하거나 @entrypoint를 만들면 그 결과물이 Pregel 인스턴스가 돼요. 이 인스턴스에 입력을 넣어 호출하는 식으로 LangGraph 애플리케이션을 실행해요.
이 가이드는 런타임을 높은 수준에서 설명하고, Pregel로 애플리케이션을 직접 구현하는 방법도 짚어줄게요.
참고:
Pregel런타임의 이름은 Google의 Pregel 알고리즘에서 따온 거예요. 그래프를 이용한 대규모 병렬 컴퓨팅의 효율적인 방법을 설명한 알고리즘이죠.
출처: 문서
본문
개요 (Overview)
LangGraph에서 Pregel은 액터(actors)와 **채널(channels)**을 하나의 애플리케이션으로 결합해요. 액터는 채널에서 데이터를 읽고 채널에 데이터를 써요. Pregel은 Pregel 알고리즘/Bulk Synchronous Parallel 모델을 따라 애플리케이션의 실행을 여러 단계로 조직해요.
각 단계는 세 가지 단계로 이루어져요.
- Plan(계획): 이번 단계에서 어떤 액터를 실행할지 정해요. 예를 들어 첫 단계에서는 특수한 입력 채널을 구독하는 액터들을 고르고, 이후 단계에서는 이전 단계에서 갱신된 채널을 구독하는 액터들을 골라요.
- Execution(실행): 선택한 모든 액터를 병렬로 실행해요. 모두 끝나거나, 하나가 실패하거나, 타임아웃에 도달할 때까지요. 이 단계에서 채널의 갱신은 다음 단계까지 액터에게 보이지 않아요.
- Update(갱신): 이번 단계에서 액터들이 쓴 값으로 채널을 갱신해요.
실행할 액터가 없어지거나, 최대 단계 수에 도달할 때까지 이 과정을 반복해요.
액터 (Actors)
액터는 PregelNode예요. 채널을 구독하고, 채널에서 데이터를 읽고, 채널에 데이터를 써요. Pregel 알고리즘의 액터라고 생각하면 돼요. PregelNode는 LangChain의 Runnable 인터페이스를 구현해요.
채널 (Channels)
채널은 액터(PregelNode) 간에 통신하는 데 쓰여요. 각 채널은 값 타입, 갱신 타입, 그리고 갱신 함수(일련의 갱신을 받아 저장된 값을 수정하는 함수)를 가져요. 채널은 한 체인에서 다른 체인으로 데이터를 보내거나, 미래의 어떤 단계에서 체인이 자기 자신에게 데이터를 보내는 데 쓸 수 있어요.
LastValue
LastValue는 기본 채널 타입이에요. 마지막으로 쓰인 값을 저장하고, 이전 값을 덮어써요. 입력·출력 값이나, 한 단계에서 다음 단계로 데이터를 넘길 때 쓰면 돼요.
from langgraph.channels import LastValue
channel: LastValue[int] = LastValue(int)
Topic
Topic은 설정 가능한 PubSub 채널이에요. 액터 사이에 여러 값을 보내거나 여러 단계에 걸쳐 출력을 누적할 때 유용해요. 값을 중복 제거하도록 설정하거나, 실행 중에 쓰인 모든 값을 누적하도록 설정할 수 있어요.
from langgraph.channels import Topic
# Accumulate all values written across steps
channel: Topic[str] = Topic(str, accumulate=True)
BinaryOperatorAggregate
BinaryOperatorAggregate는 현재 값과 각 새 갱신에 이항 연산자를 적용해 갱신되는 영속 값을 저장해요. 여러 단계에 걸친 실행 누적(running aggregates)을 계산할 때 써요.
import operator
from langgraph.channels import BinaryOperatorAggregate
# Running total: each write adds to the current value
total = BinaryOperatorAggregate(int, operator.add)
DeltaChannel
DeltaChannel은langgraph>=1.2가 필요하고 현재 베타 상태예요. API는 향후 릴리스에서 바뀔 수 있어요.
DeltaChannel은 전체 누적 값 대신 각 단계의 증분 델타만 저장해요. 자주 쓰이면서 시간이 지나며 큰 값으로 누적되는 채널에 가장 유용해요 — 예를 들어 오래 실행되는 스레드의 대화 메시지 목록 같은 경우죠. 델타 저장이 없으면 매 체크포인트마다 전체 목록을 다시 직렬화해야 해요. DeltaChannel을 쓰면 각 단계에서 새로 쓰여진 메시지만 저장돼요.
자주 쓰이면서 시간이 지날수록 커지는 채널이라면
DeltaChannel을 고려해 보세요. 좋은 신호 하나: 특정 채널의 체크포인트 크기가 스레드 길이에 비례해 선형으로 늘어나는 게 보인다면DeltaChannel이 딱 어울려요.
일반 리듀서(reducer)를 쓰듯 Annotated 타입 어노테이션에서 DeltaChannel을 사용할 수 있어요.
from typing import Annotated, Sequence
from typing_extensions import TypedDict
from langgraph.channels import DeltaChannel
def my_reducer(state: list[str], writes: Sequence[list[str]]) -> list[str]:
result = list(state)
for write in writes:
result.extend(write)
return result
class State(TypedDict):
messages: Annotated[list[str], DeltaChannel(my_reducer)]
벌크 리듀서 요건 (Bulk reducer requirement)
DeltaChannel에 전달하는 reducer는 **벌크 리듀서(bulk reducer)**예요. 현재 상태와 현재 단계에서 온 모든 쓰기의 시퀀스를 한 번의 호출로 받아요. 표준 리듀서처럼 한 번에 하나씩(pairwise) 받지 않아요. 이는 StateGraph의 Annotated와 함께 쓰는 키별 리듀서(갱신마다 한 번씩 호출)와는 다르죠.
벌크 리듀서는 **결합적(associative)**이어야 해요. 즉 배칭에 불변이어야 해요:
reducer(reducer(state, [xs]), [ys]) == reducer(state, [xs, ys])리듀서가 결합적이지 않으면, LangGraph가 여러 단계에 걸쳐 쓰기를 배칭하는 방식에 따라 재구성된 상태가 달라질 수 있어요. 그러면 일관되지 않은 동작이 생기죠.
리듀서는 쓰기 시점이 아니라 재구성 시점에 실행돼요.
BinaryOperatorAggregate의 경우 리듀서가 쓰기 시점에 호출되어 결합된 값이 체크포인트에 직렬화되는 반면,DeltaChannel의 리듀서는 채널 값이 저장된 쓰기에서 *재구성(rebuilt)*될 때 호출돼요. 직렬화되는 건 단계별 원본 쓰기이고, 리듀서는 값이 실체화될 때 — 다음 읽기, 다음 단계의 액터가 접근할 때, 히스토리를 재생할 때 — 비로소 호출돼요.리듀서를 설계할 때 실제로 신경 써야 할 점들:
(state, writes)의 순수 함수로 만들어요. 부작용, 무작위성, 벽시계 읽기(wall-clock reads, 예:uuid.uuid4(),datetime.now())는 값이 재구성될 때마다 실행되어 재생할 때마다 다른 결과를 내요. 이런 것들이 저장된 쓰기에 박혀 있지 않아요.- 들어오는 쓰기에 대한 변형이 저장된다고 믿지 마세요. 리듀서가 쓰기 객체를 변형한다면(예: ID가 없는 항목에 안정적인 ID를 할당), 그 변형은 재구성된 값에만 존재해요. 저장된 쓰기는 여전히 원래 모양이라, 다음 재구성에서는 변형되지 않은 입력을 다시 보게 돼요.
- 정체성(identity)과 그 밖의 안정 메타데이터는 상류(upstream)에서 붙여요. 다운스트림 코드가 여러 턴에 걸쳐 항목을 ID로 참조해야 한다면(예: 나중에 갱신하거나 제거하기 위해), 리듀서 안이 아니라 채널에 값이 쓰이기 전에 그 ID를 할당하세요.
가장 흔한 두 경우에 대한 벌크 리듀서를 보여줄게요.
from typing import Any, Sequence
# List: append all writes in order
def list_reducer(state: list[Any], writes: Sequence[list[Any]]) -> list[Any]:
result = list(state)
for write in writes:
result.extend(write)
return result
# Dict: merge all writes, last write wins on key conflicts
def dict_reducer(
state: dict[str, Any], writes: Sequence[dict[str, Any]]
) -> dict[str, Any]:
result = dict(state)
for write in writes:
result.update(write)
return result
둘 다 결합적이에요. 배치를 하나씩 적용하나 함께 적용하나 같은 결과가 나와요.
읽기 지연 시간을 제한하는 snapshot_frequency
스냅샷이 없으면 DeltaChannel 값을 읽으려면 전체 쓰기 히스토리를 재생해야 해요. N단계짜리 스레드라면 O(N)이죠. snapshot_frequency=K로 설정하면 매 K개의 pregel 단계마다 전체 스냅샷을 써서 읽기 깊이를 최대 K단계로 제한해요.
class State(TypedDict):
messages: Annotated[
list[str],
DeltaChannel(my_reducer, snapshot_frequency=5),
]
snapshot_frequency가 높을수록 저장 오버헤드는 줄지만 읽기 지연 시간은 늘어나요. 낮을수록 지연 시간을 더 단단히 제한하는 대신 체크포인트가 커져요. None(기본값)은 스냅샷을 아예 건너뛰어요 — 읽기가 드물거나 스레드가 짧을 때 적합해요.
버전 호환성과 롤백 (Version compatibility and rollbacks)
지속된 채널을
DeltaChannel에서 비-델타 채널로 바꾸는 건 권장하지 않아요. 체크포인트가 채널 타입을 다르게 인코딩하기 때문에, 기존 스레드의 타입을 바꾸면 상태 재구성이 불완전하거나 잘못될 수 있어요. 채널 정의는 스레드의 수명 동안 안정적으로 유지하세요. 채널 타입을 바꾸기 전에 영향받는 스레드를 새 표현으로 마이그레이션하거나, 버리고 새 스레드를 시작하세요.
DeltaChannel을 지원하지 않는 버전으로의 롤백은 지원되지 않아요.langgraph>=1.2는 델타 채널 체크포인트를 이전 버전이 읽을 수 없는 새 형식으로 써요. 한 스레드가DeltaChannel을 사용한 뒤 LangGraph를 다운그레이드하면, 이전 런타임은 델타 형식을 이해하지 못해 채널 상태를 재구성할 수 없으므로 그 체크포인트들을 읽을 수 없게 돼요. 롤백이 필요하다면 delta-channel-dump 복구 스크립트로 영향받는 스레드를 마이그레이션하거나, 다운그레이드 전에 그것들을 버리세요.
예시 (Examples)
대부분의 사용자는 StateGraph API나 @entrypoint 데코레이터를 통해 Pregel과 상호작용해요. 하지만 Pregel을 직접 다루는 것도 가능해요.
아래는 Pregel API의 감을 잡을 수 있게 해주는 몇 가지 예시예요.
단일 노드 (Single node)
from langgraph.channels import EphemeralValue
from langgraph.pregel import Pregel, NodeBuilder
node1 = (
NodeBuilder().subscribe_only("a")
.do(lambda x: x + x)
.write_to("b")
)
app = Pregel(
nodes={"node1": node1},
channels={
"a": EphemeralValue(str),
"b": EphemeralValue(str),
},
input_channels=["a"],
output_channels=["b"],
)
app.invoke({"a": "foo"})
{'b': 'foofoo'}
여러 노드 (Multiple nodes)
from langgraph.channels import LastValue, EphemeralValue
from langgraph.pregel import Pregel, NodeBuilder
node1 = (
NodeBuilder().subscribe_only("a")
.do(lambda x: x + x)
.write_to("b")
)
node2 = (
NodeBuilder().subscribe_only("b")
.do(lambda x: x + x)
.write_to("c")
)
app = Pregel(
nodes={"node1": node1, "node2": node2},
channels={
"a": EphemeralValue(str),
"b": LastValue(str),
"c": EphemeralValue(str),
},
input_channels=["a"],
output_channels=["b", "c"],
)
app.invoke({"a": "foo"})
{'b': 'foofoo', 'c': 'foofoofoofoo'}
Topic
from langgraph.channels import EphemeralValue, Topic
from langgraph.pregel import Pregel, NodeBuilder
node1 = (
NodeBuilder().subscribe_only("a")
.do(lambda x: x + x)
.write_to("b", "c")
)
node2 = (
NodeBuilder().subscribe_to("b")
.do(lambda x: x["b"] + x["b"])
.write_to("c")
)
app = Pregel(
nodes={"node1": node1, "node2": node2},
channels={
"a": EphemeralValue(str),
"b": EphemeralValue(str),
"c": Topic(str, accumulate=True),
},
input_channels=["a"],
output_channels=["c"],
)
app.invoke({"a": "foo"})
{'c': ['foofoo', 'foofoofoofoo']}
BinaryOperatorAggregate
이 예시는 BinaryOperatorAggregate 채널로 리듀서를 구현하는 방법을 보여줘요.
from langgraph.channels import EphemeralValue, BinaryOperatorAggregate
from langgraph.pregel import Pregel, NodeBuilder
node1 = (
NodeBuilder().subscribe_only("a")
.do(lambda x: x + x)
.write_to("b", "c")
)
node2 = (
NodeBuilder().subscribe_only("b")
.do(lambda x: x + x)
.write_to("c")
)
def reducer(current, update):
if current:
return current + " | " + update
else:
return update
app = Pregel(
nodes={"node1": node1, "node2": node2},
channels={
"a": EphemeralValue(str),
"b": EphemeralValue(str),
"c": BinaryOperatorAggregate(str, operator=reducer),
},
input_channels=["a"],
output_channels=["c"],
)
app.invoke({"a": "foo"})
{ 'c': 'foofoo | foofoofoofoo' }
순환 (Cycle)
이 예시는 체인이 구독하는 채널에 쓰기를 해서 그래프에 순환을 도입하는 방법을 보여줘요. 채널에 None 값이 쓰일 때까지 실행이 계속돼요.
from langgraph.channels import EphemeralValue
from langgraph.pregel import Pregel, NodeBuilder, ChannelWriteEntry
example_node = (
NodeBuilder().subscribe_only("value")
.do(lambda x: x + x if len(x) < 10 else None)
.write_to(ChannelWriteEntry("value", skip_none=True))
)
app = Pregel(
nodes={"example_node": example_node},
channels={
"value": EphemeralValue(str),
},
input_channels=["value"],
output_channels=["value"],
)
app.invoke({"value": "a"})
{'value': 'aaaaaaaaaaaaaaaa'}
고수준 API (High-level API)
LangGraph는 Pregel 애플리케이션을 만드는 두 가지 고수준 API를 제공해요. StateGraph (Graph API)와 Functional API예요.
StateGraph (Graph API)
StateGraph (Graph API)는 Pregel 애플리케이션 생성을 단순화하는 더 높은 수준의 추상화예요. 노드와 엣지로 이루어진 그래프를 정의할 수 있게 해줘요. 그래프를 컴파일하면 StateGraph API가 Pregel 애플리케이션을 자동으로 만들어줘요.
from typing import TypedDict
from langgraph.constants import START
from langgraph.graph import StateGraph
class Essay(TypedDict):
topic: str
content: str | None
score: float | None
def write_essay(essay: Essay):
return {
"content": f"Essay about {essay['topic']}",
}
def score_essay(essay: Essay):
return {
"score": 10
}
builder = StateGraph(Essay)
builder.add_node(write_essay)
builder.add_node(score_essay)
builder.add_edge(START, "write_essay")
builder.add_edge("write_essay", "score_essay")
# Compile the graph.
# This will return a Pregel instance.
graph = builder.compile()
컴파일된 Pregel 인스턴스는 노드와 채널 목록과 연결돼요. 인쇄해서 노드와 채널을 살펴볼 수 있어요.
print(graph.nodes)
대략 이런 결과가 보일 거예요.
{'__start__': <langgraph.pregel.read.PregelNode at 0x7d05e3ba1810>,
'write_essay': <langgraph.pregel.read.PregelNode at 0x7d05e3ba14d0>,
'score_essay': <langgraph.pregel.read.PregelNode at 0x7d05e3ba1710>}
print(graph.channels)
이런 결과가 보일 거예요.
{'topic': <langgraph.channels.last_value.LastValue at 0x7d05e3294d80>,
'content': <langgraph.channels.last_value.LastValue at 0x7d05e3295040>,
'score': <langgraph.channels.last_value.LastValue at 0x7d05e3295980>,
'__start__': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e3297e00>,
'write_essay': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e32960c0>,
'score_essay': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e2d8ab80>,
'branch:__start__:__self__:write_essay': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e32941c0>,
'branch:__start__:__self__:score_essay': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e2d88800>,
'branch:write_essay:__self__:write_essay': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e3295ec0>,
'branch:write_essay:__self__:score_essay': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e2d8ac00>,
'branch:score_essay:__self__:write_essay': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e2d89700>,
'branch:score_essay:__self__:score_essay': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e2d8b400>,
'start:write_essay': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e2d8b280>}
Functional API
Functional API에서는 @entrypoint를 사용해 Pregel 애플리케이션을 만들 수 있어요. entrypoint 데코레이터는 입력을 받아 출력을 반환하는 함수를 정의할 수 있게 해줘요.
from typing import TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.func import entrypoint
class Essay(TypedDict):
topic: str
content: str | None
score: float | None
checkpointer = InMemorySaver()
@entrypoint(checkpointer=checkpointer)
def write_essay(essay: Essay):
return {
"content": f"Essay about {essay['topic']}",
}
print("Nodes: ")
print(write_essay.nodes)
print("Channels: ")
print(write_essay.channels)
Nodes:
{'write_essay': <langgraph.pregel.read.PregelNode object at 0x7d05e2f9aad0>}
Channels:
{'__start__': <langgraph.channels.ephemeral_value.EphemeralValue object at 0x7d05e2c906c0>, '__end__': <langgraph.channels.last_value.LastValue object at 0x7d05e2c90c40>, '__previous__': <langgraph.channels.last_value.LastValue object at 0x7d05e1007280>}
더 알아보기 (Learn more)
- StateGraph (Graph API): 노드·엣지 기반 그래프로 Pregel 애플리케이션을 만드는 방법을 다뤄요.
- Functional API:
@entrypoint로 Pregel 애플리케이션을 만드는 방법을 다뤄요. - delta-channel-dump 복구 스크립트:
DeltaChannel사용 스레드를 롤백 전에 마이그레이션하는 방법을 다뤄요.