LLM WikiAccess-protected knowledge portal

WIKI

Apache Flink CDC 3.6: Oracle·Hudi·PostgreSQL 스키마 진화를 단일 파이프라인으로 묶는 방법

요약 Apache Flink CDC 3.6.0이 2026년 3월 30일 출시되었다. 핵심 변화는 세 가지다 Oracle Source Pipeline 커넥터 정식 지원 , Apache Hudi Sink Pipeline 커넥터 추가 , PostgreSQL Schema Evolution 지원 . 여기에 Transform 프레임워크의 VARIANT 타입과 JSON 파싱 내장, Fluss Sink의 Lenient 모드 안정화가 더해졌다

경로human/study/content/database-frontier/87-flink-cdc-3-6-oracle-hudi-schema-evolution-pipeline.md
카테고리Study
태그#cicd #evolution #hudi #mysql #oracle #pipeline #schema #study

요약

Apache Flink CDC 3.6.0이 2026년 3월 30일 출시되었다. 핵심 변화는 세 가지다: Oracle Source Pipeline 커넥터 정식 지원, Apache Hudi Sink Pipeline 커넥터 추가, PostgreSQL Schema Evolution 지원. 여기에 Transform 프레임워크의 VARIANT 타입과 JSON 파싱 내장, Fluss Sink의 Lenient 모드 안정화가 더해졌다. 단일 YAML 파일 하나로 "Oracle → 스키마 진화 → Hudi MoR 적재"까지 코드 없이 선언할 수 있는 파이프라인 모델이 비로소 프로덕션 수준에 이르렀다.

동적 폴백 적용: 이 챕터는 커리큘럼의 모든 예정 항목이 완료된 상태에서 동적 폴백 정책에 따라 선택되었다. Flink CDC 3.6.0 출시일은 2026년 3월 30일(작성 기준 130일 전)로 90일 윈도우를 초과하지만 180일 윈도우 안에 있다. 아직 다루지 않은 기술적으로 중요한 릴리스다.


배경: CDC 파이프라인이 왜 어려운가

CDC(Change Data Capture)는 소스 DB의 변경 이벤트를 싱크로 전달하는 작업이다. 소스 스키마는 살아서 변하고, 소스 시스템은 Oracle처럼 별도 설정이 필요하며, 싱크는 Hudi처럼 자체 파일 포맷 관리가 있다. 전통적 접근의 문제를 정리하면 다음과 같다.

상황기존 접근문제점
ALTER TABLE ADD COLUMN파이프라인 수동 재시작다운타임, 이벤트 유실 위험
Oracle 래거시 소스Debezium·GoldenGate 별도 운영운영 복잡도 증가
Hudi Lakehouse 싱크Spark Structured Streaming 별도Flink·Spark 두 클러스터 필요
PostgreSQL DDL 전파무시(Ignore) 처리싱크 스키마 드리프트 누적

Flink CDC 3.6은 이 네 가지 문제에 직접 대응한다.


Flink CDC 파이프라인 모델

Flink CDC 3.x의 핵심은 선언형 파이프라인 모델이다. 사용자는 YAML 한 파일로 소스·변환·라우팅·싱크를 기술하고 flink-cdc.sh CLI로 제출한다. 내부적으로 세 단계를 거친다.

  1. 소스 커넥터 — 체인지 이벤트(INSERT / UPDATE / DELETE)와 스키마 이벤트(DDL)를 분리해 하류로 전달한다.
  2. Schema Operator — DDL 이벤트를 받아 schema.change.behavior 설정에 따라 싱크 스키마를 업데이트하거나 이벤트를 무시한다.
  3. 싱크 커넥터 — 업데이트된 스키마로 레코드를 기록한다.
소스 레이어
Oracle CDC
LogMiner · redo 로그
MySQL CDC
Binlog
PostgreSQL CDC
WAL 슬롯
Flink CDC 엔진
Schema Operator
DDL 전파·변환
Transform
VARIANT · JSON
Router
정규식 테이블 매핑
싱크 레이어
Apache Hudi
CoW / MoR
Apache Iceberg
Kafka · StarRocks
Paimon · Fluss
Flink CDC 3.6 Pipeline Architecture

Oracle Source Pipeline 커넥터

Oracle CDC는 이전 버전에서도 Flink SQL 커넥터로 존재했지만, 3.6에서 파이프라인 커넥터(YAML 기반)로 승격되었다. LogMiner를 통해 Oracle redo 로그를 읽는다.

선행 조건: Oracle 측 설정

아카이브 로그 모드 활성화

SHUTDOWN IMMEDIATE;
STARTUP MOUNT;
ALTER DATABASE ARCHIVELOG;
ALTER DATABASE OPEN;

보완 로깅(Supplemental Logging) 활성화

-- 데이터베이스 전체 활성화
ALTER DATABASE ADD SUPPLEMENTAL LOG DATA;

-- 또는 테이블 단위 활성화
ALTER TABLE PROD.ORDERS ADD SUPPLEMENTAL LOG DATA;

보완 로깅이 없으면 LogMiner가 UPDATE 이전 값을 읽지 못한다. CDC에서 before-image가 필요한 경우 반드시 활성화해야 한다.

LogMiner 사용자 권한

CREATE USER cdc_user IDENTIFIED BY cdc_pw;
GRANT CREATE SESSION TO cdc_user;
GRANT SELECT ANY DICTIONARY TO cdc_user;
GRANT LOGMINING TO cdc_user;
GRANT SELECT ON V_$LOGMNR_CONTENTS TO cdc_user;

YAML 소스 선언

source:
  type: oracle-cdc
  hostname: 10.0.1.50
  port: 1521
  username: cdc_user
  password: ${ORACLE_CDC_PASS}
  database-name: ORCL
  schema-name: PROD
  table-name: "PROD\\.ORDERS,PROD\\.CUSTOMERS,PROD\\.PRODUCTS"
  scan.startup.mode: initial
  debezium.log.mining.strategy: online_catalog
  debezium.log.mining.continuous.mine: "true"

scan.startup.mode: initial은 첫 실행 시 전체 스냅샷을 찍은 뒤 실시간 redo 로그 스트리밍으로 전환한다. 재시작 시에는 마지막 체크포인트 SCN(System Change Number)부터 다시 읽는다.


Apache Hudi Sink Pipeline 커넥터

Apache Hudi는 Lakehouse 테이블 포맷으로, CoW(Copy-on-Write)와 MoR(Merge-on-Read) 두 가지 저장 전략을 제공한다.

구분CoWMoR
쓰기 방식영향받은 파일 전체 재작성델타 로그 파일만 추가
읽기 방식Parquet 직접 스캔Parquet + 델타 로그 병합
CDC 적합성낮은 갱신 빈도높은 갱신 빈도 (CDC)
읽기 지연낮음Compaction 전까지 다소 높음
Compaction백그라운드 자동자동 또는 수동

CDC 워크로드는 갱신 빈도가 높으므로 대부분 MoR이 적합하다. CoW는 Compaction 비용이 높지만 읽기 경로가 단순하다.

YAML 싱크 선언 (MoR)

sink:
  type: hudi
  catalog.path: s3://my-dw/hudi/prod
  table.type: MERGE_ON_READ
  compaction.async.enabled: "true"
  compaction.trigger.strategy: num_commits
  compaction.delta_commits: "5"
  hive.sync.enabled: "true"
  hive.sync.db: dw
  hive.sync.table: ${table_name}

hive.sync.enabled: true로 설정하면 Hive Metastore(또는 AWS Glue)와 자동 동기화해 Athena·Trino·Spark에서 즉시 쿼리할 수 있다.


PostgreSQL Schema Evolution

3.6 이전 버전에서 PostgreSQL 파이프라인 커넥터는 DDL 이벤트를 무시했다. 3.6부터 schema.change.behavior를 설정하면 PostgreSQL DDL 변경도 하류 싱크에 전파된다.

지원하는 DDL 이벤트:

메타데이터 컬럼 지원도 추가되어 WAL 오프셋이나 소스 테이블명을 Transform 절에서 참조할 수 있게 되었다.


Transform 프레임워크: VARIANT 타입과 JSON 파싱

3.6에서 Transform 절에 VARIANT 타입과 JSON 함수가 내장되었다. JSON 문자열 컬럼을 파싱해 중첩 필드를 정형 컬럼으로 올리는 작업이 파이프라인 YAML 안에서 가능해졌다.

transform:
  - source-table: "PROD\\.EVENTS"
    projection: >
      id,
      CAST(created_at AS TIMESTAMP(3)) AS created_at,
      JSON_VALUE(payload_json, '$.user_id') AS user_id,
      JSON_VALUE(payload_json, '$.action') AS action,
      CAST(amount AS DECIMAL(18, 4)) AS amount

JSON_VALUE·JSON_ARRAY 함수로 반정형 데이터를 정형 컬럼으로 변환하면 Hudi 테이블 구조를 단순하게 유지하면서도 소스의 유연한 스키마를 수용할 수 있다.


스키마 진화 모드 비교

schema.change.behavior 옵션별 동작이다.

모드DDL 전파실패 시 동작DROP/TRUNCATE추천 상황
EVOLVE전체 전파예외 발생·파이프라인 중단전달DDL 변경을 싱크에 반드시 반영해야 할 때
TRY_EVOLVE전체 시도경고 후 계속 (데이터 유실 가능)전달대부분 전파하되 일부 실패는 허용
LENIENT호환 변환 후 전파해당 이벤트 무시무시후방 호환성 유지하며 안전하게 진화
IGNORE없음무시싱크 스키마를 수동 관리하는 환경
EXCEPTION없음DDL 수신 즉시 예외스키마 고정 환경에서 이상 감지용

LENIENT 모드가 CDC 워크로드에 가장 많이 쓰이는 이유: DROP TABLE이나 TRUNCATE가 소스에서 실수로 실행되더라도 싱크에 영향을 주지 않으면서, ADD COLUMN 같은 후방 호환 변경은 그대로 전파한다. 예기치 않은 파이프라인 중단 없이 안전하게 운영할 수 있다.


완성된 파이프라인 YAML 예시

Oracle 세 테이블을 Hudi로 적재하는 전체 파이프라인이다.

################################################################################
# Flink CDC 3.6 — Oracle PROD → Hudi DW
################################################################################
source:
  type: oracle-cdc
  hostname: 10.0.1.50
  port: 1521
  username: cdc_user
  password: ${ORACLE_CDC_PASS}
  database-name: ORCL
  schema-name: PROD
  table-name: "PROD\\.ORDERS,PROD\\.CUSTOMERS,PROD\\.PRODUCTS"
  scan.startup.mode: initial
  debezium.log.mining.strategy: online_catalog

pipeline:
  name: oracle-to-hudi
  parallelism: 4
  schema.change.behavior: LENIENT

transform:
  - source-table: "PROD\\.ORDERS"
    projection: >
      id, customer_id, total_amount,
      CAST(order_time AS TIMESTAMP(3)) AS order_time,
      JSON_VALUE(ext_data, '$.channel') AS channel
  - source-table: "PROD\\.CUSTOMERS"
    projection: id, name, email, tier
  - source-table: "PROD\\.PRODUCTS"
    projection: id, sku, category, price

route:
  - source-table: "PROD\\.(.*)"
    sink-table: "dw.$1"
    description: PROD 스키마 → dw 스키마, 테이블명 소문자 유지

sink:
  type: hudi
  catalog.path: s3://my-dw/hudi/prod
  table.type: MERGE_ON_READ
  compaction.async.enabled: "true"
  compaction.delta_commits: "5"
  hive.sync.enabled: "true"
  hive.sync.db: dw

제출 명령:

./bin/flink-cdc.sh oracle-to-hudi.yaml \
  --flink-home /opt/flink \
  --jar flink-cdc-dist-3.6.0.jar

Flink 버전 호환성

3.6 파이프라인 커넥터가 지원하는 Flink 버전:

Flink 버전파이프라인 커넥터 지원JDK
1.18.x 이하✗ (지원 종료)
1.19.x
1.20.x11 이상
2.2.x11 이상

3.6으로 업그레이드 시 JDK를 11 이상으로 맞춰야 한다. Java 8 환경은 더 이상 지원하지 않는다.


운영 체크리스트

Oracle 소스 준비

Hudi Sink 준비

스키마 진화 모드 선택

모니터링


요점 정리

References