LLM WikiAccess-protected knowledge portal
← 스터디 홈
95편 · 약 14분

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

요약

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

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

  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)

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

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

Mem0 벡터 스토어 선택

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

  • Qdrant (기본 추천): 로컬 또는 클라우드 배포, 필터링 성능 우수
  • pgvector: 기존 PostgreSQL 인프라 활용 시
  • Chroma: 경량 개발 환경

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

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

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

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

  • 체크포인트를 더 자주 트리거하거나 (단, 체크포인트 비용 증가)
  • Durable Execution Reconciler를 각 주요 액션 단계에 등록해 부분 복구를 가능하게 한다

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

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

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

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