KRX 수급 데이터 연동 TDD

Background

PRD 참조. 토스 Open API에 없는 외국인·기관 수급 등 시장 미시구조 데이터를 KRX(data.krx.co.kr)에서 EOD로 수집해 추천 산정 팩터로 제공한다.

Overview

  • backend/(Kotlin, Hexagonal)에 marketflow 도메인 패키지를 신설한다.
  • KRX MDC 호출은 infrastructure/krx의 어댑터로 캡슐화한다 (토스가 infrastructure/toss인 것과 동일 구조).
  • 마감 후 1일 1배치 스케줄러가 투자자별 거래실적·외국인 보유·밸류에이션·공매도·신용을 수집해 MySQL에 upsert한다.
  • 스케줄러와 별개로 수동 갱신 API(POST /market-flow/collect)가 동일 수집 UseCase를 호출한다 (멱등). 자동·수동 두 진입점 모두 presentation에 위치.
  • 적재는 통계당 JDBC 벌크 upsert(INSERT … ON DUPLICATE KEY UPDATE)로 처리한다 (전종목 ~2,800행/통계, 행 단위 save 금지).
  • 수집 커밋 후 BE가 marketflow.collected.v1(기준일) 이벤트를 Kafka로 발행한다 (DomainEventPublisher 추상화, otelpoc Kafka 설정 패턴 재사용).
  • 신규 worker 서버(:8091, Spring Boot)가 이벤트를 소비해 추천·포트폴리오·관심종목 재계산을 오케스트레이션한다. worker는 입력을 BE read API로 읽고, 계산을 ml에 위임하고, 결과를 BE 영속화 API로 저장한다 (worker stateless).
  • presentation의 read API가 종목·기준일별 수급 팩터를 노출하고, worker 재계산이 이를 입력으로 소비한다.
  • 배치 run 결과를 marketflow_collect_log에 남기고, 실패 시 Discord 알림 + Micrometer 메트릭(행수·소요·재시도)을 계측한다 — 청크 전환 시그널 관측용. worker 소비는 재시도→DLQ, consumer lag·DLQ를 관측한다.

Terminology

용어정의
KRX MDC정보데이터시스템(data.krx.co.kr). getJsonData.cmdbld 통계 id로 POST 시 JSON 응답
bldKRX 통계 화면 식별자. 통계 종류마다 다른 값 (투자자별 거래실적·공매도 등)
OTP 토큰bld로 generate.cmd에서 발급받아 getJsonData.cmd에 전달하는 1회용 코드
투자자별 거래실적개인·외국인·기관(연기금·투신·금투·보험 등) 일별 순매수
EODEnd Of Day. 장 마감 기준 확정 통계
marketflow신설 도메인 패키지. 일별 종목 수급·미시구조 스냅샷
marketflow.collected.v1수집 완료 이벤트. payload는 기준일(baseDate)
worker 서버신규 Spring Boot(:8091). 이벤트 소비·재계산 오케스트레이션 전담, stateless
DLQDead Letter Queue. 재시도 후에도 실패한 메시지 격리 토픽
consumer lag토픽 최신 오프셋과 소비 오프셋 차이. 소비 지연 지표

Define Problem

AS-IS

추천 스코어링 입력 = 기술 모멘텀 + 뉴스 감성
- 외국인·기관 수급 데이터 없음 (토스 API 미제공)
- 밸류에이션·공매도·신용 데이터 없음
→ 중기 시그널 근거 부족

TO-BE

backend/marketflow (Kotlin, Hexagonal)
  presentation   MarketFlowApiController (read), MarketFlowCollectScheduler, MarketFlowCollectApiController
  application    CollectMarketFlowUseCase, GetMarketFlowUseCase
  domain         InvestorTrading / ForeignHolding / Valuation / ShortSelling / CreditBalance
                 *Repository, *Gateway, DomainEventPublisher (interface)
  infrastructure krx/   KrxMdcClient + Krx*GatewayImpl  (OTP·retry·backoff)
                 mysql/ *RepositoryImpl + JPA
                 kafka/ KafkaDomainEventPublisher (marketflow.collected.v1)

worker/ (신규 Spring Boot :8091)
  presentation   MarketFlowCollectedEventWorker (@KafkaListener)
  application    RefreshRecommendationUseCase / RefreshPortfolioForecastUseCase / RefreshWatchlistSignalUseCase
  domain         재계산 오케스트레이션 (입력=BE, 계산=ml, 영속화=BE)
  infrastructure BE read/persist client, ml compute client

흐름:
마감 후 1배치 → 5종 수집 → upsert → collected.v1 발행
  → worker 소비 → (BE 입력 read → ml 계산 → BE 영속화) ×3 도메인

Possible Solutions

과거 결정 참조

방안설명채택 여부
BE(Kotlin) 수집 + MySQL 영속화토스 게이트웨이·신호 영속화와 동일 위치. ml stateless 유지채택 — 기존 ADR 일관, 영속화 책임 BE 집중
ml(Python) 수집 + 자체 저장스크래핑·스코어링이 한 서비스. 구현 단순미채택 — ml stateless 방향 위배, 영속화 이원화
실시간 스크래핑장중 폴링미채택 — EOD 데이터라 가치 없음, 차단·노이즈 리스크
도메인 단일 aggregate5종을 한 Entity에미채택 — 테이블·수집주기·시차가 달라 팩터별 분리가 단순
수집 후 Kafka 이벤트 발행 → worker 재계산수집기·소비자 디커플링, 3개 도메인 fan-out, Kafka 도입 방향(ADR-007 보류분) 실현채택 — ADR-007/008
수집 후 동기 인프로세스 재계산Kafka 없이 배치가 곧바로 재계산 호출미채택 — 디커플링·fan-out·확장성 상실 (1일 1회엔 단순하나 방향성 부재)
별도 worker 서버 vs backend consumer예측 갱신 실행 주체별도 worker 서버 채택 (ADR-008) — 프로세스 격리. 1일 1회엔 과하나 사용자 결정

Detail Design

Component Diagram

flowchart LR
    subgraph Presentation
        Scheduler[MarketFlowCollectScheduler]
        CollectCtrl[MarketFlowCollectApiController]
        Controller[MarketFlowApiController]
    end
    subgraph Application
        CollectUC[CollectMarketFlowUseCase]
        GetUC[GetMarketFlowUseCase]
    end
    subgraph Domain
        DS[MarketFlowDomainService]
        Repo[(MarketFlow Repositories)]
        Gw[MarketFlow Gateways]
        Pub[DomainEventPublisher]
    end
    subgraph Infrastructure
        Krx[KrxMdcClient + Krx GatewayImpl]
        Mysql[RepositoryImpl JPA]
        Kafka[KafkaDomainEventPublisher]
    end
    Scheduler --> CollectUC
    CollectCtrl --> CollectUC
    Controller --> GetUC
    CollectUC --> DS
    GetUC --> DS
    DS --> Repo
    DS --> Gw
    DS --> Pub
    Krx -.implements.-> Gw
    Mysql -.implements.-> Repo
    Kafka -.implements.-> Pub

Component Diagram — worker 서버

flowchart LR
    subgraph Kafka
        Topic[marketflow.collected.v1]
        Dlq[collected.v1.DLQ]
    end
    subgraph Worker["worker (:8091)"]
        EW[MarketFlowCollectedEventWorker]
        RecUC[RefreshRecommendationUseCase]
        PfUC[RefreshPortfolioForecastUseCase]
        WlUC[RefreshWatchlistSignalUseCase]
    end
    subgraph External
        BE[backend read/persist API]
        ML[ml compute API]
    end
    Topic --> EW
    EW --> RecUC
    EW --> PfUC
    EW --> WlUC
    RecUC --> BE
    RecUC --> ML
    PfUC --> BE
    PfUC --> ML
    WlUC --> BE
    WlUC --> ML
    EW -.실패 격리.-> Dlq

Sequence Diagram — 수집 (자동 스케줄러 / 수동 트리거 공통)

sequenceDiagram
    participant Sched as Scheduler / ManualTrigger
    participant UC as CollectMarketFlowUseCase
    participant DS as MarketFlowDomainService
    participant Gw as Krx GatewayImpl
    participant Krx as KrxMdcClient
    participant Repo as RepositoryImpl
    Sched->>UC: execute(baseDate)
    UC->>DS: collect(baseDate)
    loop 5종 통계
        DS->>Gw: fetch(baseDate)
        Gw->>Krx: getJsonData(bld, params)
        Krx-->>Gw: rows
        Gw-->>DS: domain models
        DS->>Repo: upsertAll(models)
    end
    DS-->>UC: collected summary
    UC->>DS: publish collected.v1 (commit 후)
    UC-->>Sched: result

Sequence Diagram — 이벤트 기반 예측 갱신

sequenceDiagram
    participant Topic as collected.v1
    participant EW as EventWorker
    participant UC as Refresh*UseCase
    participant BE as backend API
    participant ML as ml compute
    Topic->>EW: consume(baseDate)
    loop 추천·포트폴리오·관심종목
        EW->>UC: refresh(baseDate)
        UC->>BE: read inputs (factors·holdings·watchlist)
        UC->>ML: compute(inputs)
        ML-->>UC: predictions
        UC->>BE: upsert predictions (멱등)
    end
    Note over EW: 실패 시 재시도→DLQ, 재처리 멱등

ERD

erDiagram
    krx_investor_trading {
        bigint id PK
        varchar symbol
        date base_date
        bigint foreign_net
        bigint institution_net
        bigint individual_net
    }
    krx_foreign_holding {
        bigint id PK
        varchar symbol
        date base_date
        bigint holding_shares
        decimal holding_ratio
        decimal limit_exhaustion_ratio
    }
    krx_valuation {
        bigint id PK
        varchar symbol
        date base_date
        decimal per
        decimal pbr
        decimal dividend_yield
    }
    krx_short_selling {
        bigint id PK
        varchar symbol
        date trade_date
        date disclosed_date
        bigint short_volume
        decimal short_balance_ratio
    }
    krx_credit_balance {
        bigint id PK
        varchar symbol
        date base_date
        bigint credit_balance
        decimal credit_ratio
    }
    marketflow_collect_log {
        bigint id PK
        date base_date
        varchar statistic
        varchar status
        int row_count
        int duration_ms
        varchar error_message
        datetime created_at
    }
  • 각 팩터 테이블 (symbol, base_date) unique 제약으로 upsert 멱등 보장 (공매도는 (symbol, trade_date)).
  • marketflow_collect_log는 run별·통계별 1행 — 상태(완료/실패/부분실패)·행수·소요·에러를 기록. 관측 스택과 무관하게 동작하는 실행 이력 원천.
  • FK 컬럼 미사용 (애플리케이션 레벨 관리), DATE/DATETIME(6) 정밀도 준수, BOOLEAN 미사용(status는 VARCHAR).
  • 이벤트는 outbox 테이블을 두지 않는다 (정책 6, 1일 1회·멱등·수동 재트리거). collected.v1은 커밋 후 직접 발행.
  • 예측 영속화 테이블: 관심종목 시그널은 기존 signal_snapshots(STK9) 재사용. 추천·포트폴리오 예측은 기준일 단위 멱등 upsert가 가능한 스냅샷 테이블이 필요 — 구현 시 main의 STK8/STK9 머지 상태를 확인해 재사용 또는 신설한다(예측별 (scope_key, base_date) unique).

Testing Plan

  • KrxMdcClient: OTP 발급→getJsonData 흐름, 4xx/5xx 시 재시도·백오프, 차단 응답 처리 (TestContainers 불필요, WireMock/MockWebServer)
  • Krx*GatewayImpl: KRX JSON → 도메인 모델 파싱, 빈 응답·필드 누락 graceful
  • 각 *RepositoryImpl: upsert 멱등(동일 symbol·base_date 재적재), 통합(TestContainers MySQL)
  • MarketFlowDomainService: 부분 실패 시 나머지 적재 지속, 공매도 거래일/공시일 구분
  • CollectMarketFlowUseCase: 5종 수집 오케스트레이션, 일부 실패 허용
  • 적재: 통계당 1회 벌크 upsert로 ~2,800행 멱등 처리 (행 단위 save 미사용 검증)
  • 모니터링: run 결과가 collect_log에 기록, 실패·부분실패 시 Discord 알림 호출, 메트릭 카운터·타이머 증가
  • MarketFlowApiController: 종목·기준일 조회, 미수집 시 빈 응답
  • 이벤트 발행: 수집 커밋 후 collected.v1 1건 발행, 수집 실패 시 미발행
  • KafkaDomainEventPublisher: 직렬화·토픽 전송 (TestContainers Kafka)
  • EventWorker: collected.v1 소비 → 3개 Refresh UseCase 호출, 소비 실패 시 재시도→DLQ (TestContainers Kafka)
  • Refresh*UseCase: BE 입력 read → ml compute → BE upsert 오케스트레이션, 동일 baseDate 재처리 멱등, ml/BE 실패 시 격리
  • scenario: 마감 후 수집 → collected.v1 발행 → worker 소비 → 3개 도메인 예측 갱신 E2E

Observability

기존 옵저버빌리티 스택(PRD)에 정렬한다 — OTel(OTLP) → 단일 Collector → Grafana(Prometheus/Loki). backend는 OTel Java agent 자동계측이라 Micrometer 메터가 그대로 export된다.

실행 이력·알림 (관측 스택 비의존)

  • marketflow_collect_log에 run별·통계별 상태·행수·소요·에러 기록.
  • 실패·부분 실패 시 DiscordNotificationGateway(alert 도메인)로 즉시 알림 (마감 후 배치).

메트릭 (청크 전환 시그널)

메트릭타입의미
marketflow.collect.rows{statistic}gauge통계별 수집 행수
marketflow.collect.duration{statistic}timer통계별 소요
marketflow.collect.total.durationtimer배치 전체 소요
marketflow.upsert.duration{statistic}timer벌크 upsert 소요
marketflow.collect.failure{statistic}counter통계별 실패·부분실패
marketflow.krx.retrycounterKRX 재시도·429 횟수

JVM heap은 OTel agent 자동계측 메트릭을 그대로 활용한다.

전환 임계(알람)

  • 단일 통계 행수 > 50,000, 전체 소요 > 5분, 단일 통계 > 2분, 벌크 upsert > 30초, heap 사용률 급증/OOM 근접, 부분 실패·429 빈발.
  • 임계 지속 초과 시 Grafana 알람 → Discord. 이것이 Spring Batch(청크·restart·partition) 전환 신호다.
  • 대시보드·알람 규칙은 옵저버빌리티 스택(STK-OBS) 완료에 의존. 메트릭 emission·DB 로그·Discord 알림은 선행 무관하게 동작.

worker 소비 관측

  • worker도 OTel Java agent 자동계측 — Kafka consumer lag·처리 시간이 자동 수집된다.
  • 커스텀: marketflow.refresh.duration{domain}(도메인별 재계산 소요), marketflow.refresh.failure{domain}(실패), marketflow.dlq.count(DLQ 적재).
  • consumer lag 임계·DLQ 적재 시 Discord 알림. collected.v1 발행 후 일정 시간 내 재계산 미완(예측 미갱신)도 알람 대상.

Release Scenario

  1. DB 마이그레이션(6개 테이블) 선반영 (Flyway).
  2. Kafka 토픽(marketflow.collected.v1, DLQ) 생성.
  3. KRX MDC client + 팩터별 Gateway·Repository 배포 (스케줄러·발행 비활성 상태).
  4. 스케줄러 활성화, 1회 수동 트리거로 적재 검증 (이벤트 발행 OFF 상태).
  5. ml 예측이 KRX 팩터를 입력으로 받도록 반영 배포.
  6. worker 서버 배포 + 이벤트 발행 활성화, 1회 수동 트리거로 collected.v1 → 재계산 → 예측 갱신 검증.
  7. 롤백: ① 발행 플래그 OFF → 이벤트 중단(예측은 직전 값 유지) ② 스케줄러 OFF → 수집 중단 ③ worker 중지 → 재계산 중단. 예측은 KRX 팩터 없이 graceful. 테이블은 역방향 DDL로 제거 가능.

Project Information

  • 저장소: backend/(com.biuea.stock.marketflow + 이벤트 발행), worker/(신규, 이벤트 소비·재계산), ml/(예측 계산)
  • 포트: backend :8080, ml :8000, aggregator :8090, worker :8091 (신규)
  • Kafka 토픽: marketflow.collected.v1 (+ DLQ)
  • 티켓 prefix: STK10