# Flink SQL Gateway Helm Chart Apache Flink SQL Gateway를 쉽게 배포하기 위한 Helm Chart입니다. 외부 Python 클라이언트 접근에 최적화되어 있습니다. ## 📋 개요 이 Helm Chart는 다음 구성 요소를 배포합니다: ### 필수 컴포넌트 (항상 배포됨) - **Flink Session Cluster**: SQL Gateway가 연결할 Flink 클러스터 - **Flink SQL Gateway**: SQL 쿼리를 실행할 수 있는 REST API 서버 - **Service**: SQL Gateway 접근을 위한 Kubernetes 서비스 - **RBAC**: Flink 리소스 관리를 위한 권한 설정 - **ConfigMap**: SQL Gateway 설정 ### 선택적 컴포넌트 (설정으로 활성화) - **SQL Client Pod**: 대화형 SQL 클라이언트 (`sqlClient.enabled: true`) - **Kafka Integration**: Kafka 연동 설정 (`kafka.enabled: true`) - **SQL Scripts ConfigMap**: SQL 스크립트 모음 (SQL Client 활성화 시) ## 🏗️ 아키텍처 ### 외부 Python 클라이언트 접근 (권장) ``` ┌──────────────────┐ │ Python Client │ │ (External) │ └──────────────────┘ │ ▼ HTTP REST API ┌──────────────────┐ │ SQL Gateway │ │ (REST API) │ │ Port: 8083 │ └──────────────────┘ │ ▼ SQL 실행 요청 ┌─────────────────────┐ │ Flink Session │ │ Cluster │ │ (Data Processing) │ └─────────────────────┘ │ ▼ 데이터 처리 (선택적) ┌─────────────────┐ │ Kafka Cluster │ │ (Streaming) │ └─────────────────┘ ``` ### 내부 SQL Client 사용 (선택적) ``` ┌──────────────────┐ │ SQL Client │ │ (Pod) │ └──────────────────┘ │ ▼ SQL 쿼리 전송 ┌──────────────────┐ │ SQL Gateway │ │ (REST API) │ └──────────────────┘ │ ▼ SQL 실행 요청 ┌─────────────────────┐ │ Flink Session │ │ Cluster │ └─────────────────────┘ ``` **주요 특징:** - **외부 Python 클라이언트 최적화**: Port Forward를 통한 직접 REST API 접근 - **유연한 구성**: 필요한 컴포넌트만 선택적 배포 - **Session Mode**: 대화형 SQL 개발 환경 제공 - **검증된 안정성**: YAML 파싱 오류 해결 및 템플릿 검증 완료 ## 🚀 빠른 시작 ### 1. 사전 요구사항 - Kubernetes 클러스터 - Helm 3.x - Flink Kubernetes Operator 설치됨 - Strimzi Kafka Operator 설치됨 (Kafka 사용 시) ### 2. 설치 #### 기본 설치 (외부 Python 클라이언트용 - 권장) ```bash # 최소 리소스로 기본 설치 helm install flink-sql-gateway ./ \ --namespace flink-sql-gateway \ --create-namespace # 외부 Python 클라이언트용 최적화 설치 (권장) helm install flink-sql-gateway ./ \ --namespace flink-sql-gateway \ --create-namespace \ --values ./custom-values.yaml ``` #### 고급 설치 옵션 ```bash # SQL Client 포함 설치 (대화형 SQL 사용) helm install flink-sql-gateway ./ \ --namespace flink-sql-gateway \ --create-namespace \ --set sqlClient.enabled=true # Kafka 연동 포함 설치 helm install flink-sql-gateway ./ \ --namespace flink-sql-gateway \ --create-namespace \ --set kafka.enabled=true \ --set kafka.user.password="your-kafka-password" # 완전한 환경 설치 (모든 컴포넌트) helm install flink-sql-gateway ./ \ --namespace flink-sql-gateway \ --create-namespace \ --values ./custom-values.yaml \ --set sqlClient.enabled=true \ --set kafka.enabled=true \ --set kafka.user.password="your-kafka-password" ``` ### 3. 외부 Python 클라이언트 접근 (권장) #### Port Forward 설정 ```bash # SQL Gateway API 접근 kubectl port-forward -n flink-sql-gateway svc/flink-sql-gateway 8083:8083 # Flink Web UI 접근 (선택적) kubectl port-forward -n flink-sql-gateway svc/flink-session-cluster-rest 8081:8081 ``` #### Python 클라이언트 예제 ```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()) ``` ### 4. 내부 SQL Client 접근 (선택적) ```bash # SQL Client에 연결 (sqlClient.enabled=true인 경우) kubectl exec -it flink-sql-client -n flink-sql-gateway -- \ /opt/flink/bin/sql-client.sh gateway \ --endpoint http://flink-sql-gateway:8083 ``` ## ⚙️ 설정 ### 주요 설정 항목 | 파라미터 | 설명 | 기본값 | custom-values.yaml | |----------|------|--------|-------------------| | `global.namespace` | 배포할 네임스페이스 | `flink-sql-gateway` | `flink-sql-gateway` | | `global.image.repository` | Flink 이미지 저장소 | `paasup/flink-sql` | `paasup/flink-sql` | | `global.image.tag` | Flink 이미지 태그 | `1.20-kafka` | `1.20-kafka` | | `sessionCluster.enabled` | Session Cluster 활성화 | `true` | `true` | | `sqlGateway.enabled` | SQL Gateway 활성화 | `true` | `true` | | `sqlClient.enabled` | SQL Client 활성화 | `false` | `false` | | `kafka.enabled` | Kafka 연동 활성화 | `false` | `false` | ### 리소스 설정 비교 #### 기본 설정 (values.yaml) ```yaml # 최소 리소스 - 테스트용 sessionCluster: jobManager: resources: memory: 1024m cpu: 0.5 taskManager: resources: memory: 2048m cpu: 1 flinkConfiguration: taskmanager.numberOfTaskSlots: "2" sqlGateway: resources: requests: memory: 512Mi cpu: 0.25 limits: memory: 1Gi cpu: 0.5 ``` #### 최적화 설정 (custom-values.yaml) ```yaml # 외부 Python 클라이언트용 최적화 - 프로덕션 권장 sessionCluster: jobManager: resources: memory: 2048m # 2배 증가 cpu: 1 # 2배 증가 taskManager: resources: memory: 4096m # 2배 증가 cpu: 2 # 2배 증가 flinkConfiguration: taskmanager.numberOfTaskSlots: "4" # 2배 증가 sqlGateway: resources: requests: memory: 1Gi # 2배 증가 cpu: 0.5 # 2배 증가 limits: memory: 2Gi # 2배 증가 cpu: 1 # 2배 증가 # SQL Client 비활성화 (외부 Python 클라이언트 사용) sqlClient: enabled: false ``` ### Kafka 설정 (선택적) ```yaml kafka: enabled: true bootstrapServers: "kafka-cluster-kafka-external-bootstrap.kafka.svc.cluster.local:9094" topics: input: name: test.input output: name: test.output user: name: flink-sql-gateway password: "your-kafka-password" # 또는 External Secret 사용 externalSecret: enabled: true secretName: "kafka-user-credentials" usernameKey: "username" passwordKey: "password" ``` ### External Secret 설정 ```yaml kafka: user: externalSecret: enabled: true secretName: "flink-sql-gateway" # Kafka namespace의 기존 secret usernameKey: "username" passwordKey: "password" ``` ## 📝 사용 예시 ### 1. 테이블 생성 ```sql -- 소스 테이블 생성 CREATE TABLE test_input ( id STRING, message STRING, timestamp_field TIMESTAMP(3), WATERMARK FOR timestamp_field AS timestamp_field - INTERVAL '5' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'test.input', 'properties.bootstrap.servers' = 'kafka-cluster-kafka-external-bootstrap.kafka.svc.cluster.local:9094', 'properties.group.id' = 'flink-sql-gateway', 'properties.security.protocol' = 'SASL_SSL', 'properties.sasl.mechanism' = 'SCRAM-SHA-512', 'properties.sasl.jaas.config' = 'org.apache.kafka.common.security.scram.ScramLoginModule required username="flink-sql-gateway" password="YOUR_PASSWORD";', 'properties.ssl.truststore.location' = '/opt/flink/certs/truststore.jks', 'properties.ssl.truststore.password' = 'changeit', 'properties.ssl.truststore.type' = 'JKS', 'format' = 'json', 'json.timestamp-format.standard' = 'ISO-8601' ); ``` ### 2. SQL 쿼리 실행 ```sql -- 간단한 필터링 쿼리 INSERT INTO test_output SELECT id, CONCAT('Processed: ', message) as processed_message, timestamp_field as original_timestamp, CURRENT_TIMESTAMP as processing_time FROM test_input WHERE LENGTH(message) > 5; ``` ## 🔧 관리 ### 업그레이드 ```bash helm upgrade flink-sql-gateway . \ --namespace flink-sql-gateway \ --values values.yaml ``` ### 제거 ```bash helm uninstall flink-sql-gateway --namespace flink-sql-gateway ``` ### 상태 확인 ```bash # Helm 릴리스 상태 helm status flink-sql-gateway -n flink-sql-gateway # Pod 상태 kubectl get pods -n flink-sql-gateway # FlinkDeployment 상태 kubectl get flinkdeployment -n flink-sql-gateway # 로그 확인 kubectl logs -n flink-sql-gateway -l app=flink-sql-gateway ``` ## 🐛 문제 해결 ### 일반적인 문제 1. **SQL Gateway 연결 실패** ```bash # Gateway 상태 확인 kubectl get deployment flink-sql-gateway -n flink-sql-gateway kubectl logs -n flink-sql-gateway -l app=flink-sql-gateway # Port Forward 확인 kubectl port-forward -n flink-sql-gateway svc/flink-sql-gateway 8083:8083 curl http://localhost:8083/v1/info ``` 2. **Python 클라이언트 연결 실패** ```python # 연결 테스트 import requests try: response = requests.get("http://localhost:8083/v1/info", timeout=5) print("SQL Gateway 연결 성공:", response.json()) except requests.exceptions.ConnectionError: print("Port Forward가 실행 중인지 확인하세요") except requests.exceptions.Timeout: print("SQL Gateway가 시작되지 않았을 수 있습니다") ``` 3. **Flink Session Cluster 시작 실패** ```bash # FlinkDeployment 상태 확인 kubectl describe flinkdeployment flink-session-cluster -n flink-sql-gateway ``` 4. **Kafka 연결 오류** - Kafka 사용자 비밀번호 확인 - 인증서 설정 확인 - 네트워크 연결 확인 ### 디버깅 명령어 ```bash # 전체 환경 상태 확인 kubectl get all -n flink-sql-gateway # 이벤트 확인 kubectl get events -n flink-sql-gateway --sort-by='.lastTimestamp' # Helm 템플릿 확인 helm template flink-sql-gateway . --values values.yaml ``` ## 📚 참고 자료 - [Apache Flink Documentation](https://nightlies.apache.org/flink/flink-docs-release-1.20/) - [Flink SQL Gateway Documentation](https://nightlies.apache.org/flink/flink-docs-release-1.20/docs/dev/table/sql-gateway/) - [Flink Kubernetes Operator](https://nightlies.apache.org/flink/flink-kubernetes-operator-docs-release-1.13/) - [Helm Documentation](https://helm.sh/docs/)