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 모드 안정화가 더해졌다. 단일 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로 제출한다. 내부적으로 세 단계를 거친다.
- 소스 커넥터 — 체인지 이벤트(INSERT / UPDATE / DELETE)와 스키마 이벤트(DDL)를 분리해 하류로 전달한다.
- Schema Operator — DDL 이벤트를 받아
schema.change.behavior설정에 따라 싱크 스키마를 업데이트하거나 이벤트를 무시한다. - 싱크 커넥터 — 업데이트된 스키마로 레코드를 기록한다.
LogMiner · redo 로그
Binlog
WAL 슬롯
DDL 전파·변환
VARIANT · JSON
정규식 테이블 매핑
CoW / MoR
Paimon · Fluss
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) 두 가지 저장 전략을 제공한다.
| 구분 | CoW | MoR |
|---|---|---|
| 쓰기 방식 | 영향받은 파일 전체 재작성 | 델타 로그 파일만 추가 |
| 읽기 방식 | 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 이벤트:
ADD COLUMN— 컬럼 추가DROP COLUMN— 컬럼 삭제 (모드에 따라 동작이 다름)ALTER COLUMN TYPE— 타입 변경RENAME TABLE— 테이블 이름 변경
메타데이터 컬럼 지원도 추가되어 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 amountJSON_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.jarFlink 버전 호환성
3.6 파이프라인 커넥터가 지원하는 Flink 버전:
| Flink 버전 | 파이프라인 커넥터 지원 | JDK |
|---|---|---|
| 1.18.x 이하 | ✗ (지원 종료) | — |
| 1.19.x | ✗ | — |
| 1.20.x | ✓ | 11 이상 |
| 2.2.x | ✓ | 11 이상 |
3.6으로 업그레이드 시 JDK를 11 이상으로 맞춰야 한다. Java 8 환경은 더 이상 지원하지 않는다.
운영 체크리스트
Oracle 소스 준비
- [ ] 아카이브 로그 모드 확인:
SELECT LOG_MODE FROM V$DATABASE; - [ ] 보완 로깅 확인:
SELECT SUPPLEMENTAL_LOG_DATA_MIN FROM V$DATABASE; - [ ] LogMiner 사용자 계정 생성 및 접속 테스트
- [ ]
scan.startup.mode결정:initial(전체 스냅샷 후 스트리밍) vslatest-offset(현재 위치부터) - [ ] Oracle RAC 환경이면 모든 노드의 redo 로그 접근 가능 여부 확인
Hudi Sink 준비
- [ ] S3/HDFS 버킷·경로 생성, Flink 실행 계정의 IAM 권한 확인
- [ ]
table.type결정: CDC는 대부분MERGE_ON_READ - [ ] Compaction 전략 결정:
num_commits(커밋 횟수 기준) vstime_elapsed(시간 기준) - [ ] Hive Metastore / AWS Glue 동기화 설정 및 메타스토어 접속 확인
스키마 진화 모드 선택
- [ ]
LENIENT— 대부분의 CDC 파이프라인 기본값으로 권장 - [ ]
EVOLVE— 싱크가 모든 DDL을 반드시 반영해야 할 때 - [ ]
IGNORE— 싱크 스키마를 수동 관리할 때 - [ ]
EXCEPTION— 스키마 고정 환경에서 이상 감지용
모니터링
- [ ] Flink Web UI에서 Schema Operator의
schemaChangeEventCount확인 - [ ] 재시작 후 체크포인트에서 올바른 SCN으로 복구되는지 확인
- [ ] Hudi Compaction 백그라운드 작업 정상 동작 확인
요점 정리
- Oracle Source Pipeline: 3.6에서 YAML 파이프라인 커넥터로 승격. LogMiner 기반, 보완 로깅 필수.
- Hudi Sink Pipeline: CDC 워크로드는 MoR 권장.
compaction.delta_commits로 Compaction 빈도 제어. - PostgreSQL Schema Evolution: 3.6 이전에는 DDL이 무시되었으나 이제 전파 모드 선택 가능.
- VARIANT Transform: JSON 컬럼에서 중첩 필드를 SQL 수준에서 추출해 Hudi 테이블 구조를 단순하게 유지 가능.
- LENIENT 모드: DROP/TRUNCATE는 무시하고 ADD COLUMN은 전파. 프로덕션 CDC에 가장 안전한 기본값.
- JDK 11 필수: Flink CDC 3.6은 Java 8을 지원하지 않는다. 업그레이드 전 JDK 버전 확인 필수.