Repository for dip
You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

18 KiB

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
    ├── ingress.yaml         # Ingress (조건부)
    ├── job.yaml             # CDC 파이프라인 제출 Job (조건부)
    ├── role.yaml            # Role (rbac.create: true 시 생성)
    ├── rolebinding.yaml     # RoleBinding (rbac.create: true 시 생성)
    ├── secrets.yaml         # truststore / gcs-key Secret (조건부)
    └── serviceaccount.yaml  # ServiceAccount (serviceAccount.create: true 시 생성)

리소스 명명 규칙

모든 리소스 이름은 Helm Release 이름을 기반으로 동적 생성됩니다.

리소스 이름 패턴
FlinkDeployment {release-name}-session
ConfigMap {release-name}-pipeline-config
ServiceAccount {release-name}
Role / RoleBinding {release-name}
Truststore Secret (chart 생성 시) {release-name}-flink-truststore

예: helm install flink-cdc-session . → FlinkDeployment 이름 flink-cdc-session-session


외부 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 검증이 가능하며, 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="paasup/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 생성

chart가 자동 생성하는 Secret 이름은 {release-name}-flink-truststore입니다.
기존 Secret을 재사용하려면 truststore.secretName에 이름을 지정합니다.

# chart 명명 규칙에 맞게 생성 (release name: flink-cdc-session)
kubectl create secret generic flink-cdc-session-flink-truststore \
  --from-file=truststore.jks=./truststore.jks \
  -n flink-test

# 또는 기존 Secret 이름 그대로 사용 (custom-values.yaml에서 secretName 지정)
kubectl create secret generic flink-truststore \
  --from-file=truststore.jks=./truststore.jks \
  -n flink-test

기존 Secret 업데이트:

kubectl create secret generic flink-cdc-session-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-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 기동 확인 (LIFECYCLE STATE가 STABLE이 될 때까지 대기)
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 Secret 재사용

네임스페이스에 이미 truststore Secret이 있는 경우 truststore.secretName으로 지정합니다.

# custom-values.yaml
truststore:
  enabled: true
  secretName: flink-truststore  # 기존 Secret 이름

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 생성 (이름: {release-name}-flink-truststore)
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 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 권한:

리소스 권한
pods get, list, watch, create, delete, patch
configmaps get, list, watch, create, update, patch, delete
services get, list, create, delete
events create, patch
apps/deployments get, list, watch, create, update, patch, delete
apps/replicasets get, list, watch, create, update, patch, delete
파라미터 기본값 설명
flinkVersion v1_20 Flink 버전
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 false 외부 JKS 사용 여부. 사설 CA 환경이면 true
truststore.secretName "" 기존 Secret 이름 지정. 비어있으면 {release-name}-flink-truststore 사용
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 이름 suffix. 전체 이름: {release-name}-{job.name}
job.restAddress flink-cdc-session-rest Session Cluster REST 서비스명
job.restPort 8081 REST 포트
job.pipelineFile ./conf/pipeline.yaml 파이프라인 설정 파일 경로

job.restAddress는 FlinkDeployment 이름에 -rest를 붙인 값입니다.
이 chart가 생성하는 FlinkDeployment 이름은 {release-name}-session이므로
REST 서비스명은 {release-name}-session-rest가 됩니다.
예: release name flink-cdc-sessionjob.restAddress: flink-cdc-session-session-rest


렌더링 확인

# 전체 렌더링 확인
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 상태 확인 (release name: flink-cdc-session)
kubectl get flinkdeployment -n flink-test
# NAME                        JOB STATUS   LIFECYCLE STATE
# flink-cdc-session-session                STABLE

# JobManager 로그 확인
kubectl logs -n flink-test -l app.kubernetes.io/instance=flink-cdc-session

# Flink REST API로 Job 상태 확인
kubectl exec -n flink-test <jobmanager-pod> -- curl -s http://localhost:8081/jobs

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