Files
service-catalog/manifests/helm/flink-cdc-session/0.1.0

flink-cdc-session Helm Chart

PostgreSQL → Apache Iceberg CDC 파이프라인을 Flink Session Mode로 배포하는 Helm chart입니다.

아키텍처

PostgreSQL (WAL)
    │  pgoutput
    ▼
Flink Session Cluster  (FlinkDeployment)
    │  Flink CDC 3.5.0
    ▼
Iceberg REST Catalog (Lakekeeper)
    │  S3FileIO
    ▼
Object Storage (MinIO / GCS / S3)

사전 요구사항

항목 버전
Helm 3.x
Flink Kubernetes Operator 1.x
Kubernetes 1.25+

파일 구조

flink-cdc-session/
├── Chart.yaml
├── values.yaml
├── custom-values.yaml       # 환경별 오버라이드 (배포 시 사용)
└── templates/
    ├── _helpers.tpl
    ├── configmap.yaml       # flink-cdc.yaml + pipeline.yaml
    ├── flinkdeployment.yaml # Flink Session Cluster
    ├── job.yaml             # CDC 파이프라인 제출 Job (조건부)
    └── secrets.yaml         # truststore / gcs-key Secret (조건부)

외부 JKS(truststore) 사용 여부 결정

이 chart는 truststore.enabled 플래그로 외부 JKS 사용 여부를 제어합니다.
배포 전에 반드시 아래 기준에 따라 값을 결정하세요.

판단 기준

Flink는 다음 엔드포인트와 HTTPS로 통신합니다.

엔드포인트 용도
Lakekeeper (sink.catalog.uri) Iceberg REST 카탈로그
Keycloak (sink.catalog.oauth2Uri) OAuth2 토큰 발급
MinIO (sink.catalog.s3Endpoint) Iceberg 데이터 파일 읽기/쓰기

truststore.enabled: true (외부 JKS 사용)
→ 위 엔드포인트 중 하나라도 사설 CA가 서명한 인증서를 사용하는 경우

truststore.enabled: false (외부 JKS 불필요)
→ 모든 엔드포인트가 공인 CA(Let's Encrypt, DigiCert 등)로 서명된 인증서를 사용하는 경우
→ JVM 기본 cacerts만으로 TLS 검증이 가능하며, flink-truststore Secret도 불필요

왜 외부 JKS가 필요한가
Flink JVM 옵션 -Djavax.net.ssl.trustStore는 JVM 기본 cacerts를 완전히 교체합니다.
따라서 사설 CA 인증서뿐만 아니라 공인 CA 체인도 포함된 truststore를 직접 만들어야 합니다.
반대로 이 옵션을 설정하지 않으면 JVM 기본 cacerts가 그대로 사용됩니다.

엔드포인트 인증서 확인 방법

로컬 또는 클러스터 내부에서 아래 명령으로 인증서 발급자를 확인합니다.

# 인증서 발급자(Issuer) 확인
openssl s_client -connect lakekeeper.example.org:443 -showcerts </dev/null 2>/dev/null \
  | openssl x509 -noout -issuer -subject

openssl s_client -connect keycloak.example.org:443 -showcerts </dev/null 2>/dev/null \
  | openssl x509 -noout -issuer -subject

출력 예시:

# 사설 CA → truststore.enabled: true 필요
issuer=O=paasup, CN=paasup-root-ca
subject=CN=lakekeeper.example.org

# 공인 CA → truststore.enabled: false 가능
issuer=C=US, O=Let's Encrypt, CN=R11
subject=CN=lakekeeper.example.org

클러스터 내부에서 확인하는 경우:

kubectl run ssl-check --rm -it --image=alpine/openssl --restart=Never -- \
  s_client -connect lakekeeper.example.org:443 -showcerts </dev/null 2>/dev/null \
  | grep -E "issuer|subject"

Secret 준비

1. truststore.jks 생성 (truststore.enabled: true 인 경우만 필요)

Step 1 — JVM 기본 cacerts 추출

Docker 이미지 내부의 JVM cacerts를 truststore 베이스로 사용합니다.

IMAGE="wbsong111/flink-cdc-pipeline:3.5.0-flink1.20-r1"

docker run --rm --entrypoint="" "${IMAGE}" \
  cat /opt/java/openjdk/lib/security/cacerts > truststore.jks

Step 2 — 사설 CA 추가

cert-manager ClusterIssuer에서 추출하는 경우:

kubectl get secret root-ca-secret -n cert-manager \
  -o jsonpath='{.data.ca\.crt}' | base64 -d > /tmp/internal-ca.pem

keytool -import -noprompt \
  -keystore truststore.jks \
  -storepass changeit \
  -alias internal-root-ca \
  -file /tmp/internal-ca.pem

rm -f /tmp/internal-ca.pem

PEM 파일에서 직접 추가하는 경우:

keytool -import -noprompt \
  -keystore truststore.jks \
  -storepass changeit \
  -alias my-ca \
  -file ./certs/my-ca.crt

추가된 인증서 목록 확인:

keytool -list -keystore truststore.jks -storepass changeit | tail -5

Step 3 — truststore 동작 검증

배포 전에 truststore로 실제 엔드포인트 TLS 연결이 성공하는지 확인합니다.

# 사설 CA PEM으로 curl 검증
curl -v --cacert /tmp/internal-ca.pem https://lakekeeper.example.org/catalog 2>&1 \
  | grep -E "SSL|TLS|certificate"

Step 4 — Kubernetes Secret 생성

kubectl create secret generic flink-truststore \
  --from-file=truststore.jks=./truststore.jks \
  -n flink-test

기존 Secret 업데이트:

kubectl create secret generic flink-truststore \
  --from-file=truststore.jks=./truststore.jks \
  -n flink-test \
  --dry-run=client -o yaml | kubectl apply -f -

create-truststore.sh 스크립트로 Step 1~4를 한번에 실행할 수 있습니다:

bash ../../create-truststore.sh

배포 후 마운트 확인

# truststore.jks가 Pod에 정상 마운트되었는지 확인
kubectl exec -n flink-test deploy/flink-cdc-session -- ls -la /opt/flink/certs/

# JVM이 truststore를 인식하는지 확인 (JobManager 로그에서 SSL 관련 출력)
kubectl logs -n flink-test -l component=jobmanager | grep -i "truststore\|ssl\|tls" | head -20

2. 체크포인트 저장소 자격증명

checkpointStorage.storageType 값에 따라 필요한 Secret이 다릅니다.

MinIO / S3 — existingSecret 방식 (권장)

kubectl create secret generic minio-credentials \
  --from-literal=access-key-id=<ACCESS_KEY> \
  --from-literal=secret-access-key=<SECRET_KEY> \
  -n flink-test

custom-values.yaml에서 참조:

checkpointStorage:
  storageType: s3
  checkpointDir: s3://flink/checkpoints
  savepointDir: s3://flink/savepoints
  s3:
    endpoint: "http://minio.minio.svc.cluster.local:9000"  # 클러스터 내부 MinIO
    pathStyleAccess: "true"
    existingSecret: minio-credentials

클러스터 내부 MinIO를 사용할 경우 plain HTTP 엔드포인트(http://)를 사용하면 TLS 없이 통신할 수 있습니다.

GCS — 서비스 계정 키 방식

kubectl create secret generic gcs-key \
  --from-file=key.json=./gcs_key/key.json \
  -n flink-test

custom-values.yaml에서 참조:

checkpointStorage:
  storageType: gcs
  checkpointDir: gs://my-bucket/flink/checkpoints
  savepointDir: gs://my-bucket/flink/savepoints
  gcs:
    credentialSecret: gcs-key
    credentialPath: /opt/flink/gcs-key/key.json

배포

custom-values.yaml 사용 (권장)

환경별 설정은 values.yaml을 수정하지 않고 custom-values.yaml에 오버라이드합니다.

helm install flink-cdc-session ./flink/helm/flink-cdc-session \
  -n flink-test \
  -f values.yaml \
  -f custom-values.yaml

Session Cluster만 먼저 배포 (Job 제외)

# 1단계: Session Cluster만 배포 (custom-values.yaml에 job.enabled: false 설정)
helm install flink-cdc-session ./flink/helm/flink-cdc-session \
  -n flink-test \
  -f values.yaml \
  -f custom-values.yaml

# 2단계: Cluster 기동 확인
kubectl get flinkdeployment -n flink-test -w

# 3단계: Job 포함해서 업그레이드
helm upgrade flink-cdc-session ./flink/helm/flink-cdc-session \
  -n flink-test \
  -f values.yaml \
  -f custom-values.yaml \
  --set job.enabled=true

truststore 불필요한 경우 (공인 CA 환경)

helm install flink-cdc-session ./flink/helm/flink-cdc-session \
  -n flink-test \
  -f values.yaml \
  -f custom-values.yaml \
  --set truststore.enabled=false

truststore.enabled=false 시 다음이 모두 제거됩니다:

  • env.java.opts.jobmanager/taskmanager의 SSL JVM 옵션
  • Pod의 truststore volumeMount 및 volume

Secret도 chart에서 생성하는 경우

# truststore Secret chart 생성
helm install flink-cdc-session ./flink/helm/flink-cdc-session \
  -n flink-test \
  --set truststore.create=true \
  --set truststore.base64jks="$(base64 < truststore.jks)"

# GCS Secret chart 생성
helm install flink-cdc-session ./flink/helm/flink-cdc-session \
  -n flink-test \
  --set checkpointStorage.storageType=gcs \
  --set checkpointStorage.gcs.create=true \
  --set "checkpointStorage.gcs.base64json=$(base64 < gcs_key/key.json)"

삭제

helm uninstall flink-cdc-session -n flink-test

values.yaml 파라미터

이미지

파라미터 기본값 설명
image.repository wbsong111/flink-cdc-pipeline 컨테이너 이미지 저장소
image.tag 3.5.0-flink1.20-r1 이미지 태그
image.pullPolicy Always 이미지 풀 정책
파라미터 기본값 설명
flinkVersion v1_20 Flink 버전
serviceAccount flink Pod 서비스 어카운트
jobManager.replicas 1 JobManager 인스턴스 수
jobManager.memory 1024m JobManager 메모리
jobManager.cpu 0.5 JobManager CPU
taskManager.replicas 3 TaskManager 인스턴스 수
taskManager.memory 2048m TaskManager 메모리
taskManager.cpu 1 TaskManager CPU
taskManager.taskSlots 2 TaskManager 슬롯 수

CheckpointStorage (저장 위치 + 자격증명)

파라미터 기본값 설명
checkpointStorage.storageType s3 저장소 유형 (gcs | s3)
checkpointStorage.checkpointDir s3://my-bucket/… 체크포인트 저장 경로
checkpointStorage.savepointDir s3://my-bucket/… 세이브포인트 저장 경로
checkpointStorage.s3.endpoint "" S3/MinIO 엔드포인트. 비워두면 AWS S3
checkpointStorage.s3.pathStyleAccess "false" Path-style 접근 여부. MinIO는 "true"
checkpointStorage.s3.existingSecret "" 자격증명 Secret 이름. 설정 시 아래 두 항목보다 우선
checkpointStorage.s3.accessKeyIdKey access-key-id Secret 내 Access Key ID 필드명
checkpointStorage.s3.secretAccessKeyKey secret-access-key Secret 내 Secret Access Key 필드명
checkpointStorage.s3.accessKeyId "" Access Key ID 직접 입력 (테스트용)
checkpointStorage.s3.secretAccessKey "" Secret Access Key 직접 입력 (테스트용)
checkpointStorage.gcs.credentialSecret gcs-key GCS 키 Secret 이름
checkpointStorage.gcs.credentialPath /opt/flink/gcs-key/key.json 컨테이너 내 키 파일 경로
checkpointStorage.gcs.create false true이면 chart가 Secret 생성
checkpointStorage.gcs.base64json "" create: true 시 필수, base64 인코딩된 key.json

체크포인트 동작 설정

파라미터 기본값 설명
checkpoint.interval 60s 체크포인트 주기
checkpoint.mode EXACTLY_ONCE 체크포인트 모드

Truststore (외부 JKS)

파라미터 기본값 설명
truststore.enabled true 외부 JKS 사용 여부. 사설 CA 환경이면 true
truststore.secretName flink-truststore truststore Secret 이름
truststore.password changeit JKS 패스워드
truststore.create false true이면 chart가 Secret 생성
truststore.base64jks "" create: true 시 필수, base64 인코딩된 .jks

PostgreSQL 소스

파라미터 기본값 설명
postgres.hostname postgres-postgresql.flink-test.svc… PostgreSQL 호스트
postgres.port 5432 PostgreSQL 포트
postgres.username flink_cdc CDC 전용 사용자
postgres.password password 사용자 패스워드
postgres.slotName flink_cdc_slot Replication slot 이름
postgres.decodingPlugin pgoutput WAL 디코딩 플러그인
postgres.tables testdb.public.orders,… CDC 대상 테이블 목록

파이프라인

파라미터 기본값 설명
flinkCdc.parallelism 1 Flink CDC 글로벌 병렬도
flinkCdc.schemaChangeBehavior EVOLVE 스키마 변경 처리 방식
pipeline.name pg-to-iceberg-test 파이프라인 이름
pipeline.parallelism 2 파이프라인 병렬도
pipeline.checkpointInterval 60s 파이프라인 체크포인트 주기

Iceberg 싱크

파라미터 기본값 설명
sink.catalog.uri https://lakekeeper.example.org/catalog Lakekeeper REST URI
sink.catalog.warehouse minio 카탈로그 warehouse 이름
sink.catalog.s3Endpoint https://minio.example.org S3 호환 오브젝트 스토리지 엔드포인트
sink.catalog.s3PathStyleAccess "true" Path-style S3 접근 여부
sink.catalog.oauth2Uri https://keycloak.example.org/… Keycloak OAuth2 토큰 URI
sink.catalog.credential lakekeeper-admin:… client_id:client_secret
sink.catalog.scope lakekeeper OAuth2 scope

라우팅

route:
  - sourceTable: public.orders      # PostgreSQL schema.table
    sinkTable: testdb_public.orders # Iceberg namespace.table
  - sourceTable: public.users
    sinkTable: testdb_public.users

Job

파라미터 기본값 설명
job.enabled true false이면 Job 리소스 미생성
job.name submit-cdc Job 리소스 이름
job.restAddress flink-cdc-session-rest Session Cluster REST 서비스명
job.restPort 8081 REST 포트
job.pipelineFile ./conf/pipeline.yaml 파이프라인 설정 파일 경로

렌더링 확인

# 전체 렌더링 확인
helm template flink-cdc-session . -n flink-test \
  -f values.yaml -f custom-values.yaml

# lint 검사
helm lint .

# S3/MinIO 렌더링 확인
helm template flink-cdc-session . -n flink-test \
  -f values.yaml -f custom-values.yaml \
  | grep -E "checkpointDir|savepointDir|AWS_ACCESS|s3.endpoint"

# truststore 비활성화 시 SSL 옵션이 제거되는지 확인
helm template flink-cdc-session . -n flink-test \
  -f values.yaml -f custom-values.yaml \
  --set truststore.enabled=false | grep -A3 "java.opts"

배포 후 확인

# Session Cluster 상태 확인
kubectl get flinkdeployment -n flink-test

# Job 로그 확인
kubectl logs -n flink-test job/submit-cdc

# Job 삭제 후 재제출
kubectl delete job submit-cdc -n flink-test
helm upgrade flink-cdc-session . -n flink-test \
  -f values.yaml -f custom-values.yaml \
  --set job.enabled=true