요약
AI 에이전트를 프로덕션에 올리면 두 가지 문제가 반드시 나타난다. 첫째, 장애가 발생했을 때 중간에 실행된 외부 액션(DB 쓰기, API 호출, 메시지 전송)을 어떻게 처리할 것인가. 둘째, 에이전트가 대화 맥락을 넘어 장기 기억을 어떻게 유지할 것인가.
Apache Flink Agents 0.3은 두 문제를 모두 다룬다. 에이전트를 Flink 데이터플로우 그래프의 노드로 배치함으로써 Flink가 10년 동안 해결해온 장애 내구성과 Exactly-Once 보장을 에이전트 액션에도 확장한다. 0.3 시리즈에서는 여기에 세 가지가 추가됐다.
- Agent Skills: 프롬프트·도구·리소스를 패키지로 묶어 에이전트가 런타임에 발견·로드하는 표준 단위
- Mem0 기반 Long-Term Memory: 벡터 스토어 직접 연결에서 Mem0로 교체해 의미론적 검색·요약·격리를 강화
- Durable Execution Reconciler: 장애 복구 시 진행 중인 외부 부작용을 중재하는 콜백 메커니즘
- Flink Agents 0.3.0 — 2026년 6월 19일 공개
- Flink Agents 0.3.1 — 2026년 7월 25일 공개 (Flink 2.3 지원 + 버그 수정)
- Confluent · Alibaba Cloud · LinkedIn · Ververica 공동 개발
- Preview 상태: API 및 설정은 향후 비호환 변경 가능
배경: 에이전트 인프라로서의 Flink
LLM 에이전트를 직접 운영할 때의 어려움
별도 에이전트 서버를 운영하면 다음 문제를 직접 풀어야 한다.
| 문제 | 설명 |
|---|---|
| 장애 복구 | 에이전트가 도중에 멈추면 어디서 재시작할지 알 수 없다 |
| Exactly-Once 액션 | 장애 후 재실행 시 이미 완료된 외부 API 호출이 중복 실행된다 |
| 상태 관리 | 멀티-턴 대화 상태·맥락을 직접 지속시켜야 한다 |
| 이벤트 기반 실행 | Kafka 이벤트가 도착할 때만 에이전트를 깨우는 구조를 별도로 구현해야 한다 |
| 장기 기억 | 대화 창 너머의 사용자 선호·도메인 지식을 저장·검색해야 한다 |
Flink는 10년에 걸쳐 이 중 앞 네 가지를 스트리밍 파이프라인을 위해 풀어왔다. Flink Agents는 에이전트 실행 모델을 Flink의 연산자(Operator)로 표현함으로써 이 인프라를 재활용한다.
Flink 데이터플로우에 에이전트를 올리는 이유
Flink에서 에이전트는 ActionExecutionOperator로 컴파일된다. 이 연산자는 Flink 작업(Job)의 DAG 안에서 다른 연산자와 동등한 위치를 차지한다.
Kafka Source → 이벤트 파싱 → [AgentOperator] → 결과 싱크
↕
Flink 체크포인트 메커니즘
External Action State Store이 구조가 갖는 즉각적인 이점은 두 가지다.
- Flink 체크포인트가 에이전트 상태(LLM 대화 이력, 상태 머신 위치)를 주기적으로 스냅샷한다. 장애 후 체크포인트에서 재시작한다.
- Kafka 기반 외부 액션 스토어가 이미 실행 완료된 외부 액션을 기록한다. 재실행 시 완료된 액션은 건너뛴다.
핵심 구성 요소
Exactly-Once 액션 보장
LLM 에이전트의 액션(외부 API 호출, DB 쓰기 등)은 기본적으로 멱등성이 없다. Flink 체크포인트만으로는 완전하지 않다. 체크포인트 이후에 외부 시스템에 쓰기가 일어났다면, 복구 후 재실행 시 동일 쓰기가 두 번 발생한다.
Flink Agents는 이를 외부 액션 스토어(External Action State Store)로 해결한다.
에이전트 액션 시작
↓
External Action State Store에 "IN_PROGRESS" 기록 (Kafka)
↓
실제 외부 시스템 호출
↓
External Action State Store에 "COMPLETED" 기록
↓
Flink 체크포인트 완료
--- 장애 발생 ---
↓
체크포인트에서 복구
↓
External Store 조회 → COMPLETED 기록 발견 → 해당 액션 스킵현재 Kafka가 유일한 공식 외부 액션 스토어 구현체다. Kafka의 트랜잭션 보장을 활용해 "기록 완료" 자체도 Exactly-Once로 처리한다.
Durable Execution Reconciler (0.3.0 신규)
Exactly-Once 스토어만으로는 부족한 경우가 있다. 예를 들어 에이전트가 파일을 업로드 중이거나, 장시간 실행 API를 폴링 중인 상황에서 장애가 발생하면 "진행 중인 부작용"을 어떻게 처리할지 알 수 없다.
0.3.0에서 도입된 Durable Execution Reconciler는 이 상황을 위한 콜백이다.
class FileUploadAction:
@durable_action
async def upload(self, file_path: str):
upload_id = await start_multipart_upload(file_path)
self.register_reconciler(self._reconcile_upload, upload_id=upload_id)
await complete_multipart_upload(upload_id)
async def _reconcile_upload(self, upload_id: str):
status = await check_upload_status(upload_id)
if status == "INCOMPLETE":
await abort_upload(upload_id)register_reconciler()를 호출하면 이후 장애 복구 시 Flink가 reconciler 콜백을 먼저 실행한다. reconciler는 부작용 완료 여부를 확인하고 정리한다. reconciler 실행 중 예외가 발생하면 해당 실패가 액션 스토어에 기록된다.
Agent Skills (0.3.0 신규)
Agent Skills는 에이전트가 런타임에 발견하고 로드하는 자기 기술 패키지(self-describing capability package)다.
기존 LLM 에이전트에서 도구는 에이전트 정의 시점에 하드코딩된다. Skills는 이를 분리한다. 에이전트 코드를 변경하지 않고 Skill을 추가·제거·버전 업할 수 있다.
# Java API 예시 (0.3.0)
AgentSkill translationSkill = AgentSkill.builder()
.name("translate")
.description("텍스트를 지정 언어로 번역합니다")
.prompt("당신은 전문 번역가입니다. {{source_lang}}에서 {{target_lang}}으로...")
.tools(List.of(translationTool, glossaryLookupTool))
.resources(List.of(terminologyDatabase))
.build();Mem0 기반 Long-Term Memory (0.3.0 신규)
이전 버전까지 Long-Term Memory는 벡터 스토어(pgvector, Qdrant 등)를 직접 연결하는 방식이었다. 에이전트가 기억을 직접 임베딩하고, 유사도 검색으로 회수하는 구조였다.
0.3.0에서 이를 Mem0 기반으로 교체했다. Mem0는 단순 벡터 유사도 외에 세 가지를 추가로 처리한다.
| 기능 | 이전 방식(벡터 스토어) | 0.3.0 (Mem0) |
|---|---|---|
| ADD | 임베딩 생성 후 저장 | 중복/모순 검사 후 신규/갱신/삭제 결정 |
| 검색 | 의미론적 유사도 | 의미론적 + BM25 키워드 + 엔티티 그래프 |
| 요약 | 없음 | 기억 자동 요약·통합 |
| 격리 | 수동 필터 | user_id · run_id · agent_id 자동 스코프 |
에이전트가 "이 사용자는 MySQL DBA다"라는 사실을 기억에 추가하면 Mem0는 이전에 "데이터베이스 관리자"로 기록된 유사 기억과 통합한다. 중복 기억이 쌓이지 않는다.
아키텍처 전체 그림
0.3.x 시리즈 변경 사항
0.3.0 (2026-06-19)
- Agent Skills 도입 — Python, Java API 모두 지원
- Durable Execution Reconciler 도입 — reconciler 콜백 등록 및 실패 지속 기록
- Mem0 기반 Long-Term Memory — 기존 벡터 스토어 직접 연결 방식 대체
- Flink 2.2 지원 확인
0.3.1 (2026-07-25)
- Flink 2.3 배포 지원 추가
- 설치 개선 (installer)
- 버그 수정 5건
운영 고려사항
프리뷰 버전의 의미
0.3.x는 공식 프리뷰(Preview) 상태다. 공식 문서에 명시된 주의사항이 있다.
- API와 설정 옵션은 향후 비호환 변경이 가능하다
- 일부 기능에 알려지거나 알려지지 않은 버그가 있을 수 있다
- 프로덕션 배포는 특성에 맞게 검토가 필요하다
실무에서는 낮은 위험도의 내부 자동화 파이프라인, 이벤트 분류, 내부 챗봇 등에서 선도적으로 검증하는 것이 적합하다.
Kafka 외부 액션 스토어 설정
Exactly-Once 보장을 받으려면 Kafka를 외부 액션 스토어로 설정해야 한다. 설정 없이 실행하면 Local 실행 모드(in-memory state)로 동작하며 Exactly-Once 보장이 없다.
# flink-agents 설정 예시
action_state_store:
type: kafka
bootstrap_servers: "kafka-broker:9092"
topic: "agent-action-state"
transaction_timeout_ms: 60000Kafka 토픽의 파티션 수는 Flink 병렬도에 맞춰야 한다. 파티션 수가 병렬도보다 적으면 일부 연산자가 스토어에 접근하지 못한다.
Mem0 벡터 스토어 선택
Mem0 백엔드는 내부적으로 벡터 스토어를 사용한다. 지원되는 옵션:
- Qdrant (기본 추천): 로컬 또는 클라우드 배포, 필터링 성능 우수
- pgvector: 기존 PostgreSQL 인프라 활용 시
- Chroma: 경량 개발 환경
프로덕션에서는 agent_id와 user_id를 명시적으로 설정해 메모리 격리를 보장해야 한다. 격리 설정 없이 여러 에이전트가 같은 Mem0 인스턴스를 쓰면 기억이 혼재된다.
체크포인트 주기와 에이전트 실행 시간
LLM 호출은 수 초에서 수십 초가 걸린다. Flink 기본 체크포인트 주기(1분)는 많은 LLM 작업이 체크포인트 사이에 완료될 만큼 충분히 길다.
그러나 매우 긴 ReAct 루프(수십 번의 LLM + 도구 호출 반복)는 체크포인트 주기 내에 완료되지 않을 수 있다. 이 경우:
- 체크포인트를 더 자주 트리거하거나 (단, 체크포인트 비용 증가)
- Durable Execution Reconciler를 각 주요 액션 단계에 등록해 부분 복구를 가능하게 한다
다른 에이전트 프레임워크와 비교
Flink Agents가 적합한 상황과 그렇지 않은 상황은 다음 기준으로 판단할 수 있다.
| 기준 | Flink Agents | LangGraph / CrewAI | Temporal / DBOS |
|---|---|---|---|
| Exactly-Once 액션 | Kafka + Flink 체크포인트 | 없음 (별도 구현 필요) | 워크플로우 단위 |
| 이벤트 드리이브 | Kafka 네이티브 | 별도 연동 필요 | 폴링 또는 시그널 |
| 상태 관리 | Flink 분산 상태 | 인메모리 / 외부 스토어 | 워크플로우 상태 |
| Long-Term Memory | Mem0 (내장) | 별도 연동 | 별도 연동 |
| 학습 곡선 | Flink 지식 필요 | 낮음 | 중간 |
| 프로덕션 성숙도 | Preview | GA | GA |
Flink Agents가 명확히 유리한 경우:
- 조직에 이미 Flink 기반 스트리밍 파이프라인이 운영 중이고, 에이전트를 그 파이프라인 안에 배치하고 싶을 때
- 고부하 이벤트 스트림(초당 수천 건)에서 에이전트를 병렬로 실행하고 Exactly-Once를 보장해야 할 때
- Kafka 이벤트가 에이전트 실행의 주 트리거인 아키텍처에서
Flink Agents가 과도한 경우:
- 대화형 챗봇처럼 낮은 동시성·짧은 세션의 에이전트
- Flink 클러스터 없이 서버리스나 단순 컨테이너 배포가 목표인 경우
요점 정리
Apache Flink Agents 0.3은 스트리밍 데이터 플랫폼 위에 AI 에이전트를 올리는 가장 구조적인 방법 중 하나다.
핵심 가치는 세 가지다. Exactly-Once 보장은 에이전트가 프로덕션에서 가장 먼저 부딪히는 장애 복구 문제를 Kafka와 Flink 체크포인트 조합으로 해결한다. Agent Skills는 에이전트 코드 변경 없이 도구·프롬프트·리소스를 교체 가능하게 한다. Mem0 기반 Long-Term Memory는 에이전트가 대화 창 너머로 기억을 유지할 때 발생하는 중복·모순 문제를 자동으로 처리한다.
다만 Preview 상태라는 점은 실질적인 운영 경계다. 0.1.0(2025-10)부터 0.3.1(2026-07)까지 빠르게 발전하고 있지만, API 변경 리스크를 감수할 수 있는 맥락에서 도입을 고려해야 한다.
References
- Apache Flink Agents 0.3.0 릴리스 노트: https://flink.apache.org/2026/06/19/apache-flink-agents-0.3.0-release-announcement/
- Apache Flink Agents 0.3.1 릴리스 노트: https://flink.apache.org/2026/07/25/apache-flink-agents-0.3.1-release-announcement/
- FLIP-531 제안서: https://cwiki.apache.org/confluence/display/FLINK/FLIP-531:+Initiate+Flink+Agents+as+a+new+Sub-Project
- Flink Agents 공식 문서: https://nightlies.apache.org/flink/flink-agents-docs-latest/
- Confluent 블로그 — Why Flink Agents Are the Future of Enterprise AI: https://www.confluent.io/blog/event-driven-flink-agents-enterprise-ai/
- Mem0 공식 문서: https://docs.mem0.ai