본문 바로가기
IT/트러블슈팅

[Trouble Shooting] 해시태그-카테고리 유사도 판단 파이프라인

by nohumb 2026. 9. 23.

아키텍처 다이어그램

개요

해시태그가 새로 생성되면(HashtagCreatedEvent), 기존 카테고리들과의 유사도를 계산해 즉시 병합 / 관리자 승인 대기 / 신규 카테고리 생성 중 하나로 연결하는 파이프라인.

  • 진입점: HashtagCreatedEventConsumer → CategoryHashtagLinkFacade.tryLink(hashtagId)
    • 실제 DB 반영(연결/신규 카테고리 승격)은 CategoryHashtagLinkWriter가 담당 - Facade는 조회+파이프라인 호출까지만
  • 유사도 계산 방식(Levenshtein/임베딩/LLM)은 언제든 교체·추가·삭제 가능하도록 설계됨

진입 이벤트 발행 (Outbox 패턴)

HashtagCreatedEvent/CategoryCreatedEvent는 생성 즉시 Kafka로 발행되는 게 아니라 Outbox 패턴을 거침 - dual-write(DB 커밋은 성공했는데 발행만 유실되는 문제) 방지 목적

  • 해시태그/카테고리 생성과 같은 트랜잭션 안에서 OutboxRepository에 이벤트를 먼저 기록(PENDING)
  • OutboxRelay(3초 주기, BATCH_SIZE=10, MAX_ATTEMPT=5)가 별도로 폴링해서 실제 Kafka 발행 - 즉 "DB에 생성됨"과 "이벤트가 실제로 발행됨" 사이에 최대 3초 갭이 하나 더 있음(파이프라인 전체 e2e 지연 계산 시 빠뜨리면 안 됨)
  • Outbox 자체의 알려진 갭(코드에 TODO로 남음):
    • 릴레이가 여러 인스턴스로 뜨는 경우 claim 쿼리(FOR UPDATE SKIP LOCKED) 없이는 중복 발행 가능성 있음
    • 발행 처리 중 죽으면 PROCESSING 상태에 멈춘 이벤트를 복구하는 로직이 없음
  • 매칭되는 publisher가 없는 이벤트 타입은 재시도 없이 즉시 실패 처리(데드레터)

파이프라인 구성 요소

스테이지 (Strategy 패턴)

  • 포트
// 유사도 판단 파이프라인의 스테이지, 기술적인 구현은 추상화
// candidates로 받은 병합 후보 각각에 대해 판정을 내려 반환한다(후보 하나당 결과 하나)
public interface CategorySimilarityStage {

    List<CategoryCandidateResult> evaluate(Hashtag hashtag, List<Category> candidates);
}
  • 구현체
    • LevenshteinCategorySimilarityStage(@Order(1), 문자열 거리, 빠름)
    • EmbeddingCategorySimilarityStage(@Order(2), 벡터 코사인 유사도)
      • 실제 임베딩 계산은 EmbeddingClientImpl이 Feign(EmbeddingFeignClient)으로 embedding-service를 호출
    • LlmCategorySimilarityStage(@Order(3))
      • 스켈레톤만 존재, 클래스 전체가 주석 처리돼 비활성 상태 (Bean으로 등록 안 됨)
  • 실행 순서: 각 구현체의 @Order로 결정 — 빠르고 단순한 것 → 느리고 정확한 것 순
  • 임계값은 Stage마다 상수(MERGE_THRESHOLD, PENDING_APPROVAL_THRESHOLD)를 따로 둠
  • 파이프라인은 스테이지의 기술적 구현(알고리즘/외부 호출 여부 등)에 전혀 관여하지 않음

카테고리 벡터 계산 파이프라인 (Embedding 스테이지 전제조건)

EmbeddingCategorySimilarityStage가 코사인 유사도를 계산하려면 카테고리마다 벡터가 미리 계산돼 있어야 함. (해시태그 벡터와 계산 방식이 비대칭적)

해시태그 벡터 카테고리 벡터

계산 시점 Embedding 스테이지 실행 중 lazy(그 자리에서) 카테고리 생성 시 비동기(Kafka 이벤트)
트리거 EmbeddingCategorySimilarityStage.ensureHashtagVectorExists() CategoryCreatedEventConsumer → CalculateCategoryVectorUseCase#calculate(categoryId)
실패 시 그 호출 자체가 Failed로 즉시 반영 Kafka 컨슈머 예외 → 재시도(3회) → 실패하면 CategoryVectorRecoveryScheduler(10초 주기, BATCH_SIZE=20, findActiveIdsWithoutVector)가 재계산
  • 카테고리는 admin API로 직접 생성되든, promoteToNewCategory()로 해시태그가 신규 승격되든 상관없이 생성 시 항상 CategoryCreatedEvent가 발행됨(Outbox) → 이 이벤트가 벡터 계산을 트리거
  • 막 생성된 카테고리는 벡터 계산이 아직 안 끝났을 수 있음 - 그 상태에서 다른 해시태그가 이 카테고리를 후보로 Embedding 스테이지에서 평가하면 Failed("카테고리 벡터 없음")가 남 (일시적 상태, CategoryVectorRecoveryScheduler가 해소)

결과 타입 (Hierarchy)

CategoryMatchResult (sealed interface, domain.model)

타입 의미

Merge 즉시 병합 가능할 만큼 유사함
PendingApproval 병합 가능성은 있으나 애매함 → 관리자 승인 대기
NotSimilar 이 스테이지 기준으론 유사한 카테고리 없음 (정상 판단)
Failed 스테이지 실행 자체가 실패함 (타임아웃/예외 등, 비정상 — "유사하지 않음"과는 다름)

결과 결합 규칙 (도메인 규칙)

CategoryMatchResult.combineWith(next) — 우선순위 규칙 자체를 타입에 캡슐화

  • Merge — 찾는 즉시 확정 (최우선, 이후 스테이지 호출 안 함)
  • PendingApproval — 한 번 찾으면 이후 어떤 결과가 와도(Merge 제외) 계속 유지됨
  • NotSimilar / Failed — 서로 우선순위 없이 최신 스테이지 결과로 자유롭게 갱신 (뒤 스테이지가 더 정확한 방법이므로 그 결론을 신뢰)
  • 임계값은 Stage마다 상수(MERGE_THRESHOLD, PENDING_APPROVAL_THRESHOLD)를 따로 둠

오케스트레이터 (Pipeline)

  • CategorySimilarityPipeline
@Component
@RequiredArgsConstructor
public class CategorySimilarityPipeline {

    // 순서는 각 구현체의 @Order로 결정 - 스테이지를 추가/제거/재정렬해도 이 클래스는 안 바뀜
    private final List<CategorySimilarityStage> stages;

    public List<CategoryCandidateResult> resolve(Hashtag hashtag, List<Category> candidates) {
        //...
    }
}

  • 후보(카테고리)별로 독립적인 누적 결과(Map<categoryId, CategoryMatchResult>)를 추적
    • "해시태그 전체에 대한 단일 결과"가 아님
  • 각 스테이지는 아직 확정 안 된(remaining) 후보만 대상으로 실행됨
    • 이번 스테이지 결과를 combineWith로 누적에 접고(fold), 그 결과로 Merge/PendingApproval이 확정된 후보는 remaining에서 빠져 다음 스테이지로 안 넘어감
    • 후보 단위 조기 확정(전체 조기 종료)가 아님
  • remaining이 비거나 모든 스테이지를 다 돌면 종료, 후보별 최종 누적 결과를 반환
  • 후보마다 결과가 다를 수 있어 한 해시태그가 서로 다른 카테고리 여러 개에 동시에(서로 다른 match_type/status로) 연결되는 것도 가능
  • 판단 규칙(우선순위)에는 관여하지 않고 "실행 + 누적"만 담당
    • 포트(CategorySimilarityStage)를 들고 있어서 application 계층에 위치

실제 연결 (Pipeline 밖, Facade/Writer가 담당)

  • CategoryHashtagLinkFacad : Pipeline을 호출하고, 결과를 CategoryHashtagLinkWriter에 넘겨 실제 DB에 반영
@Service
@RequiredArgsConstructor
public class CategoryHashtagLinkFacade implements LinkHashtagToCategoryUseCase {

    private final CategoryCommandRepository categoryCommandRepository;
    private final HashtagCommandRepository hashtagCommandRepository;
    private final CategorySimilarityPipeline categorySimilarityPipeline;
    private final CategoryHashtagLinkWriter categoryHashtagLinkWriter;

    @Override
    public void tryLink(UUID hashtagId) {
        Hashtag hashtag = hashtagCommandRepository.findById(hashtagId)
                .orElseThrow(() -> new BusinessException(CategoryErrorCode.HASHTAG_NOT_FOUND));
        List<Category> categories = categoryCommandRepository.findAllActive();

        List<CategoryCandidateResult> candidateResults = categorySimilarityPipeline.resolve(hashtag, categories);

        categoryHashtagLinkWriter.applyResults(hashtagId, hashtag, candidateResults);
    }
}
  • 후보(카테고리)별 결과를 순회하며 Merge/PendingApproval인 후보는 전부 그 카테고리에 연결
    • linkIfAbsent : ON CONFLICT DO NOTHING이라 중복 연결 시도해도 안전
    • 후보가 여러 개면 여러 카테고리에 동시 연결됨
  • 순회 후 판단(해시태그 전체 기준)

조건 처리

하나라도 Merge/PendingApproval로 매칭됨 (이미 위에서 연결 완료, 추가 처리 없음)
매칭 없음 + 실패(Failed)도 없음 신규 카테고리 생성 후 MERGED로 연결
매칭 없음 + 하나라도 Failed 있음 카테고리 연결 없이 보류 (스케줄러가 재시도 대상으로 나중에 다시 처리)