|
|
# 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/)
|
|
|
|