Issue 01중국 AI
AC POST
중국 AI 목록
掘金2026년 9월 16일 12:27중국어 → 한국어

Kafka에 AI 연결하는 세 가지 경로 정리

Kafka를 AI에 연결하는 세 경로로 1차 MCP 서버 제안, Streams 세션 기억, Flink SQL 실시간 컨텍스트를 제시하고 선정 비교와 권한 목록을 담았다.

중국어 원문을 AI로 번역했습니다. 고유명사와 수치는 원문 표기를 우선하며, 중요한 판단에는 아래 출처 원문을 함께 확인하세요.

안녕하세요, 저는 프로그래머 톈톈쿤(天天困)입니다.

Kafka 4.3 릴리스 공지를 번역하다가 잠깐 멈칫했습니다. 감사 명단에 Claude, Claude Sonnet 4.6, Copilot이 적혀 있었거든요. 메시지 큐의 버전 릴리스에 AI가 기여자로 섞여 들어온 겁니다. 그런데 Kafka가 AI를 받아들이는 데서 정말 이야기할 가치가 있는 것은, AI가 Kafka 코드를 대신 써주게 하는 게 아니라 AI 시스템에 세 가지 부류의 능력을 보강해 주는 것입니다. Kafka를 조작하고, 컨텍스트를 기억하고, 데이터를 실시간으로 처리하는 것.

이 글은 세 가지 노선을 한 번에 다룹니다. 논의 중인 1자(제1자) MCP Server 제안, Streams 세션 메모리, Flink SQL 실시간 컨텍스트입니다. 각각에 대해 지금의 실제 상태, 예제 코드의 전제 조건, 그리고 지금은 열어서는 안 되는 권한이 무엇인지까지 다 적어두겠습니다. 즐겨찾기 눌러두고 시작합니다.

Kafka 자체가 아직 익숙하지 않다면, 제가 쓴 이 Kafka 해설 글을 먼저 보면서 topic, 파티션, 컨슈머 그룹 같은 기본 개념을 훑고 오시길 권합니다. 이 글은 그 개념들이 무엇인지 이미 알고 있다고 전제합니다.

먼저 증거부터 보죠. 이건 농담이 아닙니다.

이 스크린샷은 Apache Kafka 4.3.0 Release Announcement 하단 Summary 단락에서 나온 것으로, 첫 번째 빨간 박스는 「147 contributors (and 3 AIs)」이고 두 번째는 명단에 있는 Claude, Claude Sonnet 4.6, Copilot입니다.

1. 결론부터: Kafka가 AI를 받아들이는 것은 조작, 기억, 실시간 처리 세 가지 능력을 보강하는 일이다

Kafka가 AI를 받아들이는 것을 가장 쉽게 잘못 설명하는 지점은, Kafka에 두뇌를 달아줘야 한다고 사람들이 생각하게 만드는 것입니다. 그럴 필요는 없습니다. 앞의 두 가지 능력, 즉 Agent가 조작할 수 있게 하고 기억할 수 있게 하는 것은 많은 Agent 시스템이 보강해야 하는 약점입니다. 세 번째는 모델 추론을 스트림 처리에 연결해서 데이터가 흐르는 순간 판단을 함께 지니게 하는 것입니다.

MCP(Model Context Protocol, 모델 컨텍스트 프로토콜): Anthropic이 2024년 11월 25일에 제안하고 2025년 12월에 Linux 재단 Agentic AI Foundation에 기부한 개방형 표준으로, JSON-RPC 2.0을 사용해 외부 시스템을 AI가 호출할 수 있는 도구와 리소스로 포장합니다. 「AI 도구의 USB-C 인터페이스」라고 이해하면 됩니다.

Kafka MCP Server(KIP-1318 제안): 아직 논의 중인 커뮤니티 제안으로, Kafka에 1자 MCP 서비스 프로세스를 추가해 topic, 컨슈머 그룹, ACL, Connect 같은 운영 작업을 표준 도구로 노출할 계획입니다. 클러스터에 MCP를 할 줄 아는 운영 당직자 한 명을 붙여주는 셈인데, 주의하세요. 이건 현재 제안일 뿐, 이미 전달된 기능이 아닙니다.

컨텍스트 저장소(Context Store): 이 글에서는 「세션 키 기준으로 저장되고 컨텍스트를 조회할 수 있는 뷰」를 가리키는 말로 쓰며, 아래에서는 Kafka Streams로 구현합니다. 쉽게 말해, 흩어져 있는 채팅 기록을 언제든 조회할 수 있는 표 하나로 압축하는 것입니다.

비유를 하나 해보죠. 예전의 Kafka는 회사 안내 데스크와 같았습니다. 방문객이 오면 직접 서식을 쓰고 부서를 직접 찾아가며, 안내 데스크는 접수증을 넘겨주는 일만 했습니다. MCP는 안내 데스크에 표준 Q&A 매뉴얼을 갖다 놓은 것이고, AI 비서는 그걸 보고 일을 처리할 수 있습니다. 컨텍스트 저장소는 안내 데스크에 회의록을 갖다 놓은 것으로, 지난번에 누가 무슨 말을 했는지 펼쳐보면 바로 알 수 있습니다. 스트림 처리에서의 모델 호출은 안내 데스크 옆에 분석가 한 명이 앉아 있는 것과 같아서, 접수증이 아직 순환 중인데도 이미 표시가 달려 있습니다. 매뉴얼은 「무엇을 할 수 있는가」를, 회의록은 「무엇을 기억하는가」를, 분석가는 「실시간으로 꿰뚫어 보기」를 맡습니다. 이 세 가지가 바로 Agent를 실제로 적용할 때 가장 힘든 부분입니다.

세 노선은 각각 그중 한 조각을 해결합니다.

- 노선 1, Agent가 Kafka를 조작할 수 있게 하기: 논의 중인 1자 MCP Server, 클러스터 운영이 호출 가능한 도구가 됩니다.

- 노선 2, Kafka가 Agent의 기억이 되게 하기: Streams가 세션 로그를 KTable로 집계하고, Agent는 조회해서 읽어옵니다.

- 노선 3, 모델 추론을 스트림 처리로 들여오기: Flink SQL에서 AI 함수를 호출해, 데이터가 아직 스트림 안에 있을 때 라벨을 달아둡니다.

셋은 대체 관계가 아니고 위치도 다릅니다. 우선 전체 그림부터 보죠.

2. 노선 1: 논의 중인 1자 MCP Server, Agent가 Kafka를 조작하게 하기

세 노선 중 MCP 이 노선이 「Kafka가 이미 AI를 받아들였다」는 확증으로 가장 쉽게 받아들여지므로, 먼저 찬물을 좀 끼얹어야겠습니다. 지금은 여전히 하나의 제안이고, 어떤 정식 버전에도 함께 전달되지 않았습니다.

KIP-1318의 현재 상태는 Under Discussion이며, 2026년 9월 16일 기준으로도 아직 릴리스 상태에 들어가지 못했습니다. 그 타임라인은 대략 이렇습니다.

시간 이벤트 2026-04-07 KIP-1318 제안 생성 2026-04-21 dev@kafka.apache.org 로 논의 메일 발송 (메일 원문 날짜 기준. 베이징 시간으로는 이미 4월 22일) 2026-05-22 / 2026-06-25 4.3.0, 4.3.1이 차례로 릴리스되었지만 공지에는 둘 다 없음 2026-08-01 wiki 최종 수정, 상태는 여전히 Under Discussion 2026-08-11 / 08-13 / 08-21 4.4 브랜치 분기, 코드 동결, RC0 투표 개시. 그 릴리스 계획에도 마찬가지로 KIP-1318은 없음

표 이외에도 몇 가지 덧붙일 것이 있습니다. 4.3.1은 Kafka Streams RocksDB 네이티브 메모리 누수를 수정하는 bugfix 버전이라 애초에 새 기능을 추가하지 않습니다. 4.4 쪽은 릴리스 계획의 계획 날짜가 아니라 메일링 리스트의 실제 날짜를 쓴 것입니다. KAFKA-20436은 이 제안을 추적하는 JIRA 작업으로, 아직 Open, 미할당, 수정 버전 미지정 상태입니다. 추적 작업이 있다고 해서 개발이 이미 사용 가능한 단계까지 진행되었다는 뜻은 아닙니다. 또 확인 시점 기준으로 공식 블로그의 최신 릴리스 공지는 4.3.1이었고, 저는 4.4의 정식 릴리스 공지를 보지 못했습니다.

그러므로 지금 「Kafka 공식이 AI Agent의 클러스터 조작을 지원한다」고 말하는 것은 정확하지 않습니다. 정확한 표현은 이렇습니다. 커뮤니티 구성원이 KIP를 제안했고, Kafka에 1자 MCP 인터페이스를 추가할지 말지를 논의 중이다.

그렇다면 왜 기다릴 만한가? 운영 면이 정말 넓기 때문입니다. Producer, Consumer, Streams, Connect, Admin 다섯 가지 핵심 API가 커버하는 작업이 많고(제안 작성자는 100개가 넘는다고 추정하는데, 이 숫자에는 공개된 집계 기준이 없습니다), Agent가 흔히 하는 일, 즉 topic 하나 만들기, 컨슈머 그룹 적체 확인하기, 죽은 connector 재시작하기 같은 것은 오늘날 Java/Python 클라이언트 코드를 쓰거나, kafka-topics.sh 같은 명령줄을 두드리거나, Connect의 REST 인터페이스를 curl 하는 방식 중 하나여야 합니다. 이런 능력은 물론 존재하지만, Agent가 그 안에서 일하려면 도구를 하나하나 포장하고 권한을 하나하나 설정해야 하며, 통일된 진입점이 없습니다.

KIP-1318의 해법은 독립적인 서비스 프로세스를 추가하는 것입니다(tools/mcp-server 모듈에 두고, broker 안에서는 돌지 않음). Kafka의 wire protocol은 바꾸지 않으며, 핵심 작업은 공식 Java 클라이언트로 포장하고 Connect 작업은 Kafka Connect의 REST API를 통해 수행합니다. 그것이 노출할 수 있는 것은 대략 다음과 같이 분류됩니다.

분류 대표 작업 읽기인지 쓰기인지 Topic 관리 생성, 삭제, 설정 변경, 파티션 추가 쓰기(삭제는 고위험) 메시지 읽기/쓰기 프로듀스, 일괄 프로듀스, 트랜잭션 프로듀스, 컨슘 읽기 + 쓰기 컨슈머 그룹 적체 조회(읽기), 오프셋 리셋 및 멤버 제거(쓰기) 읽기 + 쓰기 ACL 및 클러스터 ACL 추가/삭제, broker 설정 변경, leader 선출 고위험 쓰기 Connect connector와 task 생성, 일시정지, 재시작 쓰기

리소스 측은 읽기 전용이며, kafka:// 로 시작하는 URI로 노출됩니다. 예를 들어 kafka://topics , kafka://groups/{id}/lag , kafka://cluster/metadata-quorum 같은 식입니다. 클라이언트는 해당 리소스 URI를 통해 구조화된 lag 정보를 읽을 수 있고, Admin API를 직접 조립할 필요가 없습니다.

제안의 보안 설계는 별도로 한 단락 볼 가치가 있습니다. KIP는 일련의 안전장치를 나열합니다. mcp.readonly=true 는 모든 쓰기 작업을 꺼버리고 produce까지 끕니다. 주의하세요, 기본값은 false이며 기본이 읽기 전용이 아닙니다. mcp.tools.allowed 와 mcp.tools.denied 는 각각 허용 목록과 거부 목록이고, 후자가 우선합니다. mcp.allowed.topic.prefixes 는 조작 가능한 topic을 agent. , sandbox. 같은 접두사 안으로 제한하지만, 이것은 topic 수준 작업만 제약하며 클러스터와 ACL 관리 권한은 읽기 전용 스위치, 목록, broker 측 ACL이 함께 받쳐줘야 합니다. 파괴적 도구는 기본적으로 대외(아웃오브밴드) 발급된 인간 승인 토큰을 요구합니다. mcp.audit.topic 은 매 호출(신원, 도구, 파라미터, 결정, 결과)을 지정된 topic에 추가 기록하여 감사에 사용합니다.

또 하나의 방어선은 taint guard입니다. 이것은 신뢰할 수 없는 읽기 결과를 파괴적 도구의 파라미터와 정규화해 정확히 또는 부분 문자열로 매칭하려 시도하고, 적중하면 유효한 승인 토큰을 요구하며, 그렇지 않으면 -32040 을 반환합니다. 이것은 어디까지나 최선의 탐지일 뿐이며, 모델의 재작성, 재인코딩, 부분 인용은 우회할 수 있습니다. 파괴 범위를 실제로 제한하는 것은 여전히 승인, 읽기 전용 모드, 리소스 범위, broker 측 최소 권한 ACL입니다.

완전한 호출 하나는 이렇게 생겼습니다.

누군가는 이렇게 물을 수 있습니다. KIP-1318이 아직 전달되지 않았다면, 지금 나는 무엇을 쓸 수 있나?

지금도 커뮤니티와 벤더가 제공하는 Kafka MCP 구현이 있습니다. mcp-confluent 를 예로 들면, 현재 공식 설명은 Confluent Cloud, Confluent Platform, 독립형 Apache Kafka 배포를 지원한다고 되어 있습니다. 두 가지를 유의하세요. 이 오픈소스 구현은 커뮤니티 지원에 속하며 공식적으로 best-effort이고 서비스 수준 약정이 없다고 명시되어 있습니다. 그리고 원래 Confluent Cloud를 쓰고 있다면 그쪽에는 별도의 완전관리형 MCP Server가 있어 권한은 기존 RBAC를 따르지만, 커버하는 것은 Confluent Cloud 상의 리소스입니다. 구체적으로 사용 가능한 도구, 인증 방식, 배포 제약은 역시 사용 중인 버전에 따라 항목별로 확인해야 합니다. 구현마다 커버 범위 차이가 크고, KIP에 있는 서드파티 구현에 대한 초기 비교도 현황으로 그대로 받아들이기에는 적절하지 않습니다.

3. 노선 2: Kafka를 Agent의 기억으로 쓰기, Streams로 세션 컨텍스트 저장하기

시스템이 이미 Kafka를 이벤트 백본으로 삼고 있고 세션 상태를 읽어야 할 컴포넌트가 여러 개라면, 이 길의 가성비가 가장 높습니다. 세션 메모리 때문에 굳이 벡터 DB를 서둘러 도입할 필요는 없고, Kafka Streams로 KTable로 떨어뜨리면 충분합니다.

이치는 복잡하지 않습니다. Agent끼리 Kafka로 메시지를 주고받을 때, 대화 자체가 이미 재생 가능한 로그이기 때문입니다. 사용자 메시지, 봇 응답, 상담원 개입이 모두 각자의 topic에 누워 있습니다. 이를 다시 Postgres나 Redis에 옮겨 한 번 더 저장하려면 외부 저장소와 동기화 경로 한 세트를 추가로 유지해야 할 수 있습니다. Streams는 이 topic들을 그대로 읽어 세션 ID로 집계하므로, 기억이 메시지와 같은 데이터 평면에 그대로 머뭅니다.

그런데 한 가지는 먼저 분명히 해야 합니다. Kafka의 순서 보장은 파티션 내에서만 성립합니다. 세 topic 모두 sessionId를 key로 쓰고 파티션 수도 같다고 해도, merge() 는 서로 다른 입력 스트림 간의 상대적 순서를 보장하지 않습니다. 즉 Kafka는 파티션 내 이벤트에 대해 재생 가능한 순서 로그를 제공하지만, topic 간 병합 시에는 세션 순서, 지연 도착, 중복 메시지를 여전히 명시적으로 처리해야 합니다. 예를 들어 Turn에 세션 순번, 이벤트 시간, 고유 이벤트 ID를 붙이고, 순서 뒤바뀜과 중복 제거 전략을 정해두는 식입니다.

다음은 핵심 토폴로지 조각입니다( Turn , SessionContext 의 타입 정의, serde, 애플리케이션 설정은 생략했으며, 바로 실행할 수는 없습니다).

StreamsBuilder builder = new StreamsBuilder (); // 세 소스 병합: 사용자 질문, 봇 응답, 상담원 응답 KStream<String, Turn> turns = builder .stream( "cs.user.turns" , Consumed.with(Serdes.String(), turnSerde)) .merge(builder.stream( "cs.bot.replies" , Consumed.with(Serdes.String(), turnSerde))) .merge(builder.stream( "cs.agent.replies" , Consumed.with(Serdes.String(), turnSerde)));

집계 전에 명시적으로 리파티션해서, 서로 다른 소스의 key 분포가 같은 대화 세션을 쪼개버리지 않도록 합니다.

turns .repartition(Repartitioned.<String, Turn>as( "session-turns" ) .withKeySerde(Serdes.String()) .withValueSerde(turnSerde)) .groupByKey(Grouped.with(Serdes.String(), turnSerde)) .aggregate( SessionContext:: new , (sessionId, turn, ctx) -> ctx.append(turn), Materialized.<String, SessionContext, KeyValueStore<Bytes, byte []>>as( "session-context" ) .withKeySerde(Serdes.String()) .withValueSerde(contextSerde));

집계 결과는 하나의 KTable이지만, 이것을 「매 턴의 전체 이력을 자동으로 저장해 주는 것」이라고 생각하면 안 된다. KTable은 테이블의 논리적 추상화일 뿐이고, 무엇을 저장할지는 집계기가 결정한다. 로컬 상태 저장소는 보통 RocksDB이며, 다른 store로 설정할 수도 있다. 절약되는 것은 「별도의 외부 데이터베이스와 동기화 컴포넌트 한 세트」이지, 복제가 없거나 저장 비용이 없다는 뜻이 아니다. 원본 topic, 로컬 상태, changelog는 여전히 서로 다른 세 개의 복사본이다.

간과하기 쉬운 것은 상태 크기이기도 하다. changelog의 compaction이 보존하는 것은 각 key의 최신 레코드 버전이며, SessionContext에 계속 append한다고 해서 배열이 작아지지는 않는다. 무한 append는 직렬화, 갱신, 복구의 비용을 지속적으로 끌어올리고, 단일 레코드 크기 상한에 부딪힐 수도 있다. 따라서 최소한 정리 전략을 설계해야 한다. 최근 N턴만 남기기, 더 이전 기록을 요약화하기, 혹은 세션 수명 주기에 따라 삭제하기 등이며, 이는 모두 애플리케이션 로직에 의존해야 하고 compaction에만 기대면 안 된다.

「프로세스가 죽으면 유실되는가」에 관해서는, changelog와 입력 로그를 사용할 수 있고 설정이 올바르다는 전제하에 커밋된 상태는 복구되고 재생될 수 있지만, 복구 자체에는 시간이 걸린다.

「제자리에서 테이블 만들기」라는 발상은 아래 그림이 비교적 직관적으로 설명해 준다.

같은 애플리케이션 안에 두 번째 뷰를 손쉽게 하나 더 붙일 수도 있다. 5분 고정 윈도우(tumbling window, 5분 단위로 겹치지 않음)로 사용자 질문 횟수를 집계하여 빈도 모니터링에 신호를 제공한다.

builder .stream( "cs.user.turns" , Consumed.with(Serdes.String(), turnSerde)) .groupByKey() .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes( 5 ))) .count(Materialized.as( "turns-5m" ));

NoGrace는 grace를 0으로 설정하는 것으로, 윈도우가 스트림 시간 기준으로 닫힌 뒤에 도착하는 레코드는 버려지며, 벽시계(wall clock) 기준으로 계산되지 않는다. 이 카운트는 스로틀링 결정의 입력 중 하나일 뿐이고, 코드 자체가 스로틀링을 수행하지는 않는다. 「사용자가 같은 질문을 반복해서 하고 있는가」를 식별하려면 질문 내용의 정규화나 유사도 판단을 결합해야 하며, sessionId로만 카운트해서는 알 수 없다.

Agent 측의 읽기는 대화형 쿼리(interactive query)에 의존한다. Agent와 Streams가 같은 프로세스에서 실행되고, 이 세션의 상태가 마침 현재 인스턴스에 의해 보유되고 있다면 로컬 상태 저장소를 직접 조회하면 된다.

ReadOnlyKeyValueStore<String, SessionContext> store = streams.store( StoreQueryParameters.fromNameAndType( "session-context" , QueryableStoreTypes.keyValueStore())); SessionContext ctx = store.get(sessionId); // 다음 턴의 prompt에 끼워 넣기

다중 인스턴스로 배포하거나 Agent가 별도 서비스로 분리된 경우에는 상태가 분산되어 있다. 먼저 key에 따라 상태가 속한 인스턴스를 찾은 다음, 자체 구축한 RPC/REST 인터페이스로 읽어야 하며, Kafka Streams가 완전한 원격 쿼리 서비스를 대신 제공해 주지는 않는다. 쿼리는 또한 상태가 아직 준비되지 않았거나, 복구 중이거나, rebalance 중인 상황을 처리해야 한다. 입력 topic에 방금 쓴 직후에 조회해도 최신 턴을 읽을 수 있다고 보장되지 않는다.

따로 집고 넘어가야 할 함정이 하나 있다. 위의 merge().groupByKey() 토폴로지는 여러 소스를 병합했다고 해서 자동으로 파티션이 통일되지 않는다. groupByKey()는 능동적으로 재파티션을 트리거하지 않고, merge()도 각 소스의 key 분포를 통일하지 않는다. 여러 소스의 파티션 수, key 직렬화 방식 또는 파티션 전략이 일치하지 않으면, 같은 세션의 메시지가 서로 다른 task로 들어가 집계 상태가 불완전해질 수 있다.

배포 전에는 단순히 파티션 수를 세는 것이 아니라 실제 토폴로지와 파티션 전략을 확인해야 한다. 업스트림 일관성을 보장할 수 없을 때는 집계 전에 명시적으로 repartition()해야 한다. 대가는 내부 topic 하나와 읽기·쓰기 한 번이 추가되는 것이다. 또한 명시적 재파티션도 소스 간의 비즈니스 시간 순서를 자동으로 복원하지는 않는다.

이 경로가 누구에게나 필요한 것은 아니다. 단일 Agent에 트래픽도 많지 않다면 Postgres에 레코드 한 줄이면 충분하고, 억지로 Streams를 올리는 것은 스스로 운영 부담을 늘리는 일이다. 이것이 이득이 되는 전제는 시스템이 이미 Kafka를 이벤트 백본으로 삼고 있고, 양도 적지 않으며, 기억과 나머지 트래픽이 동일한 순서 보장, 재생 능력, ACL을 공유하기를 바라는 경우다.

혹시 이런 질문이 나올 수 있다. 그렇다면 벡터 검색에는 벡터 DB를 쓰는 게 더 맞지 않나?

두 가지 기억은 같은 것이 아니다. 벡터 DB는 의미 기반으로 유사한 내용을 찾는 데 능하며 지식베이스 Q&A에 적합하다. 세션 컨텍스트가 원하는 것은 「이 대화가 방금 무엇을 말했는가」에 대한 정밀한 위치 확인이며, key로 한 번 조회하면 충분하다. 실제로 RAG를 하려면 흔한 방식은 스트림 처리가 메시지의 embedding을 생성해 벡터 인덱스에 쓰고, Agent가 조회할 때 의미 검색을 하는 것이다. 벡터 검색은 일반적인 key join이 아니며, 스트림 안에서 join으로 되돌릴 수 있는지는 커넥터와 외부 시스템의 지원 상황에 달려 있다.

四、경로 3: 모델 추론을 스트림 처리로 들여오기, Flink SQL로 AI 함수 호출

앞의 두 경로는 Kafka가 AI를 섬기게 하는 것이었다. 이 경로는 반대다. 모델 추론을 스트림 처리의 링크 안으로 연결한다. 모델 자체는 여전히 외부 엔드포인트에서 실행되며, Flink는 호출을发起할 뿐이라는 점에 유의하자.

Confluent는 Confluent Cloud에서 이 능력 묶음을 제품으로 만들었고, 이를 통칭 Confluent Intelligence라고 부르며, 기반은 Flink SQL이다. 방식은 「모델 호출」을 SQL 안의 한 번의 호출로 추상화하는 것이다. 아래 고객센터 감정 태깅 예시는 공식 문서의 시그니처에 따라 작성했으며, Confluent Cloud에서 실측하지는 않았다.

-- AI_SENTIMENT는 현재 Early Access이며, 먼저 해당 능력을开通해야 함 -- 단일 감정 문자열이 아니라 구조화된 ROW를 반환 SELECT ticket_id, content, AI_SENTIMENT(content, ARRAY [ 'support' ]) AS sentiment_result FROM cs_user_turns;

시계열 이상 탐지를 하려면 타임스탬프와 watermark를 갖춘 시계열 테이블을 별도로 만들어 윈도우 문법으로 호출해야 한다.

-- server_metrics는 공식 요구에 따라 event_time과 watermark를 정의해야 함 SELECT event_time, response_ms, AI_DETECT_ANOMALIES( response_ms, event_time, JSON_OBJECT ( 'model' VALUE 'ttm' ) ) OVER ( ORDER BY event_time RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ) AS anomaly_result FROM server_metrics;

불리언 플래그가 필요할 때는 바깥쪽 SELECT에서 anomaly_result.is_anomaly 를 읽는다.

공식 문서에 있는 기존 함수로는 AI_FORECAST (시계열 예측), ML_PREDICT (원격 모델 추론), AI_COMPLETE (텍스트 생성), AI_EMBEDDING (벡터 생성), AI_TOOL_INVOKE (MCP 도구 호출)가 있으며, 그 외에 CREATE MODEL , CREATE AGENT , CREATE TOOL 세 가지 문을 갖춰 각각 모델, Agent, 도구를 정의하는 데 사용한다. 이러한 능력의 개방 상태는 각기 다르며, 구체적 파라미터는 공식 문서를 기준으로 하고 2차 블로그를 그대로 베끼면 안 된다.

Confluent Intelligence : Confluent Cloud에서 Flink SQL로 Agent 워크플로를 구축하는 능력 묶음으로, Streaming Agents, Real-Time Context Engine과 내장 AI/ML 함수를 포함한다. 「모델 호출」을 SQL 안의 한 번의 호출로 추상화한 것으로 이해할 수 있다.

여기에는 배울 만한 설계 발상이 하나 있다. 그것은 AI를 두 갈래로 나눈다. 하나는 스트림 안의 함수 호출로, 데이터가 어디에 도달하든 거기에 표시한다. 다른 하나는 Real-Time Context Engine(RTCE)으로, topic에서 데이터를 지속적으로 수집해 저지연 서빙 계층에 물질화(materialize)하고, 다시 MCP로 비즈니스 데이터를 개방한다. 공식 문서는 매우 직설적으로, Agent가 Kafka에 직접 연결하는 것이 아니라 MCP 도구를 통해 데이터를 조회하도록 적고 있다.

분명히 해 둘 것은, RTCE와 KIP-1318의 클러스터 운영 MCP Server, 그리고 Streams의 KTable 기억은 서로 다른 세 가지라는 점이다. 다만 발상적으로 모두 「통일된 도구 접근 + 상태 물질화」를 사용했을 뿐이다. MCP를 지원하는 Agent는 통일된 인터페이스로 비즈니스 컨텍스트를 얻을 수 있지만, 실시간이 제로 지연을 뜻하지는 않는다. 엔드투엔드 신선도는 여전히 소비와 처리 지연의 영향을 받는다.

대가도 솔직히 말하자. 이 경로는 관리형 서비스에 의존하며, 자체 구축 Flink라면 모델 엔드포인트를 스스로 연결해야 한다. 비용도 「메시지 한 건당 모델 비용 한 건」만이 아니다. 모델 추론이 오버헤드를 늘리며, 동시에 Flink의 연산량(Confluent의 Flink는 CFU로 계량), 선택한 모델의 요청 또는 token 사용량, 그리고 저장 및 전송 비용을 추산해야 하고, 실제 함수와 서비스의 과금 규칙에 따라 산정해야 한다. 그래서 나는 AI 함수를 「이 돈을 쓸 가치가 있는」 부분의 스트림에 사용할 것을 더 권한다. 예컨대 티켓 감정, 결제 이상 등이고, 나머지 스트림은 규칙 판단을 충실히 하게 두는 것이다.

한 가지는 반드시 따로 경고해야 한다. Confluent Cloud가 실제로 이러한 AI 능력을 제공하지만, 개방 범위와 성숙도는 기능별로 다르고 제한도 매우 구체적이다. 이 글에서 사용한 AI_SENTIMENT , AI_DETECT_ANOMALIES 는 현재 여전히 Early Access로 표기되어 있으며, 평가 및 비프로덕션 테스트에만 적합하다. RTCE의 공식 문서에는 그것이 AWS의 Basic, Standard, Enterprise, Dedicated 클러스터에서만 사용 가능하며, topic이 이미 schema에 등록되어 있어야 한다고 적혀 있다. 배포 전에는 클러스터 유형, 리전, 개통 조건을 항목별로 확인해야 한다.

五、세 경로를 어떻게 고를까

고민될 때 나는 먼저 이렇게 묻는다. 당신이 해결하려는 것이 「AI가 손을 쓸 수 있게」인가, 「AI가 기억하게」인가, 아니면 「AI가 신선한 컨텍스트를 얻고 실시간 처리를 하게」인가. 세 질문이 세 경로에 대응하며, 순서는 어느 것이 더 앞서 있는지가 아니라 비즈니스 문제가 결정해야 한다.

경로 해결하는 것 현재 상태 주요 비용 주요 위험 제1자 MCP Server(KIP-1318) Agent의 클러스터 조작 논의 중인 제안, 2026-09-16 기준 공식 딜리버리 미확인 공식 대기 + 우선 커뮤니티 구현 사용 권한 확대, 오조작 Streams 컨텍스트 저장 Agent 세션 기억 오늘 바로 자체 구축 가능 Streams 애플리케이션 한 세트의 운영과 상태 저장 파티션 또는 순서 불일치로 인한 집계 상태 불완전 Flink SQL AI 함수 실시간 태깅과 예측 기능 상태 각기 다름, 이 글의 두 함수는 Early Access Flink 연산 자원 + 해당 모델 비용 비용과 라벨 품질

남겨 둘 만한 경험은 딱 하나다. Agent에 클러스터 권한을 여는 일은 그것의 기억과 컨텍스트가 신뢰할 만한지 파악한 뒤에 하는 것이다. 어느 것을 먼저 올릴지는 비즈니스 문제에 달려 있다. 시스템이 이미 Kafka를 이벤트 백본으로 삼고 있고, 세션 상태를 읽어야 할 컴포넌트가 여럿인 팀이라면 경로 2의 이득이 가장 직접적이다. Agent를 먼저 운영에 참여시키고 싶다면 권한 경계를 먼저 설계해야 한다.

六、배포 전 네 가지

아래에서 언급하는 mcp.* 필드는 모두 KIP 초안에서 나온 것이며, 기존 커뮤니티 구현은 각자 대응하는 설정 항목을 사용하므로 필드명을 그대로 베끼면 안 된다.

첫째, 권한은 읽기 전용에서 시작한다. 읽기 전용 계정과 ACL을 사용해 조작 가능한 topic을 접두사 화이트리스트 안에 묶어 두고, 잘 돌아가면 하나씩 풀어 준다. (초안의 해당 스위치는 mcp.readonly이며 기본값은 false라서 명시적으로 켜야 한다.)

권한 개방 순서는 아래 그림과 같다.

둘째, 감사는 Kafka 자체의 topic에 남긴다. mcp.audit.topic은 한 번의 호출에 대한 신원, 도구, 파라미터, 결정, 결과를 기록하지만, Kafka의 추가(append) 쓰기 의미가 데이터가 영원히 삭제되지 않는다는 뜻은 아니다. retention, compaction, 관리 권한 모두 감사 보존에 영향을 준다. 올바른 방법은 감사에 독립적이고 제한된 쓰기 신원을 부여해, 그것이 감사 topic을 삭제하거나 보존 설정을 수정하거나 기록을 정리하지 못하게 하고, 적절한 보존 정책을 설정하는 것이다. 강한 감사 요구가 있으면 독립적인 영구 감사 시스템(예: WORM 또는 SIEM)에도 동기화해야 한다.

셋째, 파티션 전략을 확인한다. 세션 메모리는 세션 메시지와 같은 key, 같은 파티션 수, 같은 직렬화와 파티셔너를 써야 한다. 상위(업스트림) 일관성을 보장할 수 없을 때는 집계 전에 명시적으로 repartition()을 한다.

넷째, 비용을 명확히 계산한다. 모델 추론, Flink 연산량, 저장과 전송을 모두 포함해 계산하고, 먼저 일부 스트림으로 효과와 비용을 검증한 뒤 전체 적용 여부를 결정한다.

2026년 9월의 상태로 보면, Streams 세션 상태라는 이 경로는 오늘 바로 스스로 구축할 수 있고, 자사(퍼스트파티) Kafka MCP는 아직 논의 중이며, Confluent Cloud의 AI 기능은 개방 상태를 하나씩 확인해야 한다. 세 경로에는 통일된 선후 순서가 없고, 확실한 것은 단 하나다. 권한과 상태 정확성을 먼저 단단히 다진 다음, Agent가 프로덕션 클러스터를 움직이게 하라는 것이다.

저는 프로그래머 톈톈쿤(天天困)이며, 계속해서 프로그래밍 알짜 정보를 공유합니다. 유용하다고 느끼셨다면 좋아요, 즐겨찾기, 팔로우 부탁드립니다~ 댓글로도 이야기 나눠 주세요. Kafka에 AI를 접목한 뒤, 여러분은 Agent의 기억을 먼저 Streams에 두고 싶으신가요, 아니면 먼저 읽기 전용 클러스터 권한을 열어 주고 싶으신가요?

프로그래머 톈톈쿤

34

5.6k

조회

2

팔로워