flink-cdc-job Helm Chart
이미 실행 중인 Flink Session Cluster에 CDC 파이프라인 Job을 제출하는 chart입니다.
flink-cdc-session chart와 분리되어 있으므로 Session Cluster를 재시작하지 않고 파이프라인만 독립적으로 배포/재배포할 수 있습니다.
flink-cdc-session과의 역할 분리
| 항목 |
flink-cdc-session |
flink-cdc-job |
| FlinkDeployment (Session Cluster) |
✅ |
❌ |
| 파이프라인 ConfigMap |
✅ (클러스터용) |
✅ (Job 전용, 독립) |
| Job 제출 |
선택적 (job.enabled) |
항상 배포 |
| 배포 주기 |
클러스터 변경 시 |
파이프라인 변경 시 |
사전 조건
flink-cdc-session chart로 Session Cluster가 실행 중이어야 합니다.
파일 구조
배포
기본 배포
Release 이름이 Job 및 ConfigMap 이름으로 사용됩니다.
- Job:
cdc-job
- ConfigMap:
cdc-job-config
여러 파이프라인 동시 배포
각 파이프라인을 독립 release로 설치합니다.
파이프라인 재제출
Kubernetes Job은 완료 후 재실행이 불가하므로, 재제출 시 uninstall → install 합니다.
삭제
values.yaml 파라미터
이미지
| 파라미터 |
기본값 |
설명 |
image.repository |
wbsong111/flink-cdc-pipeline |
컨테이너 이미지 레지스트리/이름 |
image.tag |
3.5.0-flink1.20-r1 |
이미지 태그 |
image.pullPolicy |
Always |
이미지 풀 정책 |
serviceAccount |
flink |
Kubernetes ServiceAccount 이름 |
Session Cluster 연결
| 파라미터 |
기본값 |
설명 |
sessionCluster.restAddress |
flink-cdc-session-rest |
Session Cluster REST 서비스명 |
sessionCluster.restPort |
8081 |
REST 포트 |
restAddress는 FlinkDeployment 이름에 -rest를 붙인 값입니다.
예: FlinkDeployment flink-cdc-session → 서비스명 flink-cdc-session-rest
Flink CDC 글로벌 설정
| 파라미터 |
기본값 |
설명 |
flinkCdc.parallelism |
1 |
전체 파이프라인 기본 병렬도 |
flinkCdc.schemaChangeBehavior |
EVOLVE |
스키마 변경 처리 방식 (EVOLVE / IGNORE / EXCEPTION) |
파이프라인
| 파라미터 |
기본값 |
설명 |
pipeline.name |
cdc-pipeline |
Flink Job 이름 (Flink UI 표시명) |
pipeline.parallelism |
2 |
이 파이프라인의 병렬도 |
pipeline.checkpointInterval |
60s |
체크포인트 주기 |
PostgreSQL 소스
| 파라미터 |
기본값 |
설명 |
postgres.hostname |
postgres-postgresql.flink-test.svc.cluster.local |
PostgreSQL 호스트 |
postgres.port |
5432 |
PostgreSQL 포트 |
postgres.username |
flink_cdc |
CDC 전용 사용자 |
postgres.password |
password |
사용자 비밀번호 |
postgres.slotName |
flink_cdc_slot |
Replication Slot 이름 (파이프라인마다 고유해야 함) |
postgres.decodingPlugin |
pgoutput |
디코딩 플러그인 (pgoutput / decoderbufs) |
postgres.tables |
testdb.public.orders,testdb.public.users |
CDC 대상 테이블 목록 (database.schema.table 형식) |
Iceberg REST 카탈로그 싱크
| 파라미터 |
기본값 |
설명 |
sink.catalog.uri |
https://lakekeeper.example.org/catalog |
Iceberg REST 카탈로그 URI |
sink.catalog.warehouse |
minio |
카탈로그 warehouse 이름 |
sink.catalog.s3Endpoint |
https://minio.example.org |
S3 호환 스토리지 엔드포인트 |
sink.catalog.s3PathStyleAccess |
"true" |
Path-style 접근 여부 |
sink.catalog.oauth2Uri |
https://keycloak.example.org/... |
OAuth2 토큰 엔드포인트 |
sink.catalog.credential |
lakekeeper-admin:... |
OAuth2 client_id:client_secret |
sink.catalog.scope |
lakekeeper |
OAuth2 scope |
라우팅
소스 테이블과 싱크 테이블 매핑입니다.
values 파일 예시
values-users.yaml (사용자/상품 테이블)
values-orders.yaml (주문 테이블)
렌더링 확인
배포 후 확인