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.cmd에 bld 통계 id로 POST 시 JSON 응답 |
| bld | KRX 통계 화면 식별자. 통계 종류마다 다른 값 (투자자별 거래실적·공매도 등) |
| OTP 토큰 | bld로 generate.cmd에서 발급받아 getJsonData.cmd에 전달하는 1회용 코드 |
| 투자자별 거래실적 | 개인·외국인·기관(연기금·투신·금투·보험 등) 일별 순매수 |
| EOD | End Of Day. 장 마감 기준 확정 통계 |
| marketflow | 신설 도메인 패키지. 일별 종목 수급·미시구조 스냅샷 |
marketflow.collected.v1 | 수집 완료 이벤트. payload는 기준일(baseDate) |
| worker 서버 | 신규 Spring Boot(:8091). 이벤트 소비·재계산 오케스트레이션 전담, stateless |
| DLQ | Dead 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 signal_snapshots 테이블에 영속화 (Redis 제거)
- ML in-memory TtlCache 제거 여부 결정 — ml은 stateless, 영속화는 MySQL로
- CUD 요청의 Toss와 MySQL 처리 순서 결정 — 외부 호출 어댑터는 BE infrastructure에 배치
| 방안 | 설명 | 채택 여부 |
|---|---|---|
| BE(Kotlin) 수집 + MySQL 영속화 | 토스 게이트웨이·신호 영속화와 동일 위치. ml stateless 유지 | 채택 — 기존 ADR 일관, 영속화 책임 BE 집중 |
| ml(Python) 수집 + 자체 저장 | 스크래핑·스코어링이 한 서비스. 구현 단순 | 미채택 — ml stateless 방향 위배, 영속화 이원화 |
| 실시간 스크래핑 | 장중 폴링 | 미채택 — EOD 데이터라 가치 없음, 차단·노이즈 리스크 |
| 도메인 단일 aggregate | 5종을 한 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.duration | timer | 배치 전체 소요 |
marketflow.upsert.duration{statistic} | timer | 벌크 upsert 소요 |
marketflow.collect.failure{statistic} | counter | 통계별 실패·부분실패 |
marketflow.krx.retry | counter | KRX 재시도·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
- DB 마이그레이션(6개 테이블) 선반영 (Flyway).
- Kafka 토픽(
marketflow.collected.v1, DLQ) 생성. - KRX MDC client + 팩터별 Gateway·Repository 배포 (스케줄러·발행 비활성 상태).
- 스케줄러 활성화, 1회 수동 트리거로 적재 검증 (이벤트 발행 OFF 상태).
- ml 예측이 KRX 팩터를 입력으로 받도록 반영 배포.
- worker 서버 배포 + 이벤트 발행 활성화, 1회 수동 트리거로 collected.v1 → 재계산 → 예측 갱신 검증.
- 롤백: ① 발행 플래그 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