Kafka 이벤트 스트림을 Claude로 연결하기: 이벤트 기반 아키텍처와 LLM 연동 패턴
요약
Kafka 이벤트 스트림을 Claude를 이용한 LLM 기반 처리 파이프라인에 연결하는 방법을 다룹니다. Kafka 컨슈머가 이벤트를 가져와 Claude API 또는 Agent SDK 등을 통해 처리하고, 결과를 다른 토픽으로 발행하는 패턴을 제시합니다. 핵심은 LLM을 일반적인 데이터 처리 단계로 간주하여 이벤트 기반 아키텍처(EDA) 설계 원칙을 적용하는 것입니다.
핵심 포인트
- Kafka는 Claude를 직접 호출하지 않으며, 커스텀 컨슈머가 이벤트를 가져와야 합니다.
- 처리 결과는 반드시 다른 Kafka 토픽으로 발행해야 파이프라인이 완성됩니다.
- LLM의 비결정성과 느린 속도를 고려하여 Idempotency 및 DLQ 처리가 필수적입니다.
- 연결 방식은 Claude API, `claude -p`, Agent SDK 등 4가지 선택지가 있습니다.
‘Kafka의 이벤트를 Claude에 연결할 수 있을까?’라는 질문에서 시작되었습니다. 주문이나 상태 변경, Slack 발언 등이 Kafka로 흘러들어오는 환경에서, 그 이벤트를 트리거로 Claude(또는 Claude Code로 만든 에이전트)를 구동하고 싶다는 이야기입니다. 이벤트 기반 파이프라인에 LLM의 판단을 한 단계로 삽입하고 싶다는 것입니다.
먼저 결론부터 말하자면, Kafka가 Claude를 ‘호출’하지 않습니다. Kafka의 컨슈머(Consumer)를 직접 작성하여 거기서 Claude를 호출하는 것이 기본 형태입니다. Kafka에는 Claude용 소스 커넥터도 싱크(Sink)도 없기 때문에, ‘이벤트를 가져오는 → Claude로 처리시키는 → 결과를 다음 토픽으로 흘려보내는’ 이 연결고리 역할을 자신의 코드로 구현해야 합니다. 사소해 보이지만, 여기만 잡으면 나머지는 조합의 문제입니다.
본문에서는 그 연결고리를 만드는 방법, 연결 패턴 선택법, 그리고 이벤트 기반 아키텍처 특유의 함정 순서로 정리하겠습니다. 솔직히 말하자면, Kafka 브로커를 세우고 엔드투엔드로 프로덕션 운영까지 통과시킨 기록은 아닙니다. 각 명령어와 API는 공식 문서를 확인한 후, 바로 적용할 수 있는 형태로 구현했습니다.
먼저 3줄 요약입니다.
- Kafka가 호출하는 것이 아니다.
컨슈머(직접 작성)가 이벤트를 가져와 Claude를 호출하고, 결과를 다른 토픽으로 흘려보낸다. LLM을 파이프라인의 한 단계로 사용한다 -
연결은 4가지 선택지가 있다. 경량 순서: Claude API → claude -p (헤드리스) → Agent SDK → Managed Agents. 수행하려는 내용에 따라 선택해야 한다. 핵심은 이벤트 기반 측면의 설계다.
성공한 후에 오프셋을 커밋(commit), Idempotency 확보, 레이트 제한에 따른 트래픽 제어, 실패는 DLQ(Dead Letter Queue)로 처리한다. LLM은 느리고 비결정적이므로, 이 부분을 소홀히 하면 시스템이 망가진다.
| 용어 | 출처 | 대략적인 의미 |
|---|---|---|
| 컨슈머 / 컨슈머 그룹 | Kafka | 토픽에서 이벤트를 가져오는 쪽. 그룹 단위로 병렬 분담한다 |
| ... | claude -p | Claude Code |
| Agent SDK | Anthropic | Claude Code와 동일한 루프/툴을 코드에서 사용하는 SDK (query() 등) |
| Managed Agents | Anthropic | 에이전트 루프와 샌드박스를 Anthropic이 호스팅하는 방식 |
| 구조화 출력 | Claude | 스키마에 맞춰 JSON으로 응답하게 하는 메커니즘. 다음 단계가 파싱하기 쉽다 |
우선 그림으로 이해해 봅시다. 업무 시스템에서 Kafka의 events 토픽으로 데이터를 흘려보내고, 컨슈머가 하나씩 가져와 Claude로 처리시키고, 결과를 results 토픽으로 발행합니다. 실패한 이벤트는 events.DLQ로 격리됩니다.
핵심은 LLM을 ‘특별한 것’이 아니라, 입력 토픽과 출력 토픽 사이에 놓인 일반적인 처리 단계로 다루는 것입니다. 이렇게 생각하면 기존의 이벤트 기반 설계가 그대로 적용될 수 있습니다.
컨슈머에서 Claude를 호출하는 방법은 크게 4가지가 있으며, 수행하려는 내용에 따라 선택합니다.
| 패턴 | 호출 방식 | 적합한 경우 |
|---|---|---|
| ① Claude API | Messages API 직접 사용 | 분류/요약/추출/라우팅 판단. 가장 가볍고 빠르며 저렴하다. 파일이나 툴은 필요 없다 |
② claude -p | Claude Code 헤드리스 실행 | 코드 수정, 조사, 리포지토리 조작. 툴과 파일 편집이 필요할 때. subprocess로 간편하게 호출 가능하다 |
| ③ Agent SDK | query() / ClaudeSDKClient | 상주 에이전트 서비스. ②와 동일한 루프를 subprocess 없이 코드 내부에서 돌리고 싶을 때 |
| ④ Managed Agents | sessions.create() | 부모-자식 에이전트, 승인(Approval), 메모리 포함하여 루프와 샌드박스를 Anthropic에 호스팅하고 싶을 때 |
어디서부터 생각해야 할지 모를 때 ①에서 접근합니다. '이 이벤트가 버그인지 사양인지', '어떤 팀에 넘겨야 하는지'와 같은 판단만 필요하다면, Claude API의 구조화된 출력(structured output)으로 충분하며 빠르고 저렴합니다. 실제 리포지토리를 조작하게 만들고 싶다면 ② 또는 ③을 사용하고, 클라우드에서 상태를 가지는 무거운 자동화가 필요하다면 ④ 순서입니다. ②와 ③은 내부 내용이 같은 에이전트 루프(agent loop)이며, CLI로 호출할지 라이브러리로 호출할지의 차이일 뿐입니다.
가장 간편한 방법은 컨슈머에서 claude -p를 subprocess로 호출하는 형태입니다. Claude Code의 헤드리스 실행(headless execution)은 -p에 프롬프트를 전달하면 결과를 출력하고 종료합니다. 공식 문서에 근거한 핵심 내용은 이 정도입니다.
claude -p "프롬프트"로 비대화형 실행--output-format json으로 JSON (본문은result,session_id나total_cost_usd도 포함)- 스키마를 고정하고 싶다면
--output-format json --json-schema '...'(structured_output결과) - 표준 입력을 읽을 수 있어서
cat data | claude -p "..."가 사용 가능합니다 (최대 10MB). - CI나 스크립트에서는
--bare가 권장됩니다 (훅/MCP/CLAUDE.md 등의 자동 로드를 건너뛰고 매번 동일한 결과).--bare를 사용할 때는ANTHROPIC_API_KEY가 필요합니다. - 사용하게 할 도구는--allowedTools "Read,Edit,Bash", 권한은--permission-mode, 무인 실행은--permission-prompts none입니다.
Python 컨슈머에서 호출하는 기본 골격입니다 (이벤트 페이로드를 표준 입력으로 전달하는 예시).
import json, time, subprocess
from kafka import KafkaConsumer, KafkaProducer
consumer = KafkaConsumer(
...
핵심은 두 가지입니다. publish()를 통해 전송의 확정(future)을 기다린 후에 커밋해야 한다는 것 (기다리지 않고 커밋하면, 전송 전에 실패했을 때 이벤트가 사라집니다). 그리고 일시적인 오류(타임아웃)는 재시도하고, 모두 소진되었을 때만 DLQ로 보내야 한다는 것입니다 (즉시 DLQ로 보내면 복구 가능한 이벤트까지 버리게 됩니다). 이 형태라면, 기존의 .claude/agents나 skills를 활용하고 싶을 때도 --bare를 빼면 그대로 로드할 수 있습니다 (이 경우 작업 디렉토리 설정에 주의해야 합니다).
②와 같은 작업을 subprocess 없이 코드 내에서 돌리고 싶다면 Agent SDK입니다. Claude Code와 동일한 도구/루프/컨텍스트 관리를 Python / TypeScript에서 직접 사용할 수 있습니다. 진입점은 query() (1회성 세션)이고, 대화를 지속하고 싶다면 ClaudeSDKClient를 사용합니다.
# pip install claude-agent-sdk
from claude_agent_sdk import query
async def handle(event: dict) -> str:
...
상시 서비스로 작성할 것이라면 ③이 가장 자연스럽습니다. 도구 승인을 콜백으로 삽입하거나, 메시지를 타입으로 처리하는 등의 제어도 가능합니다. ②와 ③은 같은 기반(Agent SDK)이기 때문에, 'CLI로 호출할지, 라이브러리로 호출할지'를 선택하면 충분합니다.
파일 조작이 필요 없는 '판단・분류・요약・추출'의 경우, Messages API를 직접 호출하는 것이 가장 가볍습니다. 구조화된 출력으로 반환하게 하면 다음 토픽으로 바로 흐르게 할 수 있습니다. 대량 이벤트를 한 번에 처리하려면 Batch API도 선택할 수 있습니다.
import anthropic
client = anthropic.Anthropic()
def classify(event: dict) -> dict:
...
패턴 ④ (Managed Agents)에 대해서는, 승인・메모리・부모 에이전트까지 클라우드에 호스팅하고 싶을 때, 컨슈머에서 sessions.create()를 호출하는 형태가 됩니다. 이 부분은 별도 글로 자세히 다루었습니다.
연결 자체는 어렵지 않습니다. 어려운 것은 느리고 비결정적이며, 가끔 거부(refusal)까지 하는 LLM을 Kafka 처리 단계에 안전하게 수용하는 것입니다. 이 부분을 대충 하면 중복 처리, 누락, 비용 폭발이 발생합니다.
오프셋은 성공한 후에 커밋한다. 자동 커밋을 끄고, Claude의 처리가 끝나고 결과 발행이 완료된 후에 커밋합니다. 도중에 다운되어도 재처리가 되므로, 최소한 한 번은 처리되는(at-least-once) 형태가 됩니다.
멱등하게 만든다. at-least-once이기 때문에 같은 이벤트가 두 번 올 수 있습니다. LLM의 출력은 매번 같지 않을 수 있으며, 게다가 주문과 같은 부작용을 동반할 수 있습니다. 이벤트 키로 중복 제거하고, 다운스트림(하류) 시스템은 멱등하게 설계합니다. 주문처럼 되돌릴 수 없는 작업은 이전 Managed Agents 기사에서 언급된 '승인 게이트'를 거치는 것이 안전합니다.
레이트 제한에 흐름을 맞춘다. Kafka는 한꺼번에 몰려오지만, Claude API에는 RPM/TPM 상한이 있습니다. 컨슈머 측에서 병렬 수를 줄이고, 429가 발생하면 retry-after에 따라 백오프합니다(SDK는 기본적으로 재시도합니다). Kafka의 체류 시간으로 흡수할 수 있는 것이 이벤트 기반 아키텍처의 강점입니다.
긴 처리는 poll을 멈추지 않는다. LLM 호출은 몇 초에서 수십 초가 걸립니다. 컨슈머가 한 건에 너무 많은 시간을 쓰면 max.poll.interval.ms를 초과하여 리밸런싱이 발생합니다. 무거운 처리는 워커 풀로 분산시키거나, 인터벌을 조정하는 등의 조치가 필요합니다.
실패와 거부는 DLQ(Dead Letter Queue)로. 파싱할 수 없거나, 타임아웃, stop_reason: "refusal" 등은 이것들로 멈추지 않고, 전용 토픽(DLQ)으로 퇴거하여 다음 단계로 진행합니다. 나중에 모아서 사람이 보거나, 재투입하는 운영 방식을 적용할 수 있습니다.
비용으로 진입점을 좁힌다. 모든 이벤트를 LLM에 통과시키면 비용이 많이 듭니다. 저렴한 규칙(정규 표현식이나 플래그)으로 '진짜로 LLM이 필요한 이벤트'만 통과시키고, 나머지는 그냥 통과시킵니다. 진입점의 축소가 그대로 청구 금액에 영향을 미칩니다.
우선 하나를 통하게 하려면, Docker로 단일 노드의 Kafka를 띄우고, 토픽을 만들고, 위의 Python 컨슈머를 돌리는 것이 빠릅니다.
# 1) 단일 노드 Kafka (이미지/태그는 최신 확인)
docker run -d --name kafka -p 9092:9092 apache/kafka:latest
# 2) 토픽 만들기 (컨테이너 내 CLI 사용 예)
...
여기까지 통하면,あとは '2)의 토픽에 이벤트를 던지고, 컨슈머가 results에 쓰는 것'을 확인하는 것만 남습니다. Kafka 이미지 이름이나 태그, CLI 경로는 버전에 따라 바뀌므로, 직접 공식 최신 정보를 확인해 주세요.
Kafka와 연결하면 긴장했지만, 실제로 해보니 구조는 단순했습니다. Kafka는 이벤트의 백본(Backbone), Claude는 판단/처리 레이어, 그 사이를 컨슈머가 연결합니다. LLM을 특별 취급하지 않고, 입력 토픽과 출력 토픽 사이에 놓인 한 단계로 봅니다. 그렇게 인식하는 순간, 오프셋, 멱등성, DLQ라는 이벤트 기반 아키텍처에서 쌓아온 설계 지식이 그대로 무기가 되었습니다.
어려움은 LLM의 특성(느림・비결정적・가끔 거부함)을 처리 단계에 어떻게 수용하느냐에 집중됩니다. 그 부분은 이번에 언급한 몇 가지 방법으로 대부분 해결할 수 있습니다. Devin 같은 클라우드 에이전트도, 결국 '이벤트가 오면 누군가가 작동해서 결과를 반환한다'는 메커니즘입니다. Kafka는 이 '작동하는' 계기를 업무의 모든 장소에서 공급해 주는 도관이 됩니다. 다음에는 이 다리 놓기 작업에 ④의 Managed Agents를 삽입하여, 주문 승인까지 이벤트 기반으로 한 번 쭉 연결해 볼 생각입니다.
Run Claude Code programmatically (headless /
claude -p
) (-p
・--output-format
・--json-schema
・--bare
・--allowedTools
) - Agent SDK overview (query() / ClaudeSDKClient ・Python/TypeScript) - Tool use / structured outputs (Claude Docs)(Messages API와 구조화 출력)
- Apache Kafka 공식 문서(컨슈머・오프셋・
max.poll.interval.ms)
AI 자동 생성 콘텐츠
본 콘텐츠는 Qiita AI의 원문을 AI가 자동으로 요약·번역·분석한 것입니다. 원 저작권은 원저작자에게 있으며, 정확한 내용은 반드시 원문을 확인해 주세요.
원문 바로가기