# 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가 그대로 사용됩니다. ### 엔드포인트 인증서 확인 방법 로컬 또는 클러스터 내부에서 아래 명령으로 인증서 발급자를 확인합니다. ```bash # 인증서 발급자(Issuer) 확인 openssl s_client -connect lakekeeper.example.org:443 -showcerts /dev/null \ | openssl x509 -noout -issuer -subject openssl s_client -connect keycloak.example.org:443 -showcerts /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 ``` 클러스터 내부에서 확인하는 경우: ```bash kubectl run ssl-check --rm -it --image=alpine/openssl --restart=Never -- \ s_client -connect lakekeeper.example.org:443 -showcerts /dev/null \ | grep -E "issuer|subject" ``` --- ## Secret 준비 ### 1. truststore.jks 생성 (`truststore.enabled: true` 인 경우만 필요) #### Step 1 — JVM 기본 cacerts 추출 Docker 이미지 내부의 JVM cacerts를 truststore 베이스로 사용합니다. ```bash 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에서 추출하는 경우:** ```bash 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 파일에서 직접 추가하는 경우:** ```bash keytool -import -noprompt \ -keystore truststore.jks \ -storepass changeit \ -alias my-ca \ -file ./certs/my-ca.crt ``` 추가된 인증서 목록 확인: ```bash keytool -list -keystore truststore.jks -storepass changeit | tail -5 ``` #### Step 3 — truststore 동작 검증 배포 전에 truststore로 실제 엔드포인트 TLS 연결이 성공하는지 확인합니다. ```bash # 사설 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`에 이름을 지정합니다. ```bash # 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 업데이트: ```bash 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 > bash ../../create-truststore.sh > ``` #### 배포 후 마운트 확인 ```bash # 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 방식 (권장) ```bash kubectl create secret generic minio-credentials \ --from-literal=access-key-id= \ --from-literal=secret-access-key= \ -n flink-test ``` `custom-values.yaml`에서 참조: ```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 — 서비스 계정 키 방식 ```bash kubectl create secret generic gcs-key \ --from-file=key.json=./gcs_key/key.json \ -n flink-test ``` `custom-values.yaml`에서 참조: ```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`에 오버라이드합니다. ```bash helm install flink-cdc-session ./flink/helm/flink-cdc-session \ -n flink-test \ -f values.yaml \ -f custom-values.yaml ``` ### Session Cluster만 먼저 배포 (Job 제외) ```bash # 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`으로 지정합니다. ```yaml # custom-values.yaml truststore: enabled: true secretName: flink-truststore # 기존 Secret 이름 ``` ### truststore 불필요한 경우 (공인 CA 환경) ```bash 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에서 생성하는 경우 ```bash # 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)" ``` ### 삭제 ```bash 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 | ### 라우팅 ```yaml 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-session` → `job.restAddress: flink-cdc-session-session-rest` --- ## 렌더링 확인 ```bash # 전체 렌더링 확인 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" ``` ## 배포 후 확인 ```bash # 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 -- 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 ```