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가 실행 중이어야 합니다.
파일 구조
flink-cdc-job/
├── Chart.yaml
├── values.yaml
└── templates/
├── _helpers.tpl
├── configmap.yaml # 이 Job 전용 pipeline.yaml (Release.Name 기반 네이밍)
├── job.yaml # Kubernetes Job
├── serviceaccount.yaml # ServiceAccount (serviceAccount.create: true 시 생성)
├── role.yaml # Role (rbac.create: true 시 생성)
└── rolebinding.yaml # RoleBinding (rbac.create: true 시 생성)
배포
기본 배포
helm install cdc-job ./flink/helm/flink-cdc-job -n flink-test
Release 이름이 Job, ConfigMap, ServiceAccount, Role 이름으로 사용됩니다.
- Job:
cdc-job
- ConfigMap:
cdc-job-config
- ServiceAccount:
cdc-job
- Role / RoleBinding:
cdc-job
여러 파이프라인 동시 배포
각 파이프라인을 독립 release로 설치합니다.
# 사용자 테이블 파이프라인
helm install cdc-users ./flink/helm/flink-cdc-job -n flink-test \
-f values-users.yaml
# 주문 테이블 파이프라인
helm install cdc-orders ./flink/helm/flink-cdc-job -n flink-test \
-f values-orders.yaml
파이프라인 재제출
Kubernetes Job은 완료 후 재실행이 불가하므로, 재제출 시 uninstall → install 합니다.
helm uninstall cdc-job -n flink-test
helm install cdc-job ./flink/helm/flink-cdc-job -n flink-test
삭제
helm uninstall cdc-job -n flink-test
values.yaml 파라미터
이미지
| 파라미터 |
기본값 |
설명 |
image.repository |
paasup/flink-cdc-pipeline |
컨테이너 이미지 레지스트리/이름 |
image.tag |
3.5.0-flink1.20-r1 |
이미지 태그 |
image.pullPolicy |
Always |
이미지 풀 정책 |
ServiceAccount
| 파라미터 |
기본값 |
설명 |
serviceAccount.create |
true |
ServiceAccount 생성 여부 |
serviceAccount.name |
"" |
SA 이름. 비어있으면 Release.Name 사용 |
RBAC
| 파라미터 |
기본값 |
설명 |
rbac.create |
true |
Role / RoleBinding 생성 여부 |
rbac.create: true 시 생성되는 Role 권한:
| 리소스 |
권한 |
configmaps |
get, list |
pods |
get, list |
Session Cluster 연결
| 파라미터 |
기본값 |
설명 |
sessionCluster.restAddress |
flink-cdc-session-rest |
Session Cluster REST 서비스명 |
sessionCluster.restPort |
8081 |
REST 포트 |
restAddress는 flink-cdc-session chart의 FlinkDeployment 이름에 -rest를 붙인 값입니다.
flink-cdc-session chart는 FlinkDeployment를 {release-name}-session으로 생성하므로
REST 서비스명은 {release-name}-session-rest가 됩니다.
예: release name flink-cdc-session → sessionCluster.restAddress: flink-cdc-session-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 |
라우팅
소스 테이블과 싱크 테이블 매핑입니다.
route:
- sourceTable: public.orders # PostgreSQL schema.table
sinkTable: testdb_public.orders # Iceberg namespace.table
values 파일 예시
values-users.yaml (사용자/상품 테이블)
sessionCluster:
restAddress: flink-cdc-session-session-rest
pipeline:
name: cdc-users
parallelism: 2
checkpointInterval: 60s
postgres:
hostname: postgres.example.svc.cluster.local
port: 5432
username: flink_cdc
password: "password"
slotName: flink_cdc_users_slot
decodingPlugin: pgoutput
tables: mydb.public.users,mydb.public.products
route:
- sourceTable: public.users
sinkTable: mydb_public.users
- sourceTable: public.products
sinkTable: mydb_public.products
values-orders.yaml (주문 테이블)
sessionCluster:
restAddress: flink-cdc-session-session-rest
pipeline:
name: cdc-orders
parallelism: 4
checkpointInterval: 60s
postgres:
hostname: postgres.example.svc.cluster.local
slotName: flink_cdc_orders_slot
tables: mydb.public.orders,mydb.public.order_items
route:
- sourceTable: public.orders
sinkTable: mydb_public.orders
- sourceTable: public.order_items
sinkTable: mydb_public.order_items
렌더링 확인
helm template cdc-job . -n flink-test
helm lint .
# 여러 release 이름으로 리소스 이름이 분리되는지 확인
helm template cdc-users . -n flink-test | grep "name:"
helm template cdc-orders . -n flink-test | grep "name:"
배포 후 확인
# Job 상태 확인
kubectl get job -n flink-test
# Job 로그 확인
kubectl logs -n flink-test job/cdc-job
# Flink REST API로 Job 상태 확인
kubectl exec -n flink-test <jobmanager-pod> -- curl -s http://localhost:8081/jobs