
AI 에이전트를 처음 만들 때는 보통 아래와 같은 구조부터 떠올린다.
사용자
↓
HTTP Request
↓
Spring Controller
↓
LLM 호출
↓
Tool 실행
↓
LLM 재호출
↓
Response
개인 프로젝트나 짧은 질의응답이라면 이 정도로도 충분하다.
문제는 에이전트가 단순한 채팅봇을 넘어가기 시작할 때다.
예를 들어 사용자가 이런 요청을 한다고 해보자.
"프로젝트 전체를 분석하고
문제가 있는 코드를 수정한 뒤
테스트를 실행하고,
실패하면 원인을 분석해서 다시 수정해줘."
내부적으로는 단순한 LLM 한 번 호출이 아니다.
사용자 요청
↓
LLM 추론
↓
파일 검색
↓
Tool 호출
↓
코드 수정
↓
Build
↓
Test
↓
오류 분석
↓
LLM 재추론
↓
다시 Tool 실행
↓
최종 결과
몇 초가 아니라 몇 분, 경우에 따라 그 이상 지속될 수 있다.
여기에 수십 명, 수백 명이 동시에 사용하고, Worker 장애와 재시도, 중복 요청, Tool 취소, 실시간 진행 상황, 감사 로그까지 요구되면 기존의 "요청 하나를 함수 하나가 끝까지 처리한다"는 모델로는 점점 버티기 어려워진다.
그래서 아키텍처의 관점을 바꾼다.
함수 호출 중심
↓
Event + Queue + Worker 중심
이 글에서는 이 구조가 왜 필요한지와 Event Sourcing, PGMQ, Idempotency, Worker Pool, Outbox, WebSocket이 서로 어떤 관계인지 하나의 흐름으로 정리한다.
가장 단순한 Agent 코드는 대략 이런 모습이다.
public AgentResponse execute(AgentRequest request) {
Context context = new Context();
while (true) {
LlmResponse response = llm.call(context);
if (response.hasToolCall()) {
ToolResult result =
tool.execute(response.toolCall());
context.add(result);
continue;
}
return response.finalAnswer();
}
}
논리 자체는 자연스럽다.
LLM
↓
Tool
↓
LLM
↓
Tool
↓
LLM
↓
Final Answer
하지만 이 구조에서는 Agent가 어디까지 진행됐는지를 나타내는 중요한 Context가 JVM 메모리에 있다.
예를 들어 여기까지 실행됐다고 해보자.
User Request
↓
LLM 완료
↓
Tool 요청
↓
Tool 실행 중
↓
💥 장애
프로세스가 죽으면 메모리에 있던 Context 역시 사라진다.
재시작한 서버는 이런 질문에 답할 수 없다.
Tool을 호출했나?
Tool 결과를 받았나?
Iteration은 몇 번째인가?
파일 수정은 이미 끝났나?
다시 실행해도 안전한가?
여기에서 중요한 원칙이 나온다.
Agent 작업의 생명주기를 서버 프로세스의 생명주기와 분리해야 한다.
즉,
Server Process Lifetime
≠
Agent Workflow Lifetime
이어야 한다.
기존 방식에서는 현재 Agent 상태 자체가 중요했다.
Agent State
↓
Memory
이벤트 기반 구조에서는 조금 다르게 생각한다.
지금까지 무슨 일이 발생했는가?
를 기록한다.
예를 들면 다음과 같다.
session.created
turn.started
iteration.started
model.started
model.completed
tool.requested
tool.started
tool.progress
tool.completed
iteration.completed
turn.completed
PostgreSQL의 Event Log에는 이런 식으로 계속 기록될 수 있다.
seq | event
----+-----------------------
1 | session.created
2 | turn.started
3 | iteration.started
4 | model.started
5 | model.completed
6 | tool.requested
7 | tool.started
8 | tool.completed
9 | iteration.completed
서버가 여기서 죽어도 기록은 사라지지 않는다.
재시작한 서버는 Event Log를 읽고
"Tool까지 완료됐고
현재 Iteration 하나가 끝난 상태구나."
처럼 상태를 다시 계산할 수 있다.
이 관계를 보통 다음과 같이 표현한다.
Event Log
↓
Projection
↓
Current State
현재 상태가 이벤트들의 Projection으로 만들어진다면 이 구조는 Event Sourcing에 가까운 설계가 된다.
여기에서 한 가지는 구분해야 한다.
단순히 로그를 많이 저장한다고 모두 Event Sourcing인 것은 아니다.
로그를 참고만 한다
→ 일반적인 Event Logging
상태의 근거를 Event로 삼고
Event를 재생해 상태를 재구성한다
→ Event Sourcing
이 차이가 있다.
처음 접하면 가장 많이 헷갈리는 부분이다.
둘 다 Event가 들어가니 비슷해 보이지만 책임은 전혀 다르다.
질문은 이것이다.
"지금까지 무슨 일이 일어났는가?"
예:
turn.started
model.completed
tool.started
tool.completed
turn.completed
History이자 Audit 기록이다.
장애 분석, Replay, 상태 복구에도 활용한다.
질문은 이것이다.
"앞으로 누가 이 일을 처리해야 하는가?"
예:
LLM_INFERENCE_REQUESTED
TOOL_EXECUTION_REQUESTED
REPORT_GENERATION_REQUESTED
Worker가 가져가 실제 작업을 수행한다.
따라서 전체 흐름은 이런 식으로 이해하는 것이 좋다.
┌──────────── Event Log
│
Command/Event ┤
│
└──────────── Work Queue
│
▼
Worker
모든 Event Log 항목을 무조건 작업 Queue에 집어넣는 것이 아니다.
처리해야 하는 Command/Job과 이미 발생한 사실을 기록하는 Event는 목적이 다르다.
이 구분을 잡아두면 설계가 훨씬 깔끔해진다.
처음에는 PostgreSQL에 다음과 같은 테이블 하나만 만들어도 될 것처럼 보인다.
job_table
id | type | status
-----------------------
1 | chat | READY
2 | tool | READY
3 | rag | READY
그리고 Worker가 가져간다.
SELECT *
FROM job
WHERE status = 'READY'
LIMIT 1;
Worker가 하나일 때는 별문제가 없어 보인다.
그러나 Worker가 여러 개가 되면 이야기가 달라진다.
Worker A ─┐
├─ Job 1 발견
Worker B ─┘
잘못 설계하면 둘 다 같은 Job을 실행할 수 있다.
Worker A → 환자 예약 등록
Worker B → 환자 예약 등록
결과: 예약 2건
물론 PostgreSQL의 Row Lock이나 FOR UPDATE SKIP LOCKED 같은 기능으로 직접 Queue를 만들 수도 있다.
문제는 그 순간부터 우리가 직접 구현해야 할 것이 크게 늘어난다는 것이다.
Claim
Lock
Timeout
Retry
Worker Crash
Visibility
Duplicate
Dead Letter
Acknowledgement
그래서 이미 이런 동작을 Queue로 추상화해 놓은 시스템을 사용하는 것이 편하다.
이 아키텍처에서는 PostgreSQL 내부에서 Queue 기능을 제공하는 PGMQ를 선택한다.
전체 데이터 계층을 단순화하면 다음과 같다.
PostgreSQL
│
├─ Session
├─ Event Log
├─ Idempotency
├─ Outbox
├─ PGMQ
└─ pgvector
별도로 Kafka를 사용한다면 구조는 더 커진다.
Agent Server
├─ PostgreSQL
├─ Kafka
├─ Kafka Consumer
├─ Broker Monitoring
├─ Broker HA
└─ DB ↔ Broker 일관성 처리
반면 초기 Agent 시스템에서는
Spring Boot
↓
PostgreSQL
├─ 일반 데이터
├─ Event
├─ Queue
└─ Vector
↓
Worker
정도로 시작할 수 있다.
이것은 Kafka가 나쁘다는 뜻이 아니다.
규모가 커지고 대규모 스트리밍, 높은 메시지 처리량, 다양한 Consumer Group 같은 요구가 중요해지면 Kafka가 더 적합할 수도 있다.
다만 AI Agent 작업의 경우 실제 처리 시간이
Queue Read 5ms
DB Search 100ms
Tool 1s
LLM 5s
같다면 당장 가장 큰 병목은 Queue가 아니다.
LLM
외부 API
RAG
Tool
Process
가 훨씬 오래 걸릴 가능성이 높다.
따라서 초기에는 운영 복잡도를 줄이는 선택도 충분히 가치가 있다.
그리고 Queue 인터페이스를 분리해두면 나중에 구현체를 바꾸기 쉽다.
public interface JobQueue {
void send(Job job);
Optional<Job> receive();
void ack(String id);
void extendVisibility(String id);
}
현재는
PgmqJobQueue
를 사용하다가 필요하면
SqsJobQueue
KafkaJobQueue
같은 방식으로 교체할 수 있다.
핵심은 Application이 특정 Queue 제품에 너무 깊게 묶이지 않는 것이다.
Queue에서 메시지를 읽었다고 바로 삭제하면 안 된다.
다음 구현은 위험하다.
Queue
↓
Message Read
↓
Delete
↓
Worker 실행
↓
💥 Worker 장애
Message는 이미 사라졌다.
복구할 방법도 없다.
따라서 정상적인 흐름은 다음에 가깝다.
Message Read
↓
Claim / Invisible
↓
Worker 처리
↓
결과 저장
↓
Event 저장
↓
ACK / Delete
예를 들어 Visibility Timeout이 30초라고 하자.
0초
Worker A가 Job 획득
→ 다른 Worker에게는 숨김
정상적으로 처리되면
20초
처리 완료
→ ACK
→ Queue에서 제거
Worker가 죽으면?
0초 Worker A Claim
15초 Worker A 💥
30초 Visibility Timeout 만료
31초 Job 다시 Visible
32초 Worker B Claim
여기서 중요한 사고방식이 하나 나온다.
Worker는 죽어도 된다. Job은 잃으면 안 된다.
Worker 자체를 절대 죽지 않는 존재로 만들려고 하기보다, 죽더라도 다른 Worker가 이어받을 수 있도록 만드는 것이 Fault Tolerance다.
LLM이나 Tool 작업이 Visibility Timeout보다 길어질 수도 있다.
예:
Visibility Timeout = 30초
Tool Execution = 2분
그대로 두면 30초 뒤에 Queue는 작업이 버려졌다고 판단할 수 있다.
Worker A
↓
Job 실행 중
↓
30초
↓
Job 재노출
↓
Worker B가 동일 Job Claim
결과적으로 A와 B가 동시에 같은 작업을 실행한다.
따라서 실행 중인 Worker는 주기적으로
"아직 내가 처리하고 있다."
라는 신호를 보내 Visibility를 연장한다.
0초 Claim
20초 Visibility 연장
40초 Visibility 연장
60초 Visibility 연장
80초 완료
81초 ACK
이런 Heartbeat는 단순 상태 표시가 아니라 중복 실행을 줄이는 작업 소유권 관리 수단이 된다.
Queue 시스템을 사용할 때는 보통 다음 가정을 하는 것이 안전하다.
메시지는 한 번 이상 전달될 수 있다.
즉 at-least-once다.
왜 중복될까?
Worker 실행
↓
DB 저장 성공
↓
ACK 전송
↓
💥 ACK 유실
Queue 입장에서는 작업이 끝났다는 것을 모른다.
그러면 같은 Message를 다시 전달할 수 있다.
Message A
↓
Worker 실행
↓
완료
↓
ACK 유실
↓
Message A 재전달
따라서 시스템은
"중복 메시지가 발생하지 않도록"
만드는 데만 의존하면 안 된다.
대신
"중복 메시지가 와도 결과가 망가지지 않도록"
만들어야 한다.
그게 Idempotency, 멱등성이다.
예를 들어 사용자가 다음 작업을 요청했다고 하자.
transactionKey = TX100
환자 예약 등록
첫 요청이 실제로 성공했지만 응답이 네트워크에서 유실됐다.
Client는 성공 여부를 모른다.
"응답이 없네."
그리고 재시도한다.
TX100
환자 예약 등록
서버가 아무 확인 없이 실행하면 예약이 두 건 생성된다.
그래서 처리 전에 멱등성 테이블을 본다.
예:
CREATE TABLE idempotency_request (
transaction_key VARCHAR(100) PRIMARY KEY,
request_hash VARCHAR(64) NOT NULL,
status VARCHAR(20) NOT NULL,
result JSONB,
created_at TIMESTAMP NOT NULL
);
처음 TX100이 들어오면 Insert를 시도한다.
TX100 없음
↓
INSERT 성공
↓
내가 최초 처리자
두 번째 요청이 들어오면 이미 존재한다.
상태에 따라 다르게 처리한다.
┌─ 없음
│ → 신규 실행
│
transactionKey ─────┼─ PROCESSING
│ → 중복 실행하지 않음
│
├─ COMPLETED
│ → 기존 결과 반환
│
└─ Key는 같은데 Payload 다름
→ 요청 오류
멱등성은 Key만 비교하면 부족하다
예를 들어 첫 번째 요청이 이렇다.
{
"transactionKey": "ABC",
"amount": 10000
}
그런데 나중에:
{
"transactionKey": "ABC",
"amount": 50000
}
가 들어왔다.
Key는 같다.
하지만 같은 요청의 Retry가 아니다.
그래서 Payload의 Hash를 같이 저장한다.
transactionKey = ABC
requestHash = SHA-256(payload)
검사는 이렇게 된다.
Key 없음
→ 신규 요청
Key 같음 + Hash 같음
→ 같은 요청의 Retry
Key 같음 + Hash 다름
→ Idempotency Key 오사용
API라면 이런 경우 409 Conflict 같은 방식으로 거절할 수 있다.
왜 처리 결과도 같이 저장하는가
멱등성이
"이미 했습니다."
에서 끝나면 Client는 곤란하다.
Client가 원하는 것은 원래 작업의 결과다.
따라서 완료 후 결과를 저장한다.
TX100
status = COMPLETED
result =
{
"reservationId": 152,
"message": "예약되었습니다."
}
같은 요청이 재시도되면 실제 Business Logic을 다시 실행하지 않고
기존 result 반환
으로 끝낸다.
Event Log, Idempotency, Queue를 구분하면 전체 구조가 보인다
이 세 가지는 역할을 한 줄씩 정리하면 가장 쉽다.
Queue
→ 앞으로 누가 일을 할 것인가?
Idempotency
→ 이 일을 이미 했는가?
Event Log
→ 이 일이 어떻게 진행됐는가?
예를 들어 하나의 요청이 다음과 같이 처리됐다고 하자.
Idempotency:
TX100 → COMPLETED
Event Log:
101 request.accepted
102 turn.started
103 model.started
104 model.completed
105 tool.started
106 tool.completed
107 turn.completed
Queue는 처리 중에 다음 Job들을 Worker에게 전달했을 것이다.
LLM_REQUESTED
TOOL_REQUESTED
LLM_REQUESTED
셋은 비슷해 보이지만 각자 완전히 다른 문제를 해결한다.
Outbox가 필요한 이유: DB 저장과 Event 전송의 틈
여기서 분산 시스템의 고전적인 문제가 하나 더 생긴다.
Tool 실행이 끝났다.
Worker가 DB에는 완료 상태를 저장했다.
UPDATE tool_execution
SET status = 'COMPLETED';
그리고 다음 Event를 Queue에 보내려고 한다.
TOOL_COMPLETED
그 사이 서버가 죽는다.
결과:
DB
Tool = COMPLETED
그런데
Queue
TOOL_COMPLETED 없음
시스템 흐름이 끊어진다.
반대로 Queue부터 보내고 DB 저장이 실패해도 문제가 생긴다.
그래서 Transactional Outbox Pattern을 사용할 수 있다.
BEGIN;
UPDATE tool_execution
SET status = 'COMPLETED';
INSERT INTO outbox_event (...);
COMMIT;
Business 상태와 "이 Event를 나중에 발행해야 한다"는 사실을 같은 DB Transaction으로 저장한다.
그다음 별도 Publisher가 Outbox를 읽는다.
Outbox Table
↓
Publisher
↓
PGMQ / Stream
이렇게 하면
DB에는 성공했는데 Event는 사라짐
같은 Dual Write 문제를 줄일 수 있다.
그리고 여기에도 재전송 가능성이 있기 때문에 Consumer 측 멱등성이 다시 필요하다.
결국 다음 요소들이 서로 맞물린다.
At-Least-Once Delivery
+
Idempotency
+
Transactional Outbox
Event Envelope는 Event의 공통 봉투다
Event가 몇 종류 안 될 때는 다음처럼 만들어도 된다.
{
"type": "tool.completed"
}
하지만 실제 Agent에서는 수많은 Session과 Turn, Tool 실행이 동시에 움직인다.
Event만 보고도 최소한 이런 것을 알아야 한다.
어떤 Session인가?
어느 Turn인가?
몇 번째 Iteration인가?
무엇 때문에 발생했는가?
요청인가 결과인가?
현재 상태는 무엇인가?
순서는 몇 번째인가?
그래서 모든 Event에 공통 Envelope를 둔다.
{
"eventKey": "evt-123",
"transactionKey": "tx-100",
"parentEventKey": "evt-120",
"sessionKey": "session-1",
"runKey": "run-1",
"turnKey": "turn-3",
"iterationKey": "iteration-2",
"stepKey": "step-4",
"toolCallKey": "tool-call-10",
"kind": "progress",
"scope": "tool",
"state": "running",
"sequence": 13,
"createdAt": "2026-09-09T01:00:00Z",
"payload": {
"tool": "patientSearch"
}
}
Spring/Java에서는 이런 형태로 옮길 수 있다.
public record EventEnvelope<T>(
UUID eventKey,
UUID transactionKey,
UUID parentEventKey,
UUID sessionKey,
UUID runKey,
UUID turnKey,
UUID iterationKey,
UUID stepKey,
UUID toolCallKey,
EventKind kind,
EventScope scope,
EventState state,
long sequence,
Instant createdAt,
T payload
) {}
eventKey와 transactionKey는 다르다
둘 다 ID라 헷갈리지만 목적이 완전히 다르다.
eventKey
Event 하나의 고유 식별자
예:
evt-001
evt-002
evt-003
Event Store의 PK 같은 역할을 할 수 있다.
transactionKey
논리적으로 같은 요청인지 판단하는 식별자
예:
Client Request
transactionKey = TX100
네트워크 Retry로 같은 요청이 다시 들어와도
TX100
을 유지한다.
따라서 멱등성 판단에 사용할 수 있다.
Parent Event와 Correlation 정보가 필요한 이유
Agent 실행은 단순한 직선이 아니라 Tree나 Graph 형태가 되기 쉽다.
Turn
└─ Iteration
├─ LLM
├─ Search Tool
├─ DB Tool
└─ File Tool
예를 들어:
E1 turn.started
└─ E2 model.completed
└─ E3 tool.requested
└─ E4 tool.completed
각 Event가
parentEventKey
를 가지고 있으면 인과관계를 복원할 수 있다.
E1
└─ E2
└─ E3
└─ E4
그러면 장애 분석할 때도 다음 질문에 답할 수 있다.
이 Tool 실행은 어느 LLM 응답 때문에 발생했나?
이 오류는 어느 Iteration에서 발생했나?
이 결과는 어느 Turn의 결과인가?
Distributed Tracing의 trace/span 관계와도 비슷한 부분이 있다.
Session → Turn → Iteration → Step → ToolCall 계층
Agent의 실행 문맥은 하나의 ID만으로 표현하기 어렵다.
다음처럼 계층적으로 생각할 수 있다.
Session
└─ Run
└─ Turn
└─ Iteration
└─ Step
└─ ToolCall
대략적인 의미는 다음과 같다.
개념 의미 Session 전체 대화 또는 작업 세션 Run 하나의 Agent 실행 단위 Turn 사용자의 한 번의 입력과 그 처리 Iteration LLM 판단 → Tool 실행 → 재판단의 한 사이클 Step Iteration 안의 세부 단계 ToolCall 구체적인 Tool 호출 하나
예를 들어:
Session S1
Turn T3
│
├─ Iteration I1
│ ├─ LLM
│ └─ HospitalSearch Tool
│
└─ Iteration I2
├─ LLM
└─ Final Answer
Tool을 병렬 실행한다면 더 중요해진다.
Iteration I3
├─ ToolCall TC1
├─ ToolCall TC2
└─ ToolCall TC3
완료 순서는:
TC2
TC1
TC3
일 수 있다.
따라서 Event 자체에 자신이 어디에 속하는지를 설명하는 Context가 있어야 한다.
Kind / State / Scope로 Event 의미를 분리한다
Event 이름 하나에 모든 의미를 욱여넣는 대신 몇 개의 축으로 나눌 수 있다.
Kind
Event가 어떤 역할인가?
request
response
progress
notice
heartbeat
State
현재 상태는?
queued
accepted
claimed
running
waiting
completed
failed
cancelled
timed_out
Scope
어느 영역의 Event인가?
session
run
turn
iteration
step
tool
process
event
예를 들어:
{
"kind": "progress",
"scope": "tool",
"state": "running"
}
이라면
Tool 실행이 진행 중이라는 Event
라는 의미를 얻는다.
순서를 믿지 말고 Sequence와 인과관계를 기록한다
여러 Worker가 동시에 실행되는 Event Driven System에서는 이것을 가정하면 안 된다.
먼저 보낸 Event가
반드시 먼저 처리된다.
예를 들어:
A
B
C
순서로 보냈는데 처리 시간이 다르면:
B → 5초
C → 1초
A → 3초
완료 순서는:
C
A
B
가 된다.
즉,
발행 순서
≠
도착 순서
≠
실행 순서
≠
완료 순서
가 될 수 있다.
따라서 Stream을 다시 조립해야 하는 데이터에는 Sequence를 둔다.
tool.progress sequence=1
tool.progress sequence=2
tool.progress sequence=3
Client가 실제로
2
1
3
순서로 받더라도 재정렬할 수 있다.
또:
1
2
4
가 들어오면
sequence 3 누락
도 감지할 수 있다.
다만 Sequence 하나로 시스템 전체의 모든 이벤트 순서를 보장하려고 하면 복잡해진다.
실무에서는 보통 Session, Turn, ToolCall처럼 필요한 범위 안에서 순서를 관리하는 것이 현실적이다.
Event Family가 필요한 이유
Agent의 Event Type은 시간이 지나면 많아진다.
예를 들면:
Session
Turn
Iteration
Model
Tool
Process
State
Approval
Artifact
Guardrail
등이다.
Session:
session.created
session.cancelled
session.resumed
session.deleted
Turn:
turn.started
turn.progress
turn.completed
turn.failed
turn.cancelled
Iteration:
iteration.started
iteration.progress
iteration.completed
iteration.failed
Tool:
tool.requested
tool.started
tool.progress
tool.completed
tool.failed
tool.cancelled
tool.timed_out
Process:
process.started
process.stdout
process.stderr
process.exited
process.cancelled
Approval:
approval.requested
approval.granted
approval.rejected
approval.expired
Artifact:
artifact.file
artifact.diff
artifact.patch
artifact.log
artifact.screenshot
이렇게 Family별로 분리하면 Event Model의 역할이 명확해진다.
Java의 sealed interface를 활용하면 다음처럼 표현할 수도 있다.
public sealed interface AgentEvent
permits SessionEvent,
TurnEvent,
ToolEvent,
ProcessEvent,
ApprovalEvent {
}
그리고:
public sealed interface ToolEvent
extends AgentEvent
permits ToolStarted,
ToolProgress,
ToolCompleted,
ToolFailed {
}
처럼 세분화할 수 있다.
Queue나 Event Router는 공통 AgentEvent를 다루면서도 실제 Payload의 Type Safety를 유지할 수 있다.
Tool과 Process는 구분하는 편이 좋다
예를 들어:
RunTestTool
이라는 Tool이 있다고 하자.
논리적인 Tool 호출이다.
그 Tool 내부에서 실제로:
mvn test
라는 OS Process가 실행될 수 있다.
따라서:
Tool
└─ Process
관계다.
Tool Event:
tool.started
tool.progress
tool.completed
tool.failed
Process Event:
process.stdout
process.stderr
process.exit
process.timeout
process.cancel
으로 분리할 수 있다.
예:
tool.started
↓
process.started
↓
process.stdout
↓
process.stdout
↓
process.stderr
↓
process.exit
↓
tool.failed
이렇게 만들어두면 Tool 내부 구현을:
Java
Python
Node
Shell
중 무엇으로 바꾸더라도 상위 Agent Workflow는 크게 영향을 받지 않는다.
LLM만 Streaming하는 것이 아니다
AI 서비스에서 Streaming이라고 하면 보통 Token Streaming부터 떠올린다.
안
안녕
안녕하
안녕하세요
하지만 Agent에서는 Tool의 진행 상태도 Streaming하는 것이 좋다.
예를 들어 문서 500개를 분석한다면:
tool.started
"문서 목록을 읽고 있습니다."
tool.progress
100 / 500
tool.progress
250 / 500
tool.progress
420 / 500
tool.completed
500 / 500
Payload도 이런 형태가 될 수 있다.
{
"type": "tool.progress",
"toolCallKey": "TC200",
"sequence": 3,
"payload": {
"message": "관련도가 높은 문서 15건을 추렸습니다.",
"progress": 0.65
}
}
사용자 입장에서는 30초 동안 화면이 멈춰 있는 것과:
문서를 검색하고 있습니다...
152개의 문서를 찾았습니다.
관련도 계산 중입니다...
상위 15개를 분석하고 있습니다...
가 표시되는 것은 경험 차이가 크다.
그리고 이 Progress Event는 단순 UI 장식이 아니다.
운영 측에서도:
작업이 살아 있는가?
어디에서 오래 걸리는가?
어느 Tool이 병목인가?
를 파악하는 Observability 정보가 된다.
Cancel은 화면에서 멈추는 것으로 끝나면 안 된다
사용자가 Agent에:
"그만해."
라고 요청했다.
그런데 화면만 멈추고 뒤에서:
Docker Build
Git Clone
500개 PDF 분석
mvn test
가 계속 실행된다면 제대로 취소된 것이 아니다.
취소는 전체 흐름을 타고 내려가야 한다.
Client
↓
turn.cancel.requested
↓
Agent
↓
Worker
↓
Tool
↓
External Process
결과적으로:
tool.cancelled
↓
turn.cancelled
같은 Event가 발생한다.
이를 위해 작업 자체도 취소 가능한 구조여야 한다.
이런 코드는 취소하기 어렵다.
processEverything();
내부에서 10분 동안 제어권이 돌아오지 않기 때문이다.
반면 다음과 같이 작업을 작은 단위로 처리하면 취소에 반응하기 쉬워진다.
while (hasNext()) {
cancellationToken.throwIfCancelled();
processNextChunk();
}
즉 Event Driven Architecture를 사용한다고 자동으로 Non-blocking이 되는 것이 아니다.
각 작업이 실제로 제어권을 반환할 수 있는 형태인지까지 설계해야 한다.
Worker는 무엇을 해야 하는가
Worker는 기본적으로 Queue를 소비하는 실행자다.
Polling
↓
Claim
↓
Handler 선택
↓
실제 처리
↓
Result/Event 저장
↓
ACK
개념적으로는:
while (running) {
Optional<Job> job = queue.receive();
if (job.isEmpty()) {
waitForNextPoll();
continue;
}
try {
handler.handle(job.get());
queue.ack(job.get().id());
} catch (RetryableException e) {
// ACK하지 않음
// visibility 만료 후 재시도
}
}
여기서 Worker가 직접 모든 종류의 Agent 업무를 알게 만들면 또 하나의 거대한 클래스가 된다.
그래서 Event Type에 따라 Handler를 Routing한다.
LLM_REQUESTED
↓
LlmHandler
TOOL_REQUESTED
↓
ToolHandler
TURN_STARTED
↓
TurnHandler
Spring에서는 Strategy Registry 형태로 가져갈 수도 있다.
Map<EventType, EventHandler> handlers;
모르는 Event가 들어오면 실행하지 않는다.
특히 Tool은 미리 허용한 명령만 실행해야 한다.
Map<String, ToolHandler> tools = Map.of(
"patient.search", patientSearchHandler,
"fhir.read", fhirReadHandler,
"report.generate", reportGenerateHandler
);
Queue에서 받은 임의 문자열을 그대로 Shell에 넘기는 방식은 피해야 한다.
Agent에 Tool 실행 권한이 커질수록 Allow List, Approval, Audit이 중요해진다.
Worker Pool은 왜 필요한가
작업이 들어올 때마다 Thread를 새로 만드는 것은 비용이 있다.
그래서 일정 수의 Worker를 유지한다.
Worker Pool
Worker 1 Idle
Worker 2 Busy
Worker 3 Busy
Worker 4 Idle
Worker 5 Idle
Job이 오면 Idle Worker가 처리한다.
Java라면 익숙한 형태다.
ExecutorService executor =
Executors.newFixedThreadPool(20);
조금 더 세밀하게는:
ThreadPoolExecutor executor =
new ThreadPoolExecutor(
10,
50,
60,
TimeUnit.SECONDS,
new LinkedBlockingQueue<>(1000)
);
같은 식으로 구성할 수 있다.
다만 Agent 서버에서 중요한 것은 API 자체가 아니라 Pool 정책이다.
Core Worker 수
Maximum Worker 수
Queue 크기
작업 Timeout
Retry
Cancellation
Idle Worker 정리
Rejected Policy
등을 정해야 한다.
"모든 작업을 Worker Thread로 보내자"도 정답은 아니다
여기에서 Java/Spring 개발자는 CPU Bound와 I/O Bound를 구분할 필요가 있다.
CPU 연산이 많은 작업:
대규모 Parsing
Compression
Graph 계산
Embedding 후처리
은 Worker Thread를 적극적으로 활용하는 것이 의미가 있다.
반면:
LLM HTTP 요청
외부 API 호출
DB Query
Object Storage Download
같은 작업은 대부분 기다리는 시간이 길다.
따라서:
CPU Bound
→ 적절한 Worker / Thread Pool
I/O Bound
→ Async / Non-blocking 또는 충분히 관리된 Blocking I/O
로 구분하는 편이 좋다.
Spring이라면 환경에 따라:
Spring MVC + Virtual Threads
CompletableFuture
WebClient / WebFlux
전용 Executor
등을 선택할 수 있다.
Event Driven = 무조건 WebFlux도 아니고, Worker = 무조건 Thread 하나도 아니다.
논리적인 작업 분리와 실제 실행 모델은 별도로 판단해야 한다.
Worker를 너무 잘게 쪼개는 것도 답은 아니다
여기서 균형도 필요하다.
파일 하나 읽기
→ Event
문자열 하나 변환
→ Event
DB Row 하나 읽기
→ Event
처럼 지나치게 잘게 나누면 오히려:
Queue I/O
직렬화
Event 저장
Transaction
Tracing
Context 전달
비용만 늘어난다.
좋은 기준은 보통:
실패·재시도·취소·확장 단위가 달라지는 곳에서 경계를 나누는 것
이다.
예를 들어:
LLM 호출
외부 Tool 실행
대용량 분석
Human Approval
장시간 Process
는 분리 가치가 크다.
반대로 단순 Java 메서드 몇 줄까지 모두 Event로 바꾸는 것은 오히려 과설계가 될 수 있다.
CPS는 이 구조를 이해하기 위한 사고방식이다
강의에서는 CPS, Continuation-Passing Style이라는 표현이 등장한다.
일반적인 함수 호출은:
Result result = execute(request);
next(result);
처럼 execute()가 끝날 때까지 기다린다.
Continuation을 사용하면 개념적으로:
execute(
request,
result -> next(result)
);
처럼 "끝난 뒤 무엇을 할지"를 전달한다.
Event Driven 구조에서는 이 Continuation을 직접 함수로 전달하는 대신 다음 Event로 표현한다고 이해할 수 있다.
LLM_REQUESTED
↓
LLM Worker
↓
LLM_COMPLETED
↓
Agent Worker
↓
TOOL_REQUESTED
즉:
A Event
↓
B Event
↓
C Event
↓
D Event
가 하나의 Continuation Chain처럼 동작한다.
엄밀한 프로그래밍 언어 이론의 CPS 구현과 완전히 같은 의미로 볼 필요는 없고, 여기서는 긴 호출 스택에 의존하지 않고 "다음에 해야 할 일"을 명시적으로 이어나간다는 관점으로 이해하는 것이 편하다.
Client와 Server의 계약도 Endpoint 중심에서 Event 중심으로 바뀐다
REST에서는 보통 API가 계약이다.
POST /sessions
POST /turns
POST /tools
GET /sessions/{id}
Event Driven 구조에서는 중요한 계약이:
Event Type
Event Envelope
Payload Schema
로 이동한다.
예:
{
"type": "tool.execution.requested",
"payload": {
"toolName": "build",
"project": "agent-server"
}
}
Tool Worker는:
tool.execution.requested
라는 Event를 처리할 수 있어야 한다.
완료하면:
tool.execution.completed
를 발행한다.
즉 클라이언트와 서버, Worker 사이에서
누가 어떤 Event를 발행하는가?
누가 어떤 Event를 소비하는가?
가 중요한 계약이 된다.
그래서 Event Schema Versioning도 결국 중요한 문제가 된다.
실시간 통신은 왜 필요한가
Agent 요청은 동기 HTTP 요청처럼 즉시 결과가 나오지 않는다.
Client가 요청한다.
Client
↓
REQUEST
↓
Server
서버는 작업을 Queue에 넣고 다른 일을 할 수 있다.
몇 초 뒤:
Server
↓
tool.progress
↓
Client
다시 몇 초 뒤:
Server
↓
model.progress
↓
Client
마지막으로:
Server
↓
turn.completed
↓
Client
가 된다.
따라서 Client가 Server로부터 비동기 이벤트를 지속적으로 받을 수 있는 채널이 필요하다.
SSE도 가능하지만 WebSocket은 양방향 Agent에 잘 맞는다
여기서 흔히:
"WebSocket만 가능하고 SSE는 안 된다."
라고 단정하면 조금 과하다.
다음 구조도 충분히 가능하다.
Client → POST /turn
Client ← SSE /sessions/{id}/events
즉 요청은 일반 HTTP로 보내고, Server Event는 SSE로 받을 수 있다.
하지만 Agent가 다음처럼 양방향 제어를 많이 필요로 한다면 WebSocket이 자연스럽다.
Client → input
Client → cancel
Client → approval
Client → confirmation
Server → progress
Server → token
Server → tool result
Server → completed
즉:
Client ⇄ WebSocket ⇄ Agent Server
형태가 된다.
기술 선택은 "Agent니까 무조건 WebSocket"이 아니라 양방향성과 연결 수, 운영 환경을 보고 판단하면 된다.
Session을 Channel로 생각하면 Routing이 단순해진다
예를 들어 Session ID가:
session-100
이라면 이를 Channel ID로 사용할 수 있다.
channel = session-100
Client는 해당 Channel을 구독한다.
SUBSCRIBE session-100
Agent Worker가 Session 100 관련 Event를 만든다.
{
"channelId": "session-100",
"type": "tool.progress"
}
Channel Router는:
Event
↓
channelId 확인
↓
session-100 Subscriber 검색
↓
WebSocket Push
한다.
물리적으로 Session마다 Queue를 하나 만드는 것이 아니라 Event의 Routing 속성으로 Channel을 사용할 수 있다.
Central Event Stream
E1 channel=A
E2 channel=B
E3 channel=A
E4 channel=C
Router가 적절한 Client로 보내면 된다.
서버가 여러 대가 되면 중앙 Stream이 필요해진다
Agent Server를 Scale Out했다고 하자.
Browser A → Server 1
Browser B → Server 2
Worker → Server 3
Worker가 Server 3에서 Event를 만들었는데 Browser A의 WebSocket은 Server 1에 연결되어 있다.
Server 3 메모리에만 Event가 있으면 Server 1은 알 수 없다.
따라서 서버 간 공통 Event 전달 계층이 필요해진다.
Worker
↓
Central Stream / Broker
├─ Server 1 → Browser A
├─ Server 2 → Browser B
└─ Server 3 → Browser C
규모나 요구사항에 따라:
Redis Pub/Sub
Redis Streams
NATS
Kafka
PostgreSQL LISTEN/NOTIFY
같은 선택지가 있을 수 있다.
여기서 PGMQ와 실시간 Stream을 구분하는 것이 좋다.
PGMQ
→ 누가 Job 하나를 가져가 처리할 것인가
Pub/Sub Stream
→ 같은 Event를 여러 연결에 어떻게 전달할 것인가
Queue와 Broadcast는 목적이 다르다.
WebSocket 연결이 너무 많아지면 Stream Server를 Shard한다
직원이 수만 명이고 각 Client가 WebSocket을 유지한다고 생각해보자.
Agent Server
├─ Client 1
├─ Client 2
├─ Client 3
├─ ...
└─ Client 10000
Agent Workflow 처리뿐 아니라 WebSocket Push도 상당한 부하가 된다.
이때 Agent 업무와 Socket 전달 업무를 분리할 수 있다.
Agent Gateway
│
Event Stream
│
┌────────────────┼────────────────┐
▼ ▼ ▼
WS Shard 1 WS Shard 2 WS Shard 3
│ │ │
Clients Clients Clients
Shard는 Agent 비즈니스 로직을 알 필요가 없다.
역할은 단순하다.
Event 받기
channelId 확인
Subscriber 찾기
WebSocket 전송
따라서 상대적으로 Scale Out하기 쉽다.
다만 이런 Sharding은 처음부터 만들 필요는 없다.
초기에는:
Agent Server
+
WebSocket
으로 시작하고 실제 연결 수와 송신 부하가 문제가 될 때 분리하는 것이 현실적이다.
큰 파일은 Event 안에 넣지 않는다
Event 기반 시스템을 처음 만들면 모든 것을 Payload에 넣고 싶어진다.
예:
{
"type": "document.analysis.requested",
"file": "...10MB Base64..."
}
이 방식은 좋지 않다.
Event가 Queue를 이동할 때마다:
Serialization
DB Write
DB Read
Network Transfer
Deserialization
이 반복된다.
큰 바이너리라면 더 심하다.
보통은:
Client
↓
File Upload
↓
S3 / MinIO
↓
Object Key
↓
Event
처럼 만든다.
Event에는:
{
"type": "document.analysis.requested",
"payload": {
"fileId": "FILE-100",
"objectKey": "sessions/S100/input/report.pdf"
}
}
정도만 전달한다.
Event는 작고, 이동하기 쉽고, 재처리하기 쉬운 메시지로 유지하는 편이 좋다.
Approval과 Guardrail도 Event Flow에 포함시킬 수 있다
기업 Agent가 모든 일을 자동 실행하는 것은 위험하다.
예:
메일 전송
운영 DB 수정
파일 삭제
환자 정보 변경
배포 실행
같은 작업은 사람 승인이 필요할 수 있다.
흐름을 Event로 표현하면:
tool.requested
↓
approval.requested
↓
사용자 확인
↓
approval.granted
↓
tool.started
승인을 거절하면:
approval.rejected
↓
tool.cancelled
로 끝낼 수 있다.
Guardrail도 같은 방식으로 넣을 수 있다.
tool.requested
↓
guardrail.requested
↓
guardrail.passed
↓
tool.started
또는:
guardrail.blocked
기업 환경에서 Event Log를 남기는 이유는 단순 장애 복구에만 있지 않다.
누가 요청했는가?
어떤 Tool을 실행했는가?
누가 승인했는가?
어떤 Guardrail이 적용됐는가?
어떤 파일이 수정됐는가?
를 추적하는 Security Audit에도 사용할 수 있다.
Event Replay와 Projection
Event를 모두 기록했다면 특정 Session을 다시 추적할 수 있다.
10:00:00 session.created
10:00:03 turn.started
10:00:05 model.started
10:00:08 model.completed
10:00:09 tool.started
10:00:15 tool.failed
이것만으로:
어느 단계에서 실패했는지
확인할 수 있다.
더 나아가 Event를 재생하여 현재 상태를 만든다.
Event Store
↓
Replay
↓
Projection
예:
SessionProjection
status = RUNNING
currentTurn = T20
currentTool = TEST
lastSequence = 152
Event가 너무 많아지면 처음부터 전부 재생하는 비용이 커질 수 있다.
그래서:
Snapshot
Checkpoint
Compaction
같은 전략도 등장한다.
예:
Event 1 ~ 10000
↓
Checkpoint 저장
이후에는
Checkpoint
+
Event 10001 이후
만 읽는다.
Agent의 최종 Life Cycle을 하나로 연결해보자
사용자가 다음 요청을 보낸다고 하자.
"프로젝트 테스트를 돌리고 실패하면 문제 코드까지 수정해줘."
전체 흐름은 이런 식이 된다.
Client
│
│ USER_MESSAGE
▼
Agent Gateway
│
├─ Idempotency 확인
├─ Event Log 저장
└─ Job Queue 등록
│
▼
PGMQ
│
Claim
▼
Turn Worker
│
turn.started
│
▼
LLM Worker
│
model.started
│
▼
LLM
│
│ "mvn test 필요"
▼
tool.requested
│
▼
Tool Worker
│
tool.started
│
▼
mvn test 실행
│
┌─────┴─────────────┐
│ │
process.stdout process.stderr
│ │
└─────┬─────────────┘
│
tool.failed
│
▼
LLM Worker
│
오류 원인 분석
│
▼
file.patch.requested
│
▼
File Worker
│
file.patch.completed
│
▼
Test 다시 실행
│
▼
tool.completed
│
▼
turn.completed
│
├─ 결과 저장
├─ Idempotency COMPLETED
└─ Queue ACK
│
▼
Event Stream
│
session Channel
│
▼
Client
이 흐름에서 서버가 중간에 죽어도 Event Log와 Queue가 살아 있다면 전체 작업을 잃지 않을 가능성이 높아진다.
같은 요청이 다시 들어와도 Idempotency가 중복 실행을 막는다.
Client는 WebSocket을 통해 진행 상태를 계속 받을 수 있다.
사용자가 Cancel하면 Tool이나 Process까지 중단할 수 있다.
Spring Boot 기준으로 나누면
이 구조를 Spring 프로젝트로 옮기면 대략 다음 정도로 모듈을 분리할 수 있다.
agent-server
│
├─ protocol
│ ├─ EventEnvelope
│ ├─ EventType
│ ├─ EventKind
│ ├─ EventState
│ └─ EventScope
│
├─ event
│ ├─ EventStore
│ ├─ EventPublisher
│ ├─ EventRouter
│ └─ ProjectionService
│
├─ queue
│ ├─ JobQueue
│ ├─ PgmqJobQueue
│ └─ QueueConsumer
│
├─ worker
│ ├─ TurnWorker
│ ├─ LlmWorker
│ ├─ ToolWorker
│ └─ ProcessWorker
│
├─ tool
│ ├─ ToolRegistry
│ ├─ ToolHandler
│ └─ ProcessManager
│
├─ idempotency
│ └─ IdempotencyService
│
├─ outbox
│ ├─ OutboxRepository
│ └─ OutboxPublisher
│
├─ websocket
│ ├─ AgentWebSocketHandler
│ ├─ SessionChannelRegistry
│ └─ EventStreamRouter
│
├─ approval
│ └─ ApprovalService
│
└─ projection
└─ AgentContextProjection
DB 역할은 다음처럼 나눌 수 있다.
PostgreSQL
agent_event
→ Append-only Event Store
idempotency_request
→ 중복 실행 제어
outbox_event
→ Transactional Event Publishing
PGMQ
→ Worker Job Queue
session / projection table
→ 빠른 조회용 현재 상태
pgvector
→ Memory / History / RAG 검색
처음 구현한다면 어디부터 시작해야 할까
이 구조를 처음부터 완성형으로 만들려고 하면 거의 확실히 지친다.
현실적으로는 단계적으로 만드는 것이 좋다.
EventEnvelope
↓
Event Store
↓
PGMQ Queue
↓
단일 Worker
↓
Idempotency
↓
Worker Pool
↓
Tool Progress
↓
WebSocket
↓
Cancellation
↓
Outbox
↓
Replay / Projection
↓
Scale Out / Stream Sharding
먼저 LLM_REQUESTED → LLM_COMPLETED 하나만 Queue와 Worker로 처리해 보는 것이 좋다.
그다음:
TOOL_REQUESTED
→ TOOL_STARTED
→ TOOL_COMPLETED
를 붙인다.
그 후 Worker를 여러 개로 늘리고:
중복 메시지
Worker 강제 종료
Timeout
Retry
를 실제로 발생시켜 본다.
이 구조는 코드가 "정상적으로 돌아가는 것"보다 장애 상황에서도 어떻게 행동하는지를 테스트하는 것이 훨씬 중요하다.
이 구조가 해결하려는 진짜 문제
Event Driven Architecture를 쓰는 목적이 단순히 최신 기술을 많이 넣는 것은 아니다.
각 기술이 해결하는 문제는 명확하다.
구성 요소 해결하려는 문제 Event Sourcing / Event Log 실행 이력, 복구, Replay, Audit PGMQ Worker 간 안전한 작업 분배 Worker Pool 동시 처리와 자원 관리 Visibility Timeout Worker 장애 후 Job 복구 Heartbeat 장시간 Job의 중복 Claim 방지 Idempotency Retry와 중복 실행 방지 Request Hash 같은 Key의 잘못된 재사용 검증 Outbox DB와 Event 발행의 일관성 문제 Sequence Streaming 순서와 누락 감지 Parent Event 인과관계 추적 WebSocket 양방향 비동기 Event 전달 Channel Session 단위 Routing Streaming LLM·Tool의 실시간 진행 표시 Cancellation 실제 Tool/Process까지 중단 Projection Event에서 빠른 현재 상태 조회 Sharding 대규모 WebSocket 연결 분산
이 표만 제대로 이해해도 전체 구조의 절반 이상은 잡힌다.
결국 핵심은 "Agent를 긴 함수로 보지 않는 것"
기존 Agent는 이렇게 생각하기 쉽다.
Request
↓
execute()
↓
LLM
↓
Tool
↓
LLM
↓
return
기업용 Agent Runtime에서는 관점을 바꾼다.
Request Event
↓
Queue
↓
Worker
↓
Result Event
↓
Queue
↓
Worker
↓
Progress Event
↓
Stream
↓
Client
그리고 이 구조를 안전하게 만들기 위해:
Event Log
+
PGMQ
+
Worker
+
Idempotency
+
Outbox
+
WebSocket
이 서로 연결된다.
한 줄로 압축하면 이렇다.
기업용 AI Agent를 하나의 오래 실행되는 함수로 만들지 않고, 상태 변화와 작업을 이벤트로 기록하고 Queue와 Worker를 통해 독립적으로 처리하며, 중복·장애·재시도·취소·실시간 스트리밍까지 견딜 수 있는 Workflow Runtime으로 설계하는 것이다.
이 관점이 잡히면 Event Sourcing, CPS, PGMQ, Worker Pool, Transaction Key, Sequence, WebSocket, Outbox가 따로 노는 기술이 아니라 하나의 문제를 각자 다른 위치에서 해결하고 있는 구성 요소라는 게 보이기 시작한다.
그리고 그 지점부터가 단순히 "LLM API를 붙인 서버"와 실제로 운영 가능한 Agent Server의 차이다.