Skip to content
Draft
4 changes: 4 additions & 0 deletions build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,9 @@ dependencies {
implementation 'org.springframework.boot:spring-boot-starter-security-oauth2-client'
implementation 'org.springdoc:springdoc-openapi-starter-webmvc-ui:3.0.3'
implementation 'org.springframework.boot:spring-boot-starter-flyway'
// Context→AI 요청 경로의 전달 수단(BD-48). 재시도 토픽 체인·DLT는 spring-kafka의
// @RetryableTopic을 쓴다. 버전은 Boot BOM이 관리한다.
implementation 'org.springframework.boot:spring-boot-starter-kafka'
implementation 'io.micrometer:micrometer-registry-prometheus'
implementation 'org.flywaydb:flyway-database-postgresql'
compileOnly 'org.projectlombok:lombok'
Expand All @@ -58,6 +61,7 @@ dependencies {
testImplementation 'org.springframework.boot:spring-boot-testcontainers'
testImplementation 'org.testcontainers:testcontainers-junit-jupiter'
testImplementation 'org.testcontainers:testcontainers-postgresql'
testImplementation 'org.testcontainers:testcontainers-kafka'
testCompileOnly 'org.projectlombok:lombok'
testRuntimeOnly 'org.junit.platform:junit-platform-launcher'
testAnnotationProcessor 'org.projectlombok:lombok'
Expand Down
28 changes: 28 additions & 0 deletions compose.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -33,5 +33,33 @@ services:
timeout: 3s
retries: 10

kafka:
# Context→AI 요청 큐(BD-48). 테스트 컨테이너(IntegrationContainerSupport)도 같은 태그를 쓴다.
# KRaft 단일 노드. 호스트(애플리케이션)는 HOST 리스너(19092)로, 컨테이너 내부는 9092로 붙는다 —
# Kafka는 접속 후 advertised listener로 재접속하므로 호스트용 주소를 따로 광고해야 한다.
image: apache/kafka:4.1.0
labels:
# Boot의 compose 자동 연결을 끄고 application-local.yml의 명시 주소를 쓴다. 자동 연결은
# 9092의 매핑 포트를 기대하는데 이 구성은 호스트 리스너가 19092라 서로 어긋난다.
org.springframework.boot.ignore: true
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@localhost:9093
KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093,HOST://:19092
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,HOST://localhost:${KAFKA_PORT:-19092}
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT,HOST:PLAINTEXT
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
ports:
- "${KAFKA_PORT:-19092}:19092"
healthcheck:
test: ["CMD-SHELL", "/opt/kafka/bin/kafka-broker-api-versions.sh --bootstrap-server localhost:9092 > /dev/null 2>&1"]
interval: 5s
timeout: 10s
retries: 10
start_period: 15s

volumes:
postgres-data:
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
# BD-17. 비동기 AI 처리를 메시지 큐 없이 DB State + Scheduler로

- **상태**: Accepted
- **상태**: Superseded by [BD-48](BD-48-context-ai-kafka-queue.md) — 요청 전달 수단(인메모리 큐→Kafka)에 한한 대체. DB State가 진실의 원본이라는 결정과 재스캔 안전망은 유지된다
- **날짜**: 2026-07-22 (AI 설계 확립 시점. 단일 커밋으로 특정 불가)
- **작성 시점**: 2026-07-27 — 결정 이후에 정리
- **관련**: S15P11A705-76
Expand Down
70 changes: 70 additions & 0 deletions docs/backend/decisions/BD-48-context-ai-kafka-queue.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
# BD-48. Context→AI 요청 경로의 인메모리 큐를 Kafka로 교체한다

- **상태**: Accepted
- **날짜**: 2026-08-04
- **관련**: S15P11A705-290 · [back#129](https://github.com/Team-PinLog/back/issues/129) · [BD-17](BD-17-async-without-message-queue.md)

## 맥락

Record 저장 커밋 후 FastAPI 호출은 `@Async` 리스너의 인메모리 큐(`aiCallExecutor`)를 탔다.
이 큐는 포화되면 요청을 버리고(재스캔이 다음 주기에 줍는다), 프로세스가 죽으면 담긴 힌트가
사라지며, 브로커 수준의 재시도·백오프·DLQ 같은 실패 격리 장치가 없다.

back#129에서 안 A(요청 방향 전환)·안 B(완료 이벤트 역방향)·안 C(back 내부에만 도입)를 비교했고,
AI 파트가 안 C 착수에 이견 없음을 밝히며 다섯 가지 질문을 남겼다. 핵심 지적은 정확하다 —
**유실 방지는 이미 DB State가 해결하고 있으므로, 이 경로에서 MQ가 실제로 개선하는 것은
"재스캔 주기(최대 5분)만큼의 복구 지연"이다.** 재스캔 주기 단축·큐 용량 확대 같은 더 싼
대안도 함께 제시됐다.

동시에 이 도입의 명분 절반은 처음부터 기술 경험·포트폴리오였고, back#129 본문과 AI 파트
코멘트 모두 "학습 가치는 그 자체로 팀이 도입을 결정할 수 있는 사유"로 인정했다.

## 선택지

| 안 | 장점 | 단점 |
|---|---|---|
| (a) 현상 유지 + 설정 조정(재스캔 주기 단축 등) | 비용 0에 가까움 | 복구 지연 개선뿐, 실패 격리·학습 목표는 얻지 못함 |
| (b) RabbitMQ | 작업 큐 패턴과 도구 성격이 가장 일치(ack·DLX 기본 제공) | 컨테이너 추가. 이력서·확장성 관점 매력은 Kafka보다 낮다고 판단 |
| (c) Redis Streams | 인프라 추가 0 | 백오프·DLQ 직접 구현 — "MQ를 배웠다"기보다 "Redis로 흉내"에 가까움 |
| (d) Kafka | 학습·포트폴리오 가치 최대. 안 B(완료 이벤트)를 넘어 이벤트 백본으로 키울 여지 | 이 워크로드에는 과투자. 개별 재전달이 없어 재시도 토픽 체인을 직접 구성. 운영 부담 최대 |

## 결정

**(d) Kafka를 채택한다. 능동적 선택이며, back#129가 열어 둔 "학습 목적으로 과투자를 감수한다"는
선택지를 집는 것이다.** 신뢰성 필수론으로 포장하지 않는다 — 아키텍처가 얻는 실익은 복구 지연
단축(5분 → 백오프 수 초)과 실패 격리(DLT)이고, 나머지는 학습 가치다.

구현 경계는 다음을 지킨다.

- **DB State(`ai.context_ai_state`)가 진실의 원본이라는 BD-17의 핵심은 유지한다.** Kafka는 전달
수단만 교체하며, 소비자의 멱등 가드도 브로커가 아니라 상태 행을 근거로 판정한다.
- **FastAPI와의 계약은 바꾸지 않는다.** 소비자는 재스캔과 같은 조립기·클라이언트로 같은 HTTP를
보낸다. AI 파트 코드 변경 없음(안 C의 정의).
- **재스캔은 최종 안전망으로 유지한다.** 브로커와 독립인 HTTP 직접 호출 경로다. 발행 실패·브로커
장애·DLT 격리 어느 경우에도 PENDING이 남아 재스캔·Finalizer가 최종 처분한다.
- 실패 분류: 5xx·연결 실패는 재시도 토픽 체인(`@RetryableTopic`, 지수 백오프), 4xx·역직렬화
실패는 DLT 직행. DLT 격리는 ERROR 로그로 남긴다.
- 발행은 커밋 후 리스너에서 하며 실패를 삼킨다(기존 의미론). `max.block.ms=1000`으로 브로커
장애가 요청 스레드를 붙잡는 상한을 건다.

back#129의 다섯 질문에 대한 답은 이 문서와 티켓 S15P11A705-290 본문이 겸한다 — ① 포화 실측
없음(예방+학습), ② 5분 지연은 시연 체감 문제, ③ 주기 단축 대비 추가 이득은 실패 격리와 학습,
④ 단일 인스턴스에서 브로커가 주는 것은 프로세스 독립 내구성, ⑤ DLQ 관측은 ERROR 로그가
최소선이고 지표는 `/metrics` 승인 뒤 후속.

## 결과

- 이 결정으로 감수하는 것:
- **재시도 주체가 둘이 된다.** 체인(수 초 단위 순간 장애)과 재스캔(그 밖의 모든 유실, 예산
`retry_count`)의 역할을 위처럼 갈랐고, FAILED 종결 권한은 여전히 재스캔·Finalizer에만 있다.
- **INFRA에 브로커 운영 부담이 생긴다.** dev/prod 배포는 INFRA 소유이며 이 작업의 머지
선행 조건이다.
- **DLT 관측이 로그뿐이다.** 지표·알림은 `/metrics` prod 승격 승인(infra `docs/ai-serving.md`
검증 7) 이후 후속 작업.
- 공용 계약 `10_MVP_기능범위` §2(MQ 제외)와 어긋난다 — back#129 합의로 대체하며 docs 레포
갱신이 후속으로 필요하다.
- 재검토 트리거(이 조건이 오면 다시 논의):
- 안 B(FastAPI 완료 이벤트 역방향)를 진행하게 되면 — 이벤트 스키마·발행 실패 정책을 AI 파트와
합의하고, 이 브로커를 그 백본으로 쓸지 결정한다.
- 브로커 운영 부담이 실익을 넘어서면 — 안 (a)로 회귀하는 것도 열려 있다. 재스캔이 그대로
있으므로 회귀 비용은 발행·소비 코드 제거뿐이다.
26 changes: 26 additions & 0 deletions docs/backend/worklog/2026-08-04-S15P11A705-290-context-ai-kafka.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
# Context→AI 요청 경로의 인메모리 큐를 Kafka로 교체

- **날짜**: 2026-08-04
- **추적**: S15P11A705-290
- **관련**: [BD-48](../decisions/BD-48-context-ai-kafka-queue.md) · [back#129](https://github.com/Team-PinLog/back/issues/129) · [BD-17](../decisions/BD-17-async-without-message-queue.md)

back#129의 안 C를 구현했다. 커밋 후 리스너가 `@Async` 인메모리 큐 대신 Kafka 토픽에 발행하고,
back 내부 컨슈머가 소비해 기존과 동일한 FastAPI 호출을 보낸다. FastAPI 계약과 재스캔 안전망은
그대로다. 브로커를 Kafka로 정한 근거(학습 가치로 과투자를 감수)와 그때 검토한 트레이드오프,
AI 파트의 다섯 질문에 대한 답은 BD-48에 있다.

내린 판단들:

- 발행 실패는 기존처럼 삼킨다. 대신 `max.block.ms=1000`으로 브로커 장애가 요청 스레드를 붙잡는
상한을 걸었다 — 커밋 후 리스너가 요청 스레드에서 돌기 때문이다.
- 실패 분류를 예외 타입으로 나눴다. 5xx·연결 실패는 `AiProcessRetryableException`(재시도 토픽
체인), 4xx·역직렬화 실패는 `AiProcessFatalException`(DLT 직행). `AiProcessClient`의 삼키는
`process()`는 재스캔용으로 남기고, 던지는 `processOrThrow()`를 소비자용으로 추가했다.
- 멱등 가드는 브로커가 아니라 상태 행 기준이다 — 두 단계 모두 PENDING을 지난 Context는 생략한다.
- `aiCallExecutor`·`@EnableAsync`는 이 리스너 전용이었으므로 함께 제거했다.
- 5xx·연결 끊김 통합 테스트는 재시도 체인이 생기며 호출이 3회로 늘었다. 각 테스트가 자기
재시도를 끝까지 소비하도록 고쳤다 — 남기면 뒤 테스트의 "호출 없음" 검증 구간에 흘러든다.

머지 선행 조건: INFRA의 Kafka 브로커 배포(dev). AI 파트 소유 명세
`docs/ai/spec/ai-integration.md` 4장(인메모리 큐 서술)과 공용 계약 `10_MVP_기능범위` §2(MQ 제외)가
이 변경과 어긋난다 — 각각 갱신 요청과 docs 레포 후속 PR이 필요하다.
Original file line number Diff line number Diff line change
@@ -1,38 +1,44 @@
package com.pinlog.pinlogback.domain.ai;

import java.util.concurrent.Executor;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.kafka.clients.admin.NewTopic;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.http.client.SimpleClientHttpRequestFactory;
import org.springframework.kafka.config.TopicBuilder;
import org.springframework.scheduling.TaskScheduler;
import org.springframework.scheduling.annotation.EnableAsync;
import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
import org.springframework.web.client.RestClient;

import com.pinlog.pinlogback.domain.ai.queue.AiQueueProperties;
import com.pinlog.pinlogback.domain.ai.service.AiRescanProperties;

/**
* FastAPI 연동에 필요한 인프라 조립. 설정 클래스를 {@code global/config}가 아니라 소비자와 같은
* 패키지에 두는 기준은 {@code docs/development/package-structure.md}의 보안 설정과 같다 — 설정과
* 그 설정이 조립하는 구현이 떨어져 있으면 한쪽만 고치게 된다.
*
* <p>{@link EnableScheduling}이 여기 있는 것은 {@link EnableAsync}와 같은 사정이다. 둘 다 애플리케이션
* <b>전역</b> 스위치인데, 켜야 하는 이유가 이 연동에만 있다. {@code global/config}로 올리면 스위치와
* 그 스위치의 유일한 소비자가 떨어져 앉는다.
* <p>{@link EnableScheduling}이 여기 있는 것도 같은 사정이다. 애플리케이션 <b>전역</b> 스위치인데,
* 켜야 하는 이유가 이 연동(재스캔)에만 있다. {@code global/config}로 올리면 스위치와 그 스위치의
* 유일한 소비자가 떨어져 앉는다. 한때 여기 있던 {@code @EnableAsync}와 {@code aiCallExecutor}
* (커밋 후 FastAPI 호출용 인메모리 큐)는 Kafka 발행으로 대체되어 사라졌다(BD-48).
*/
@Configuration
@EnableAsync
@EnableScheduling
@EnableConfigurationProperties({AiProperties.class, AiRescanProperties.class})
@EnableConfigurationProperties({AiProperties.class, AiRescanProperties.class, AiQueueProperties.class})
public class AiIntegrationConfig {

private static final Logger log = LoggerFactory.getLogger(AiIntegrationConfig.class);
/**
* Context→AI 요청의 본 토픽(BD-48). 재시도 토픽과 DLT는 {@code @RetryableTopic}이 여기서
* 파생해 만들므로 본 토픽만 선언한다. 파티션 1인 이유: 처리량이 Record 저장 빈도(초당 수 건)를
* 넘지 않고, 컨슈머도 단일 인스턴스라 병렬화로 얻을 것이 없다. 복제 1은 단일 브로커 전제다 —
* 브로커 구성이 커지면 INFRA와 함께 올린다.
*/
@Bean
public NewTopic contextAiProcessTopic(AiQueueProperties queueProperties) {
return TopicBuilder.name(queueProperties.topic()).partitions(1).replicas(1).build();
}

/** {@code process}용 전용 인스턴스(AI 파트 소유 명세 {@code docs/ai/spec/ai-integration.md} 2·3장). */
@Bean
Expand Down Expand Up @@ -102,26 +108,4 @@ public TaskScheduler taskScheduler() {
return scheduler;
}

/**
* AI 호출 전용 풀(명세 4.2). 공용 executor를 쓰지 않는 이유는 FastAPI 장애가 다른 비동기 작업까지
* 굶기지 않게 하기 위해서다.
*
* <p><b>큐가 차면 버린다.</b> {@code CallerRunsPolicy}를 쓰면 요청 스레드가 외부 호출을 대신
* 수행해 FastAPI 장애가 그대로 Core 처리량 저하가 된다. 버려진 요청은 {@code PENDING}으로 남아
* 재스캔 대상이 되므로 유실이 아니다 — 그래서 버리는 쪽이 안전하다.
*/
@Bean
public Executor aiCallExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setThreadNamePrefix("ai-call-");
executor.setCorePoolSize(2);
executor.setMaxPoolSize(4);
executor.setQueueCapacity(100);
// 버린 사실은 남긴다. 조용히 버리면 "왜 PENDING만 쌓이는가"를 알 길이 없다. 예외를 던지는
// AbortPolicy는 쓰지 않는다 — 커밋 이후 리스너를 타고 올라가 정상 응답을 오류로 만든다.
executor.setRejectedExecutionHandler((task, pool) ->
log.debug("AI 호출 큐 포화로 요청을 버렸다. PENDING이 남아 재스캔이 다시 집는다."));
executor.initialize();
return executor;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,8 @@
import org.springframework.web.client.RestClientResponseException;

import com.pinlog.pinlogback.domain.ai.AiProperties;
import com.pinlog.pinlogback.domain.ai.exception.AiProcessFatalException;
import com.pinlog.pinlogback.domain.ai.exception.AiProcessRetryableException;
import com.pinlog.pinlogback.global.web.TraceIdFilter;

/**
Expand All @@ -23,11 +25,12 @@
* 통보용 웹훅·콜백이 없으므로 이 클래스는 응답 본문을 읽지 않는다. 정합성의 근거는 이 호출이
* 아니라 DB에 영속된 {@code ai.context_ai_state}다 — 이 호출은 "지금 처리하면 조금 빨라지는 힌트"다.
*
* <p><b>모든 실패를 삼킨다.</b> 예외를 밖으로 던지면 호출자(커밋 이후 리스너)를 타고 올라가
* 사용자 응답을 오류로 만들 수 있는데, Core 데이터는 이미 커밋되어 정상이므로 그것은 거짓말이다.
* 상태를 {@code FAILED}로 바꾸지도, {@code retry_count}를 올리지도 않는다 — 호출 실패는 FastAPI
* 내부 작업의 실패가 아니고, 두 주체가 같은 사유로 FAILED를 각각 기록하면 원인 추적이 불가능해진다.
* {@code PENDING}이 남아 있으므로 재스캔이 같은 Context를 다시 집는다.
* <p>실패의 취급이 호출자마다 다르므로 진입점이 둘이다. {@link #process}는 <b>모든 실패를
* 삼킨다</b> — 재스캔 경로에서 예외는 아무것도 복구하지 못하고, {@code PENDING}이 남아 다음
* 회차가 같은 Context를 다시 집는다. {@link #processOrThrow}는 실패를 분류해 던진다 — 큐 소비
* 경로에서는 던지는 것이 곧 재시도 체인·DLT 격리의 신호다(BD-48). 어느 쪽도 상태를
* {@code FAILED}로 바꾸거나 {@code retry_count}를 올리지 않는다 — 호출 실패는 FastAPI 내부
* 작업의 실패가 아니고, 두 주체가 같은 사유로 FAILED를 각각 기록하면 원인 추적이 불가능해진다.
*/
@Component
public class AiProcessClient {
Expand Down Expand Up @@ -75,6 +78,20 @@ private static String requireSecret(String secret, Environment environment) {

/** 호출 결과를 반환하지 않는다. 호출자가 분기할 수 있으면 그 분기가 곧 Core 트랜잭션 결과에 스며든다. */
public void process(ContextProcessRequest request) {
try {
processOrThrow(request);
} catch (AiProcessRetryableException | AiProcessFatalException e) {
// 로그는 processOrThrow가 이미 남겼다. 재스캔 경로에서 실패는 여기서 끝난다 —
// PENDING이 남아 있으므로 다음 회차가 같은 Context를 다시 집는다.
}
}

/**
* 실패를 분류해 던진다. 5xx·연결 계열은 {@link AiProcessRetryableException}(다시 보내면 성공할
* 수 있다), 4xx는 {@link AiProcessFatalException}(요청 자체의 문제라 몇 번을 보내도 같다).
* 401·403도 fatal이다 — 시크릿·헤더 설정이 고쳐지기 전에는 재시도가 전부 같은 답을 받는다.
*/
public void processOrThrow(ContextProcessRequest request) {
String requestId = currentRequestId();
try {
restClient.post()
Expand All @@ -86,12 +103,20 @@ public void process(ContextProcessRequest request) {
.toBodilessEntity();
log.debug("AI process 접수됨: contextId={}, requestId={}", request.contextId(), requestId);
} catch (RestClientResponseException e) {
// 4xx는 요청 payload 형식 문제일 가능성이 있어 사람이 봐야 한다. 5xx는 상대 장애이므로
// 재스캔이 흡수한다. 응답 본문은 남기지 않는다 — 내부 API라도 로그로 새어 나갈 이유가 없다.
// 4xx는 요청 payload 형식 문제일 가능성이 있어 사람이 봐야 한다. 응답 본문은 남기지
// 않는다 — 내부 API라도 로그로 새어 나갈 이유가 없다.
logByStatus(e, request, requestId);
if (e.getStatusCode().is4xxClientError()) {
throw new AiProcessFatalException(
"AI process 호출이 " + e.getStatusCode().value() + "로 거절됐다: contextId=" + request.contextId(), e);
}
throw new AiProcessRetryableException(
"AI process 호출이 " + e.getStatusCode().value() + "를 받았다: contextId=" + request.contextId(), e);
} catch (RuntimeException e) {
log.warn("AI process 호출 실패(연결·타임아웃): contextId={}, requestId={}, cause={}",
request.contextId(), requestId, e.toString());
throw new AiProcessRetryableException(
"AI process 호출이 연결 단계에서 실패했다: contextId=" + request.contextId(), e);
}
}

Expand Down
Loading
Loading