|
|
3 weeks ago | |
|---|---|---|
| .. | ||
| templates | 3 weeks ago | |
| BUILD-README.md | 3 weeks ago | |
| CUSTOM-README.md | 3 weeks ago | |
| Chart.yaml | 3 weeks ago | |
| README.md | 3 weeks ago | |
| custom-values.yaml | 3 weeks ago | |
| dip-questions.yaml | 3 weeks ago | |
| dip-values.yaml | 3 weeks ago | |
| values.yaml | 3 weeks ago | |
README.md
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 |
Flink 클러스터
| 파라미터 | 기본값 | 설명 |
|---|---|---|
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 nameflink-cdc-session→job.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