Skip to content

MySQL ‐ Operational‐Level System Design

woojin edited this page Jul 20, 2026 · 7 revisions

실시간 변동 데이터 처리를 위한 스트리밍 데이터 처리 기법

1. Data Collection — 데이터가 발생하는 곳

  • Web/App Server : 사용자가 앱에서 클릭·검색·결제 → 서버가 "클릭 이벤트" 발생한다.
  • IoT Device: 센서가 온도·위치 값을 계속 측정해서 발생한다.
  • Database CDC : DB에 주문이 INSERT되면 그 변경을 감지해 이벤트가 발생(CDC = Change Data Capture, "DB 바뀐 것 실시간 감지")한다.

2. Stream Storage — 일단 큐에 모임(Kafka)

  • 발생한 이벤트를 곧바로 처리하지 않고 메시지 큐(Kafka)에 잠깐 쌓는다.
    • 완충(버퍼) : 갑자기 트래픽이 폭주해도 큐가 받아내고, 뒷단은 자기 속도로 소비할 수 있다.
    • 분리(디커플링) : 데이터 만드는 쪽과 처리하는 쪽이 서로 몰라도 됨 → 독립적으로 확장·교체 가능하다.
    • 여러 소비자 : 같은 이벤트를 여러 시스템이 나눠 읽을 수 있다.

3. Stream Processing — 실시간으로 가공(Flink)

  • 스트림 프로세서가 큐에서 이벤트를 도착하는 즉시 꺼내서 계산한다.
  • 데이터가 오는 순간 바로 처리해서 지연이 초 단위이다.

4. Data Serving — 결과를 여러 곳에서 사용

  • 가공한 결과를 필요한 곳으로 팬아웃(fan-out)한다.

즉, 여러 소스에서 끊임없이 발생한 데이터를 Kafka에 모아 완충하고 스트림 프로세서가 즉시 가공해서 대시보드·알림·ML·DW 등 여러 곳에 실시간으로 공급하는 구조를 만들 수 있다.

수천억대 데이터 처리를 위한 초대용량 배치 처리 방법

  • 기존 배치는 스케줄러가 하나의 INSERT ... SELECT를 던지고, 스캔·그룹핑·집계·저장을 전부 운영 MySQL이 처리하는 구조다.
  • 데이터가 수천만 건일 때는 버티지만, 수억 ~ 수천억 건으로 늘면 세 갈래로 무너진다.
    • 처리 시간 증가 : 단일 스레드·단일 노드로 전 구간을 순차 스캔·집계 → 데이터가 커질수록 배치 윈도우를 초과한다.
    • DB 부하 → 서비스 영향 : 대량 스캔이 버퍼풀·I/O·CPU를 점유해 같은 DB를 쓰는 실서비스 응답이 함께 느려진다.
    • 장애 대응의 어려움 : 한 트랜잭션으로 8시간 돌다 실패하면 처음부터. 재시작·부분복구·멱등성 보장이 어렵다.

해결 개념 - 연산을 DB 밖으로 위임한다

  • DB에게 집계를 시키지 말고, 데이터만 병렬로 꺼내와 밖에서 나눠 계산한다.
  • Spark가 이 문제에 주는 두 가지 실질 이점이 있다.
    • 계산 위임 — GROUP BY·SUM 같은 무거운 연산을 MySQL이 아니라 Spark 클러스터가 수행한다.
    • 병렬 데이터 로딩 — 하나의 큰 테이블을 order_id 범위로 쪼개 여러 Executor가 동시에 나눠 읽는다. 순차 스캔이 병렬 스캔으로 바뀐다.

1. 병렬 읽기용 뷰 정의

  • MySQL의 orders를 partitionColumn 기준으로 쪼개 여러 Executor가 동시에 읽도록 임시 뷰를 만든다.
  • numPartitions만큼의 범위 쿼리가 병렬로 나간다.
CREATE OR REPLACE TEMPORARY VIEW orders_view
USING "jdbc" OPTIONS (
  url             "jdbc:mysql://mysql-server:3306/service_db",
  dbtable         "orders",
  partitionColumn "order_id",   -- 나눌 기준 (숫자·PK 권장)
  lowerBound      "1",          -- 최소값
  upperBound      "100000000",  -- 최대값
  numPartitions   "100"         -- 100개로 분할 → 병렬 읽기
);
-- Executor 1  : WHERE order_id >= 1        AND order_id < 1000000
-- Executor 2  : WHERE order_id >= 1000000  AND order_id < 2000000
-- Executor 100: WHERE order_id >= 99000000 AND order_id <= 100000000

2. 집계는 Spark 엔진이 수행한다.

  • 이 쿼리는 더 이상 MySQL이 아니라 Spark 엔진이 인메모리로 실행한다.
  • WHERE created_at 조건은 Predicate Pushdown으로 읽기 단계에서 DB에 밀어넣어 스캔량을 줄이고, GROUP BY는 노드 간 셔플링을 거쳐 집계된다.
CREATE OR REPLACE TEMPORARY VIEW daily_summary AS
SELECT product_id,
       SUM(sale_price * quantity) AS total_amount,
       COUNT(*)                    AS total_count
FROM orders_view
WHERE created_at >= '2024-05-20 00:00:00'
  AND created_at <  '2024-05-21 00:00:00'   -- ← Predicate Pushdown
GROUP BY product_id;                        -- ← Shuffle 발생

3. 결과를 분석 DB에 저장한다.

  • 집계 결과는 운영 DB가 아닌 별도의 분석용 DB에 적재해 서비스와 격리한다.
INSERT INTO analytics_db.daily_product_sales
SELECT * FROM daily_summary;

실시간 데이터 동기화를 위한 디자인 설계 패턴

1. 문제 : 하나의 변경, 여러 목적지

  • Dual Write : 애플리케이션이 DB에 쓰고 또 Kafka에도 직접 쓴다. 두 쓰기는 한 트랜잭션이 아니라, 하나가 실패하면 시스템 간 데이터가 영구히 어긋난다.
  • 주기적 배치 폴링 : updated_at을 N분마다 훑기. 실시간이 아니고, 매 폴링이 운영 DB에 부하를 주며, 삭제(DELETE)는 감지조차 못 한다.

핵심 : 이미 정합성이 보장된 단 하나의 소스. DB의 트랜잭션 로그를 진실의 원천으로 삼는다.

CDC 파이프라인 — 4개 계층

1. Source — MySQL (Binlog)

  • 모든 변경이 이미 binlog에 순서대로 기록된다.
  • CDC의 전제 조건은 binlog_format=ROW 활성화이다.

2. CDC — Debezium Connector

  • MySQL의 복제 슬레이브인 척 binlog를 구독해 변경을 표준 이벤트로 변환한다.
  • 애플리케이션은 이 존재를 모른다.

3. Streaming — Kafka Topic

  • 변경 이벤트의 완충·보존·순서 보장 버퍼. 소비자가 잠시 죽어도 이벤트는 토픽에 남아 재생 가능하다.

4. Consumers

  • 검색 · 분석 · 캐시 · 각자 독립적으로 같은 토픽을 구독한다.
  • 소비자를 추가해도 생산 측은 전혀 바뀌지 않는다. 이것이 이벤트 브로커의 디커플링 이점이다.

왜 Kafka를 중간에 두나? - Debezium이 소비자에 직접 밀어주면, 소비자 하나가 느리면 전체가 막히고 재처리도 어렵다. Kafka가 이벤트를 보존하기에 소비자별 속도 차이·장애·나중 합류(replay)를 모두 흡수한다.

1. Debezium 동작 원리 — 스냅샷 → 스트리밍

  • 커넥터가 처음 붙을 때 풀어야 하는 문제 : 지금까지 쌓인 데이터(과거)와 앞으로 바뀔 데이터(미래)를 빠짐도 중복도 없이 이어붙이는 것이다. 이 때, 순서가 핵심이다.
  • 짧은 Lock으로 기준점 고정 : 아주 짧게 GLOBAL READ LOCK을 잡고 현재 binlog 위치를 저장한 뒤 즉시 해제한다. 이 시점이 과거와 미래의 경계선이 된다.
  • Lock 해제 후에도 변경은 안전 : 해제 순간부터의 변경은 그대로 binlog에 쌓이므로 유실되지 않는다. 그래서 Lock을 오래 잡을 필요가 없다(서비스 영향 최소).
  • 초기 스냅샷 : 전건을 'c' 이벤트로 · 테이블 전체를 SELECT해 “현재 상태”를 Create 이벤트로 Kafka에 채운다. 소비자 입장에선 모든 행이 “새로 생성된” 것처럼 보인다.
  • 저장한 위치부터 스트리밍 : 스냅샷이 끝나면 저장해둔 binlog 위치로 돌아가 'u'/'d' 이벤트를 순서대로 흘린다. 경계선 이후 것만 이어지므로 이중 반영이 없다.

최신 Debezium은 Lock을 더 줄인다. 위 흐름은 Lock 기반 스냅샷의 표준 모델이다. 실제로는 Incremental Snapshot(무Lock, 청크 단위 워터마크)으로 스냅샷과 스트리밍을 동시에 진행해 운영 부담을 더 낮추는 방식이 널리 쓰인다.

2. 변경 이벤트 메시지 구조

  • Debezium의 이벤트는 “무엇이, 어떻게 바뀌었는지”를 담는 표준 envelope다. op는 연산 종류, before/after는 변경 전후 상태이다.
{
  "op": "u",                 // c=create, u=update, d=delete, r=snapshot read
  "ts_ms": 1721457600000,
  "before": { "id": 42, "price": 1000, "status": "OPEN" },
  "after":  { "id": 42, "price": 1200, "status": "OPEN" },
  "source": { "db": "service_db", "table": "orders",
              "file": "binlog.000042", "pos": 15832 }
}
  • Kafka 메시지 Key = PK : 같은 레코드의 이벤트가 항상 같은 파티션으로 가도록 키를 PK로 둔다 → 그 레코드에 한해 순서가 보장된다.
  • 삭제는 'd' + tombstone : Delete 이벤트 뒤에 value=null 툼스톤을 보내, 키 기반 log compaction 토픽에서 실제로 지워지게 한다.
  • 'r' vs 'c' : 초기 스냅샷은 read(r), 실시간 생성은 create(c). 소비자는 대개 둘을 동일하게 “upsert”로 처리한다.
  • 주의 1 : 순서 보장은 “파티션 단위”까지만. 전역 순서는 없다. PK를 파티션 키로 삼아 레코드별 순서를 확보하고, 전역 순서가 필요하면 파티션 1개(=처리량 희생)를 감수한다.
  • 주의 2 : 멱등 소비자(Idempotent Consumer) : Kafka는 기본 at-least-once로 같은 이벤트가 재전송될 수 있다. 소비 측을 UPSERT·버전 비교로 멱등하게 만들어 중복을 무해화한다.
  • 주의 3 : Eventual Consistency 수용 : 소비자는 소스보다 항상 조금 뒤처진다. “강한 일관성”이 필요한 화면은 CDC 대상에서 제외하거나 소스를 직접 읽는다.
  • 주의 4 : 스키마 진화 대비 · ALTER TABLE은 이벤트 스키마를 바꾼다. Schema Registry + 호환성 규칙(backward)으로 소비자가 깨지지 않게 한다.
  • 주의 5 : 실패 격리 DLQ : 특정 이벤트가 소비 중 계속 실패하면 파이프라인 전체가 멈춘다. Dead Letter Queue로 밀어내 나머지를 계속 흐르게 하고 나중에 재처리한다.
  • 주의 6 : Transaction Outbox로 Dual Write 봉쇄 : 애플리케이션이 이벤트를 발행해야 한다면, 비즈니스 데이터와 함께 outbox 테이블에 한 트랜잭션으로 넣고 그 테이블을 CDC로 흘린다.
  • 주의 7 : 재적재(Backfill) 경로 확보 : 소비자 로직 변경·유실 시 스냅샷을 다시 태워 전량 재구성할 수 있어야 한다. Kafka 보존기간·compaction 정책을 이에 맞춰 설계한다.

견고한 비동기 작업을 위한 MySQL 작업 큐 활용법

MySQL + NoSQL을 결합한 시스템 아키텍처 설계

📖 Code Philosophy

📖 Java

📖 Kotlin

📖 Coroutine

📖 Spring

📖 Spring Security

📖 Spring Batch

📖 Database

📖 MySQL

📖 Redis

📖 JPA

📖 QueryDsl

📖 MSA

📖 Kafka

📖 Apache Flink

  • [Apache Flink - Apache Flink Architecture]
  • [Apache Flink - Stream Processing]
  • [Apache Flink - Data Stream API & Window]
  • [Apache Flink - State Management]

📖 HTTP

📖 AWS

📖 Docker

📖 Kubernetes

📖 Github Actions

📖 Jenkins

📖 Nginx

📖 Monitoring

📖 Test(feat. Load Testing)

  • [Test - Load Testing Fundamentals]
  • [Test - Identifying Bottlenecks with Load Testing]
  • [Test - Resolving Bottlenecks and Improving Performance]

📖 Test(feat. Java)

📖 Spring AI

📖 gRPC

  • [gRPC - Writing .proto Files with Protocol Buffers]
  • [gRPC - Various Communication Patterns in gRPC]
  • [gRPC - gRPC Optimization Techniques and Advanced Features]

📖 TDD(Test-Driven-Development)

📖 PostgreSQL

  • [PostgreSQL - Docker만을 사용하는 경량화된 환경 구성 방법]
  • [PostgreSQL - PostgreSQL에서 제공하는 데이터 타입]
  • [PostgreSQL - PostgreSQI의 JSONB, 역인덱싱과 활용 방법]
  • [PostgreSQL - 데이터베이스 성능을 위한 최적화 패턴 및 전략]
  • [PostgreSQL - 트랜잭션과 ACID, Isolation 수준별 차이]
  • [PostgreSQL - Database Lock 교착상태와 읽기/쓰기 성능을 보장하는 MVCC 모델]
  • [PostgreSQL - pgvector와 벡터 저장, 유사도 검색 패턴 개념]
  • [PostgreSQL - 벡터 인덱스 최적화와 벡터 검색과 전문 검색 결합 패턴]
  • [PostgreSQL - PostgreSQL 플러그인]
  • [PostgreSQL - PostGIS - 공간 쿼리와 GIST 인덱스, 지리 타입과 공간 쿼리를 위한 타입과 기본 함수]
  • [PostgreSQL - pg_search - 검색 엔진 없이 텍스트 검색 구현과 주의사항]
  • [PostgreSQL - 단일 인스턴스 한계를 극복하는 분산 패턴과 스케줄링, 분산 환경 구축 방법]
  • [PostgreSQL - Citus - 분산 테이블과 분산 쿼리를 위한 Extension과 데이터 분산 처리]
  • [PostgreSQL - pg_cron - PostgreSQL로 구성하는 CronJob]
  • [PostgreSQL - 스케줄러 + 분산 처리를 동시에 도입하는 주기적 집계 쿼리 패턴]

📖 Workflow-Driven Techniques for Large-Scale Traffic Processing

  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Kafka + Debezium을 활용한 CDC 패턴 설계]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Temporal을 활용한 워크플로우 패턴]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Docker와 경량 이미지를 활용한 환경 구축 방법]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Kafka에서의 메시지 Delivery Guarantee]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - 실시간 동기화의 핵심 CDC]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - MySQL Binary Log 기반의 CDC]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Binary Log 기반의 CDC 구현 플랫폼 Debezium이란?]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Debezium Architecture]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Debezium Architecture Best Practice와 주의사항]

📖 Reactive Programming

📖 ElasticSearch

  • [ElasticSearch - Search Fundamentals & How Search Works]
  • [ElasticSearch - Korean-Optimized Search]
  • [ElasticSearch - Mapping & Data Types]
  • [ElasticSearch - Common Search Features]
  • [ElasticSearch - Building a Product Search Engine with Elasticsearch]
  • [ElasticSearch - Deploying Elasticsearch with Elastic Cloud]
  • [ElasticSearch - Managing Documents]
  • [ElasticSearch - Analysis & Mapping]
  • [ElasticSearch - Search Fundamentals]
  • [ElasticSearch - Query Joins]
  • [ElasticSearch - Processing Search Results]
  • [ElasticSearch - Aggregations]
  • [ElasticSearch - Tips for Improving Search Results]
  • [ElasticSearch - Elasticsearch Clients]
  • [ElasticSearch - Understanding How Elasticsearch Works]
  • [ElasticSearch - Monitoring Elasticsearch]
  • [ElasticSearch - Elasticsearch Troubleshooting]

Clone this wiki locally