SQL Gateway 개요

SQL Gateway 개요

SQL Gateway는 원격의 여러 클라이언트가 동시에(concurrency) SQL을 실행할 수 있게 해주는 서비스입니다. 이를 통해 Flink 작업을 제출하고, 메타데이터를 조회하고, 데이터를 온라인으로 분석하는 쉬운 방법을 제공합니다. SQL Gateway는 플러그형(pluggable) 엔드포인트들과 SqlGatewayService로 구성되며, SqlGatewayService는 엔드포인트들이 요청을 처리할 때 재사용하는 프로세서입니다. 엔드포인트는 사용자가 연결할 수 있는 진입점(entry point)으로, 엔드포인트의 종류에 따라 사용자는 서로 다른 도구로 연결할 수 있습니다.

출처: 문서

본문

소개 (Introduction)

SQL Gateway는 원격의 여러 클라이언트가 동시에 SQL을 실행할 수 있게 해주는 서비스입니다. Flink 작업을 제출하고, 메타데이터를 조회하고, 데이터를 온라인으로 분석하는 쉬운 방법을 제공합니다. SQL Gateway는 플러그형 엔드포인트들과 SqlGatewayService로 구성되어 있습니다. SqlGatewayService는 엔드포인트들이 요청을 처리할 때 재사용하는 프로세서이며, 엔드포인트는 사용자가 연결할 수 있는 진입점입니다. 엔드포인트 종류에 따라 사용자는 다른 유틸리티로 연결할 수 있습니다.

SQL Gateway 아키텍처

시작하기 (Getting Started)

이 섹션은 명령줄에서 첫 번째 Flink SQL 프로그램을 설정하고 실행하는 방법을 설명합니다.

SQL Gateway는 일반 Flink 배포(distribution)에 포함되어 있어 바로 실행할 수 있습니다(out-of-the-box). 테이블 프로그램을 실행할 수 있는 실행 중인 Flink 클러스터만 있으면 됩니다. Flink 클러스터 설정에 대한 자세한 내용은 클러스터 & 배포 부분을 참고하세요. 단순히 SQL Gateway를 시험해 보고 싶다면 다음 명령으로 워커 하나를 가진 로컬 클러스터를 시작할 수도 있습니다.

$ ./bin/start-cluster.sh

SQL Gateway 시작

SQL Gateway 스크립트도 Flink의 binary 디렉터리에 있습니다. 다음 호출로 시작할 수 있습니다.

$ ./bin/sql-gateway.sh start -Dsql-gateway.endpoint.rest.address=localhost

이 명령은 주소 localhost:8083에서 수신 대기하는 REST Endpoint로 SQL Gateway를 시작합니다. curl 명령으로 REST Endpoint가 사용 가능한지 확인할 수 있습니다.

$ curl http://localhost:8083/v1/info
{"productName":"Apache Flink","version":"2.3.0"}

SQL 쿼리 실행

설정과 클러스터 연결을 검증하려면 다음 단계를 따르세요.

1단계: 세션 열기

$ curl --request POST http://localhost:8083/v1/sessions
{"sessionHandle":"..."}

반환 결과의 sessionHandle은 SQL Gateway가 모든 활성 사용자를 고유하게 식별하는 데 사용합니다.

2단계: 쿼리 실행

$ curl --request POST http://localhost:8083/v1/sessions/${sessionHandle}/statements/ --data '{"statement": "SELECT 1"}'
{"operationHandle":"..."}

반환 결과의 operationHandle은 SQL Gateway가 제출된 SQL을 고유하게 식별하는 데 사용합니다.

Flink SQL Gateway는 클라이언트가 작업을 제출할 Flink 클러스터를 지정할 수 있게 해줍니다. 이를 통해 SQL 문을 원격으로 실행할 수 있고 REST API를 통해 Flink 클러스터와 더 쉽게 상호작용할 수 있습니다. Flink 클러스터 주소를 설정하려면 POST 요청 본문에서 executionConfig 변수 안에 rest.addressrest.port를 넣으세요. 예를 들면 다음과 같습니다.

$ curl --request POST http://localhost:8083/v1/sessions/${sessionHandle}/statements/ --data '{"executionConfig": {"rest.address":"jobmanager-host", "rest.port":8081},"statement": "SELECT 1"}'
{"operationHandle":"..."}

3단계: 결과 가져오기

위의 sessionHandleoperationHandle으로 해당 결과를 가져올 수 있습니다.

$ curl --request GET http://localhost:8083/v1/sessions/${sessionHandle}/operations/${operationHandle}/result/0
{
  "results": {
    "columns": [
      {
        "name": "EXPR$0",
        "logicalType": {
          "type": "INTEGER",
          "nullable": false
        }
      }
    ],
    "data": [
      {
        "kind": "INSERT",
        "fields": [
          1
        ]
      }
    ]
  },
  "resultType": "PAYLOAD",
  "nextResultUri": "..."
}

결과의 nextResultUrinull이 아니면 다음 배치 결과를 가져오는 데 사용합니다.

$ curl --request GET ${nextResultUri}

스크립트 배포 (Deploying a Script)

SQL Gateway는 Application Mode에서 스크립트 배포를 지원합니다. Application Mode에서는 JobManager가 스크립트 컴파일을 담당합니다. 예를 들어 Kafka Source 같은 사용자 정의 리소스를 스크립트에서 사용하려면 ADD JAR 명령으로 필요한 아티팩트를 다운로드하세요.

다음은 cluster id가 CLUSTER_ID인 Flink 네이티브 K8S 클러스터에 스크립트를 배포하는 예시입니다.

$ curl --request POST http://localhost:8083/sessions/${SESSION_HANDLE}/scripts \
--header 'Content-Type: application/json' \
--data-raw '{
    "script": "CREATE TEMPORARY TABLE sink(a INT) WITH ( '\''connector'\'' = '\''blackhole'\''); INSERT INTO sink VALUES (1), (2), (3);",
    "executionConfig": {
        "execution.target": "kubernetes-application",
        "kubernetes.cluster-id": "'${CLUSTER_ID}'",
        "kubernetes.container.image.ref": "'${FLINK_IMAGE_NAME}'",
        "jobmanager.memory.process.size": "1000m",
        "taskmanager.memory.process.size": "1000m",
        "kubernetes.jobmanager.cpu": 0.5,
        "kubernetes.taskmanager.cpu": 0.5,
        "kubernetes.rest-service.exposed.type": "NodePort"
    }
}'

참고: PyFlink로 스크립트를 실행하려면 PyFlink가 설치된 이미지를 사용하세요. 자세한 내용은 docker에서 PyFlink 활성화를 참고하세요.

구성 (Configuration)

SQL Gateway 시작 옵션

현재 SQL Gateway 스크립트에는 다음과 같은 선택적 명령이 있습니다. 자세한 내용은 뒤의 단락에서 설명합니다.

$ ./bin/sql-gateway.sh --help

Usage: sql-gateway.sh [start|start-foreground|stop|stop-all] [args]
  commands:
    start               - Run a SQL Gateway as a daemon
    start-foreground    - Run a SQL Gateway as a console application
    stop                - Stop the SQL Gateway daemon
    stop-all            - Stop all the SQL Gateway daemons
    -h | --help         - Show this help message

"start" 또는 "start-foreground" 명령으로 CLI에서 SQL Gateway를 구성할 수 있습니다.

$ ./bin/sql-gateway.sh start --help

Start the Flink SQL Gateway as a daemon to submit Flink SQL.

  Syntax: start [OPTIONS]
     -D <property=value>   Use value for given property
     -h,--help             Show the help message with descriptions of all
                           options.

SQL Gateway 구성

SQL Gateway를 시작할 때 아래처럼 구성하거나, 유효한 Flink 구성 항목을 넣을 수 있습니다.

$ ./sql-gateway -Dkey=value
키 (Key) 기본값 (Default) 타입 (Type) 설명 (Description)
sql-gateway.session.check-interval 1 min Duration 유휴 세션 타임아웃의 확인 간격. 0으로 설정하면 비활성화할 수 있습니다.
sql-gateway.session.idle-timeout 10 min Duration 세션이 해당 간격 동안 접근되지 않았을 때 세션을 닫는 타임아웃 간격. 0으로 설정하면 세션이 닫히지 않습니다.
sql-gateway.session.max-num 1000000 Integer sql gateway 서비스의 활성 세션 최대 개수.
sql-gateway.session.plan-cache.enabled false Boolean true이면 sql gateway가 세션별 쿼리 플랜을 캐시하고 재사용합니다.
sql-gateway.session.plan-cache.size 100 Integer 플랜 캐시 크기. table.optimizer.plan-cache.enabled가 true일 때만 적용됩니다.
sql-gateway.session.plan-cache.ttl 1 hour Duration 플랜 캐시의 TTL. 쓰기 이후 캐시가 얼마나 오래 유지될지 제어하며, table.optimizer.plan-cache.enabled가 true일 때만 적용됩니다.
sql-gateway.worker.keepalive-time 5 min Duration 유휴 워커 스레드의 keepalive 시간. 워커 수가 최소 워커 수를 초과하면 이 시간 간격 후 초과 스레드가 종료됩니다.
sql-gateway.worker.threads.max 500 Integer sql gateway 서비스의 워커 스레드 최대 개수.
sql-gateway.worker.threads.min 5 Integer sql gateway 서비스의 워커 스레드 최소 개수.

지원되는 엔드포인트 (Supported Endpoints)

Flink는 네이티브로 REST EndpointHiveServer2 Endpoint를 지원합니다. SQL Gateway는 기본적으로 REST Endpoint와 함께 제공됩니다. 유연한 아키텍처 덕분에 사용자는 다음 호출로 지정된 엔드포인트로 SQL Gateway를 시작할 수 있습니다.

$ ./bin/sql-gateway.sh start -Dsql-gateway.endpoint.type=hiveserver2

또는 Flink 구성 파일에 다음 구성을 추가합니다.

sql-gateway.endpoint.type: hiveserver2

참고: Flink 구성 파일에도 sql-gateway.endpoint.type 옵션이 있다면 CLI 명령이 더 높은 우선순위를 가집니다.

특정 엔드포인트에 대해서는 해당 페이지를 참고하세요.

더 알아보기 (Learn more)