LLM WikiAccess-protected knowledge portal

WIKI

Apache Flink Agents 0.3: 스트리밍 이벤트 드리이브 AI 에이전트에 Exactly-Once 보장·Agent Skills·Mem0 기억을 더한 운영 아키텍처

요약 AI 에이전트를 프로덕션에 올리면 두 가지 문제가 반드시 나타난다. 첫째, 장애가 발생했을 때 중간에 실행된 외부 액션 DB 쓰기, API 호출, 메시지 전송 을 어떻게 처리할 것인가. 둘째, 에이전트가 대화 맥락을 넘어 장기 기억을 어떻게 유지할 것인가. Apache Flink Agents 0.3 은 두 문제를 모두 다룬다. 에이전트를 Flink 데이터플로우 그래프의 노드로 배치함으로써 Flink가 10년 동안 해결해온

경로human/study/content/database-frontier/95-flink-agents-0-3-event-driven-ai-checkpointing-skills-mem0.md
카테고리Study
태그#ai-review #airflow #checkpointing #driven #event #infra #mem0 #mysql #skills #study

요약

AI 에이전트를 프로덕션에 올리면 두 가지 문제가 반드시 나타난다. 첫째, 장애가 발생했을 때 중간에 실행된 외부 액션(DB 쓰기, API 호출, 메시지 전송)을 어떻게 처리할 것인가. 둘째, 에이전트가 대화 맥락을 넘어 장기 기억을 어떻게 유지할 것인가.

Apache Flink Agents 0.3은 두 문제를 모두 다룬다. 에이전트를 Flink 데이터플로우 그래프의 노드로 배치함으로써 Flink가 10년 동안 해결해온 장애 내구성과 Exactly-Once 보장을 에이전트 액션에도 확장한다. 0.3 시리즈에서는 여기에 세 가지가 추가됐다.


배경: 에이전트 인프라로서의 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

이 구조가 갖는 즉각적인 이점은 두 가지다.

  1. Flink 체크포인트가 에이전트 상태(LLM 대화 이력, 상태 머신 위치)를 주기적으로 스냅샷한다. 장애 후 체크포인트에서 재시작한다.
  2. 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)다.

Skill 패키지
프롬프트 + 도구 + 리소스
Skill Registry
중앙 발견 서비스
에이전트 런타임
필요 시 동적 로드
실행
컨텍스트 주입 후 호출
Agent Skills 생명주기

기존 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는 이전에 "데이터베이스 관리자"로 기록된 유사 기억과 통합한다. 중복 기억이 쌓이지 않는다.


아키텍처 전체 그림

Apache Flink Agents 0.3 — 프로덕션 배포 구조 이벤트 소스 Kafka / HTTP Flink Job (장애 내구성 · Exactly-Once 영역) Source Op 이벤트 역직렬화 AgentOperator ActionExecutionOperator Agent Skills 로드 (런타임) ReAct / 동적 토폴로지 실행 Durable Execution Reconciler Sink Op 결과 발행 Flink 체크포인트 상태 대화 이력 · 상태 머신 · 스킬 캐시 LLM 프로바이더 OpenAI / Anthropic Claude / GPT / Gemini MCP 도구 서버 파일 · DB · API Model Context Protocol External Action Store Kafka (현재 유일) Exactly-Once 액션 기록 Mem0 메모리 Long-Term Memory 의미론적 + 키워드 검색 Skill Registry Agent Skills 카탈로그 프롬프트·도구 패키지 Flink 분산 체크포인트 스토어 (S3 / HDFS / GCS) 주기적 스냅샷 — 장애 후 자동 재시작 지점 에이전트 상태 · 대화 이력 · 스킬 캐시 포함
Flink Agents 0.3 프로덕션 배포 아키텍처

0.3.x 시리즈 변경 사항

0.3.0 (2026-06-19)

0.3.1 (2026-07-25)


운영 고려사항

프리뷰 버전의 의미

0.3.x는 공식 프리뷰(Preview) 상태다. 공식 문서에 명시된 주의사항이 있다.

실무에서는 낮은 위험도의 내부 자동화 파이프라인, 이벤트 분류, 내부 챗봇 등에서 선도적으로 검증하는 것이 적합하다.

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: 60000

Kafka 토픽의 파티션 수는 Flink 병렬도에 맞춰야 한다. 파티션 수가 병렬도보다 적으면 일부 연산자가 스토어에 접근하지 못한다.

Mem0 벡터 스토어 선택

Mem0 백엔드는 내부적으로 벡터 스토어를 사용한다. 지원되는 옵션:

프로덕션에서는 agent_id와 user_id를 명시적으로 설정해 메모리 격리를 보장해야 한다. 격리 설정 없이 여러 에이전트가 같은 Mem0 인스턴스를 쓰면 기억이 혼재된다.

체크포인트 주기와 에이전트 실행 시간

LLM 호출은 수 초에서 수십 초가 걸린다. Flink 기본 체크포인트 주기(1분)는 많은 LLM 작업이 체크포인트 사이에 완료될 만큼 충분히 길다.

그러나 매우 긴 ReAct 루프(수십 번의 LLM + 도구 호출 반복)는 체크포인트 주기 내에 완료되지 않을 수 있다. 이 경우:


다른 에이전트 프레임워크와 비교

Flink Agents가 적합한 상황과 그렇지 않은 상황은 다음 기준으로 판단할 수 있다.

기준Flink AgentsLangGraph / CrewAITemporal / DBOS
Exactly-Once 액션Kafka + Flink 체크포인트없음 (별도 구현 필요)워크플로우 단위
이벤트 드리이브Kafka 네이티브별도 연동 필요폴링 또는 시그널
상태 관리Flink 분산 상태인메모리 / 외부 스토어워크플로우 상태
Long-Term MemoryMem0 (내장)별도 연동별도 연동
학습 곡선Flink 지식 필요낮음중간
프로덕션 성숙도PreviewGAGA

Flink Agents가 명확히 유리한 경우:

Flink Agents가 과도한 경우:


요점 정리

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