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.
 
 
 
 
 
 

9.7 KiB

Flink SQL Gateway 배포

1. 배포 방법

  • 명령어

    $ helm upgrade --install flink-sql-gateway ./ \
      --namespace flink-sql-test \
      --create-namespace \
      -f custom-values.yaml
    

2. custom-values.yaml 설명

  • custom-values.yaml에 정의된 값에 대한 설명이다.

1) 글로벌 설정

Name 설명 기본값
global.namespace 배포할 네임스페이스 flink-sql-test
global.image.repository Flink SQL Gateway 이미지 저장소 paasup/flink-sql
global.image.tag Flink SQL Gateway 이미지 태그 2.0.1
global.image.pullPolicy 이미지 풀 정책 IfNotPresent
Name 설명 기본값
sessionCluster.flinkVersion Flink 버전 v2_0
sessionCluster.jobManager.resources.memory JobManager 메모리 할당량 2048m
sessionCluster.jobManager.resources.cpu JobManager CPU 할당량 1
sessionCluster.taskManager.resources.memory TaskManager 메모리 할당량 4096m
sessionCluster.taskManager.resources.cpu TaskManager CPU 할당량 2
sessionCluster.flinkConfiguration.taskmanager.numberOfTaskSlots TaskManager 슬롯 수 "4"

3) SQL Gateway 설정

Name 설명 기본값
sqlGateway.resources.requests.memory SQL Gateway 요청 메모리 1Gi
sqlGateway.resources.requests.cpu SQL Gateway 요청 CPU 0.5
sqlGateway.resources.limits.memory SQL Gateway 제한 메모리 2Gi
sqlGateway.resources.limits.cpu SQL Gateway 제한 CPU 1

4) S3/MinIO 설정

Name 설명 기본값
sessionCluster.flinkConfiguration.fs.s3.impl S3 파일시스템 구현체 org.apache.hadoop.fs.s3a.S3AFileSystem
sessionCluster.flinkConfiguration.fs.s3a.impl S3A 파일시스템 구현체 org.apache.hadoop.fs.s3a.S3AFileSystem
sessionCluster.flinkConfiguration.fs.s3a.endpoint MinIO 엔드포인트 "https://minio.example.org"
sessionCluster.flinkConfiguration.fs.s3a.path.style.access Path Style Access 사용 여부 "true"

5) SQL Client 설정

Name 설명 기본값
sqlClient.enabled SQL Client 활성화 여부 true

3. 사설 인증서 등록

3.1 Kubernetes 클러스터

  • Kubernetes 1.20+ 권장
  • 충분한 리소스 (최소 10GB 메모리, 5 CPU 코어)

3.2 필수 Operator 설치

  • Flink Kubernetes Operator 설치 필요
    helm repo add flink-operator-repo https://downloads.apache.org/flink/flink-kubernetes-operator-1.8.0/
    helm install flink-kubernetes-operator flink-operator-repo/flink-kubernetes-operator
    

3.3 SSL/TLS 인증서 설정

  • 외부 시스템(Kafka, MinIO 등)과의 HTTPS 통신을 위한 사설 인증서 등록이 필요합니다.
  • Flink의 모든 컴포넌트(JobManager, TaskManager, SQL Gateway)에서 동일한 Truststore를 사용합니다.

3.3.1 인증서 파일 준비

포함되어야 할 인증서들:

Namespace Secret Name 파일 설명
kafka-cluster kafka-cluster-cluster-ca-cert ca.crt Kafka 클러스터 CA 인증서
cert-manager root-ca-secret ca.crt Ingress에 사용되는 Root CA 인증서

3.3.2 인증서 추출 및 Truststore 생성

# 1. Kafka CA 인증서 추출
kubectl get secret kafka-cluster-cluster-ca-cert -n kafka-cluster -o jsonpath='{.data.ca\.crt}' | base64 -d > kafka-ca.crt

# 2. Root CA 인증서 추출
kubectl get secret root-ca-secret -n cert-manager -o jsonpath='{.data.ca\.crt}' | base64 -d > root-ca.crt

# 3. 빈 truststore 생성
keytool -genkeypair -alias temp -keystore ca.p12 -storetype PKCS12 -storepass YOUR_TRUSTSTORE_PASSWORD -keypass YOUR_TRUSTSTORE_PASSWORD -dname "CN=temp" -keyalg RSA
keytool -delete -alias temp -keystore ca.p12 -storetype PKCS12 -storepass YOUR_TRUSTSTORE_PASSWORD

# 4. Kafka CA 인증서 추가
keytool -import -trustcacerts -keystore ca.p12 \
        -storetype PKCS12 \
        -storepass YOUR_TRUSTSTORE_PASSWORD \
        -alias kafka-ca \
        -file kafka-ca.crt \
        -noprompt

# 5. Root CA 인증서 추가
keytool -import -trustcacerts -keystore ca.p12 \
        -storetype PKCS12 \
        -storepass YOUR_TRUSTSTORE_PASSWORD \
        -alias root-ca \
        -file root-ca.crt \
        -noprompt

# 6. truststore 내용 확인
keytool -list -keystore ca.p12 -storetype PKCS12 -storepass YOUR_TRUSTSTORE_PASSWORD

# 7. Kubernetes Secret 생성
kubectl create secret generic truststore-secret \
  --from-file=ca.p12=ca.p12 \
  --from-literal=ca.password=YOUR_TRUSTSTORE_PASSWORD \
  -n flink-sql-test

3.3.3 custom-values.yaml에서 SSL 설정 적용

# Flink Session Cluster SSL 설정
sessionCluster:
  flinkConfiguration:
    # JVM SSL 시스템 속성 설정
    env.java.opts.jobmanager: "-Djavax.net.ssl.trustStore=/opt/flink/certs/ca.p12 -Djavax.net.ssl.trustStoreType=PKCS12 -Djavax.net.ssl.trustStorePassword=YOUR_TRUSTSTORE_PASSWORD"
    env.java.opts.taskmanager: "-Djavax.net.ssl.trustStore=/opt/flink/certs/ca.p12 -Djavax.net.ssl.trustStoreType=PKCS12 -Djavax.net.ssl.trustStorePassword=YOUR_TRUSTSTORE_PASSWORD"
  
  # 환경 변수 설정
  env:
    - name: TRUSTSTORE_PASSWORD
      valueFrom:
        secretKeyRef:
          name: truststore-secret
          key: ca.password
  
  # 볼륨 마운트 설정
  volumeMounts:
    - name: truststore-certs
      mountPath: /opt/flink/certs
      readOnly: true
  
  # 볼륨 정의
  volumes:
    - name: truststore-certs
      secret:
        secretName: truststore-secret

# SQL Gateway SSL 설정
sqlGateway:
  flinkConfiguration:
    # JVM SSL 시스템 속성 설정
    env.java.opts: "-Djavax.net.ssl.trustStore=/opt/flink/certs/ca.p12 -Djavax.net.ssl.trustStoreType=PKCS12 -Djavax.net.ssl.trustStorePassword=YOUR_TRUSTSTORE_PASSWORD"
  
  # 환경 변수, 볼륨 마운트, 볼륨 설정 (sessionCluster와 동일)
  env:
    - name: TRUSTSTORE_PASSWORD
      valueFrom:
        secretKeyRef:
          name: truststore-secret
          key: ca.password
  volumeMounts:
    - name: truststore-certs
      mountPath: /opt/flink/certs
      readOnly: true
  volumes:
    - name: truststore-certs
      secret:
        secretName: truststore-secret

# SQL Client SSL 설정 (선택적)
sqlClient:
  enabled: true
  env:
    - name: TRUSTSTORE_PASSWORD
      valueFrom:
        secretKeyRef:
          name: truststore-secret
          key: ca.password
  volumeMounts:
    - name: truststore-certs
      mountPath: /opt/flink/certs
      readOnly: true
  volumes:
    - name: truststore-certs
      secret:
        secretName: truststore-secret

3.3.4 SSL 설정 동작 원리

  1. Secret 생성: Truststore 파일(ca.p12)과 패스워드를 Kubernetes Secret으로 생성
  2. 볼륨 마운트: Secret을 각 컴포넌트의 /opt/flink/certs 경로에 마운트
  3. 환경 변수 주입: Secret의 패스워드를 TRUSTSTORE_PASSWORD 환경 변수로 주입
  4. JVM 옵션 설정: Java 시스템 속성을 통해 SSL Truststore 경로와 설정 지정
  5. SSL 통신: 외부 시스템(Kafka, MinIO 등)과의 HTTPS 통신 시 Truststore 사용

4. 배포 후 확인

4.1 Pod 상태 확인

kubectl get pods -n flink-sql-test

4.2 서비스 접근

# SQL Gateway API 접근을 위한 포트 포워딩
kubectl port-forward svc/flink-sql-gateway 8083:8083 -n flink-sql-test

# Flink Web UI 접근을 위한 포트 포워딩
kubectl port-forward svc/flink-session-cluster-rest 8081:8081 -n flink-sql-test

4.3 API 테스트

  • pod 내 실행
# SQL Gateway 정보 확인
curl http://localhost:8083/v1/info

# 세션 생성 테스트
curl -X POST http://localhost:8083/v1/sessions \
  -H "Content-Type: application/json" \
  -d '{"properties": {"execution.runtime-mode": "streaming"}}'

5. Python 클라이언트 사용 예제

import requests
import json

# SQL Gateway 연결
gateway_url = "http://localhost:8083"

# 1. 세션 생성
response = requests.post(f"{gateway_url}/v1/sessions")
session_handle = response.json()["sessionHandle"]
print(f"Session created: {session_handle}")

# 2. SQL 실행
sql_request = {
    "statement": "SHOW TABLES"
}
response = requests.post(
    f"{gateway_url}/v1/sessions/{session_handle}/statements",
    json=sql_request
)
operation_handle = response.json()["operationHandle"]

# 3. 결과 조회
result_response = requests.get(
    f"{gateway_url}/v1/sessions/{session_handle}/operations/{operation_handle}/result/0"
)
print("Query result:", result_response.json())

6. 주요 특징

  • 외부 Python 클라이언트 최적화: Port Forward를 통한 직접 REST API 접근
  • SSL/TLS 지원: 사설 인증서를 사용하는 외부 시스템 연동 지원
  • S3/MinIO 연동: 객체 스토리지 연동을 위한 설정 포함
  • 유연한 구성: 필요한 컴포넌트만 선택적 활성화 가능