LangGraph 런타임
LangGraph 런타임
Pregel은 LangGraph의 런타임을 구현해서, LangGraph 애플리케이션의 실행을 관리해요.
StateGraph를 컴파일하거나 @entrypoint를 만들면 Pregel 인스턴스가 생기는데, 이 인스턴스를 입력과 함께 호출할 수 있어요. 이 가이드에서는 런타임을 높은 수준에서 설명하고, Pregel로 애플리케이션을 직접 구현하는 방법도 함께 다룰게요.
참고:
Pregel런타임의 이름은 Google의 Pregel 알고리즘에서 따왔어요. 이 알고리즘은 그래프를 이용해 대규모 병렬 계산을 효율적으로 처리하는 방법을 설명하죠.
개요
LangGraph에서 Pregel은 actors와 **채널(channel)**을 하나의 애플리케이션으로 결합해요. Actor는 채널에서 데이터를 읽고 채널에 데이터를 씁니다. Pregel은 실행을 여러 단계로 나눠서 진행하는데, 이는 Pregel 알고리즘 / Bulk Synchronous Parallel 모델을 따르죠.
각 단계는 세 가지 단계로 나뉘어요:
- Plan(계획): 이번 단계에서 어떤 actor를 실행할지 정해요. 예를 들어 첫 단계에서는 특별한 입력 채널을 구독하는 actor를 고르고, 이후 단계에서는 이전 단계에서 갱신된 채널을 구독하는 actor를 골라요.
- Execution(실행): 선택한 actor 전부를 병렬로 실행해요. 전부 끝나거나 하나가 실패하거나 타임아웃에 이를 때까지 진행되죠. 이 단계에서 채널 갱신은 다음 단계 전까지 actor에게 보이지 않아요.
- Update(갱신): 이번 단계에서 actor들이 쓴 값으로 채널을 갱신해요.
실행할 actor가 없어지거나 최대 단계 수에 도달할 때까지 이 과정을 반복해요.
Actor
Actor는 PregelNode예요. 채널을 구독하고 데이터를 읽고 씁니다. Pregel 알고리즘의 actor처럼 생각하면 돼요. PregelNode는 LangChain의 Runnable 인터페이스를 구현해요.
채널 (Channel)
채널은 actor(PregelNode) 사이의 통신에 사용돼요. 각 채널은 값 타입(value type), 갱신 타입(update type), 그리고 update 함수를 가져요. update 함수는 일련의 갱신을 받아 저장된 값을 수정하죠. 채널은 한 체인에서 다른 체인으로 데이터를 보내거나, 미래 단계의 자기 자신에게 데이터를 보내는 데 쓸 수 있어요.
LastValue
LastValue는 기본 채널 타입이에요. 마지막으로 쓴 값만 저장하고 이전 값은 덮어써요. 입력·출력 값이나 한 단계에서 다음 단계로 데이터를 넘길 때 사용해요.
from langgraph.channels import LastValue
channel: LastValue[int] = LastValue(int)
Topic
Topic은 설정 가능한 PubSub 채널이에요. actor 사이에 여러 값을 보내거나 단계에 걸친 출력을 누적하는 데 유용하죠. 값을 중복 제거하도록 하거나, 실행 중에 쓴 값을 전부 누적하도록 설정할 수 있어요.
from langgraph.channels import Topic
# 단계를 걸쳐 쓴 모든 값을 누적해요
channel: Topic[str] = Topic(str, accumulate=True)
BinaryOperatorAggregate
BinaryOperatorAggregate는 이항 연산자를 현재 값과 각 새 갱신에 적용해 영속 값을 갱신해요. 단계에 걸친 누적 집계를 계산할 때 쓰죠.
import operator
from langgraph.channels import BinaryOperatorAggregate
# 진행 합계: 각 쓰기가 현재 값에 더해져요
total = BinaryOperatorAggregate(int, operator.add)
DeltaChannel
DeltaChannel은 전체 누적 값 대신 각 단계의 증분 델타만 저장해요. 자주 쓰이면서 시간이 지나 큰 값으로 누적되는 채널에 가장 유용해요. 예를 들어 오래 지속되는 스레드의 대화 메시지 목록이 그렇죠. 델타 저장이 없으면 전체 목록이 모든 체크포인트에 다시 직렬화되지만, DeltaChannel을 쓰면 각 단계에서 새로 쓴 메시지만 저장돼요.
DeltaChannel은 일반 리듀서를 쓰듯 Annotated 타입 어노테이션 안에 사용해요:
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)]
벌크 리듀서 요구사항
DeltaChannel에 넘기는 reducer는 벌크 리듀서예요. 현재 상태와 현재 단계의 모든 쓰기 시퀀스를 한 번에 받죠. StateGraph에서 Annotated와 함께 쓰는 케이별 리듀서처럼 갱신마다 한 번씩 호출되는 게 아니에요.
reducer(reducer(state, [xs]), [ys]) == reducer(state, [xs, ys])
리듀서가 결합 법칙을 지키지 않으면, LangGraph가 쓰기를 단계별로 어떻게 배치하느냐에 따라 재구성된 상태가 달라질 수 있어서 일관성이 깨질 수 있어요.
리듀서를 설계할 때 실질적인 함의는 다음과 같아요:
(state, writes)의 순수 함수로 만들어요. 부수 효과, 무작위성, 벽시계 읽기(uuid.uuid4(),datetime.now())는 값이 재구성될 때마다 실행되어 재생마다 다른 결과를 내요. 이런 것들은 저장된 쓰기에 박히지 않아요.- 들어오는 쓰기에 대한 변경이 영속된다고 가정하지 마세요. 리듀서가 쓰기 객체를 변경했다면(예: ID 없는 항목에 안정적인 ID를 부여), 그 변경은 재구성된 값에만 남아요. 저장된 쓰기는 원래 모양 그대로라 다음 재구성 때 변경 전 입력을 다시 봐요.
- 아이디entical 및 기타 안정적인 메타데이터는 업스트림에서 붙이세요. 다운스트림 코드가 이후 턴에서 항목을 ID로 참조해야 한다면(나중에 수정·삭제하려고), 그 ID는 값이 채널에 쓰이기 전에 할당해야 해요. 리듀서 안에서 말고요.
가장 흔한 두 경우의 벌크 리듀서를 보여드릴게요:
from typing import Any, Sequence
# 리스트: 모든 쓰기를 순서대로 이어붙여요
def list_reducer(state: list[Any], writes: Sequence[list[Any]]) -> list[Any]:
result = list(state)
for write in writes:
result.extend(write)
return result
# 딕트: 모든 쓰기를 병합하고, 키 충돌 시 마지막 쓰기가 이겨요
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(기본값)은 스냅샷을 아예 건너뛰어요. 읽기가 드물거나 스레드가 짧을 때 적합해요.
버전 호환성과 롤백
DeltaChannel을 지원하지 않는 버전으로 롤백하는 것은 지원되지 않아요. langgraph>=1.2는 이전 버전이 읽지 못하는 새 형식으로 델타 채널 체크포인트를 써요. 스레드가 DeltaChannel을 쓴 다음 LangGraph를 다운그레이드하면, 이전 런타임은 델타 형식을 이해하지 못해 채널 상태를 재구성할 수 없어서 그 체크포인트를 읽을 수 없게 돼요. 롤백이 필요하면 delta-channel-dump 복구 스크립트로 영향받는 스레드를 마이그레이션하거나, 다운그레이드 전에 버리세요.
예제
대부분의 사용자는 StateGraph API나 @entrypoint 데코레이터로 Pregel을 다루지만, Pregel을 직접 다룰 수도 있어요. 아래는 Pregel API의 느낌을 잡을 수 있는 몇 가지 예제예요.
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"})
```
```con
{'b': 'foofoo'}
```
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"})
```
```con
{'b': 'foofoo', 'c': 'foofoofoofoo'}
```
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"})
```
```pycon
{'c': ['foofoo', 'foofoofoofoo']}
```
```python
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"})
```
```console
{ 'c': 'foofoo | foofoofoofoo' }
```
```python
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"})
```
```pycon
{'value': 'aaaaaaaaaaaaaaaa'}
```
고수준 API
LangGraph는 Pregel 애플리케이션을 만드는 두 가지 고수준 API를 제공해요: StateGraph (Graph API)와 Functional API.
```python
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")
# 그래프를 컴파일해요.
# 이렇게 하면 Pregel 인스턴스가 반환돼요.
graph = builder.compile()
```
컴파일된 Pregel 인스턴스는 노드와 채널 목록과 연결돼요. `print`로 노드와 채널을 확인할 수 있어요.
```python
print(graph.nodes)
```
이런 식으로 보여요:
```pycon
{'__start__': <langgraph.pregel.read.PregelNode at 0x7d05e3ba1810>,
'write_essay': <langgraph.pregel.read.PregelNode at 0x7d05e3ba14d0>,
'score_essay': <langgraph.pregel.read.PregelNode at 0x7d05e3ba1710>}
```
```python
print(graph.channels)
```
이런 식으로 나와요:
```pycon
{'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>}
```
```python
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)
```
```pycon
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)
- Graph API와 Functional API로 Pregel 애플리케이션을 만드는 고수준 접근법
- Pregel API 레퍼런스