Flink SQL 튜토리얼

Flink SQL은 표준 SQL을 사용해 스트리밍 애플리케이션을 쉽게 개발할 수 있게 해줘요. 데이터베이스나 SQL 유사 시스템을 다뤄본 적이 있다면, ANSI-SQL 2011을 준수하는 Flink는 배우기 쉬워요.

출처: 문서

본문

배우게 될 내용 (What You'll Learn)

이 튜토리얼에서 다음을 배울 수 있어요:

  • Flink SQL Client 시작하기
  • 대화형 SQL 쿼리 실행하기
  • 외부 데이터에서 소스 테이블(source table) 만들기
  • 스트리밍 데이터를 처리하는 연속 쿼리(continuous query) 작성하기
  • 결과를 싱크 테이블(sink table)로 출력하기

코딩 불필요! SQL Client는 SQL 쿼리를 입력하고 즉시 결과를 볼 수 있는 대화형 환경을 제공해요.

사전 요구 사항 (Prerequisites)

  • SQL에 대한 기본 지식
  • 실행 중인 Flink 클러스터 (First Steps 참고)

SQL Client 시작하기 (Starting the SQL Client)

SQL Client는 Flink에 SQL 쿼리를 제출하고 결과를 시각화하는 대화형 클라이언트예요.

First Steps에서 Docker를 사용했다면:

$ docker compose run sql-client

First Steps에서 **로컬 설치(local installation)**를 사용했다면:

$ ./bin/sql-client.sh

SQL Client를 종료하려면 exit;을 입력하고 Enter를 눌러요.

Hello World

쿼리 편집기인 SQL 클라이언트가 실행되면 쿼리 작성을 시작할 차례예요. 다음과 같은 간단한 쿼리로 'Hello World'를 출력하는 것부터 시작해 보죠:

SELECT 'Hello World';

HELP 명령을 실행하면 지원되는 전체 SQL 문 집합이 나열돼요. 그런 명령 중 하나인 SHOW를 실행해 Flink의 내장 함수 전체 목록을 확인해 보죠.

SHOW FUNCTIONS;

이 함수들은 SQL 쿼리를 개발할 때 사용자에게 강력한 도구 상자를 제공해요. 예를 들어 CURRENT_TIMESTAMP는 실행되는 머신의 현재 시스템 시간을 출력해요.

SELECT CURRENT_TIMESTAMP;

소스 테이블 (Source Tables)

모든 SQL 엔진과 마찬가지로 Flink 쿼리는 테이블 위에서 동작해요. Flink는 로컬에서 정적 데이터(데이터 at rest)를 관리하지 않고, 대신 쿼리가 외부 테이블 위에서 지속적으로 동작한다는 점에서 전통적인 데이터베이스와 달라요.

Flink 데이터 처리 파이프라인은 소스 테이블에서 시작돼요. 소스 테이블은 쿼리 실행 중에 처리되는 행을 생성하며, 쿼리의 FROM 절에서 참조되는 테이블이에요. 이것은 Kafka 토픽, 데이터베이스, 파일시스템 또는 Flink가 소비하는 방법을 아는 다른 시스템이 될 수 있어요.

테이블은 SQL 클라이언트 또는 환경 구성 파일을 통해 정의할 수 있어요. SQL 클라이언트는 전통적인 SQL과 유사한 SQL DDL 명령을 지원해요. 표준 SQL DDL은 테이블을 생성, 변경, 삭제하는 데 사용돼요.

Flink는 테이블과 함께 사용할 수 있는 다양한 커넥터포맷을 지원해요. 다음은 샘플 데이터를 자동으로 생성하는 DataGen 커넥터를 사용해 소스 테이블을 정의하는 예시예요.

CREATE TABLE employee_information (
    emp_id INT,
    name STRING,
    dept_id INT
) WITH (
    'connector' = 'datagen',
    'rows-per-second' = '10',
    'fields.emp_id.kind' = 'sequence',
    'fields.emp_id.start' = '1',
    'fields.emp_id.end' = '1000000',
    'fields.name.length' = '8',
    'fields.dept_id.min' = '1',
    'fields.dept_id.max' = '5'
);

이것은 초당 10행을 생성하는 테이블을 만들며, 순차적인 직원 ID, 무작위 8자 이름, 1에서 5 사이의 부서 ID를 가져요.

이 테이블에서 새 행이 생성될 때마다 읽고 즉시 결과를 출력하는 연속 쿼리를 정의할 수 있어요. 예를 들어 부서 1에서 일하는 직원만 필터링할 수 있어요.

SELECT * from employee_information WHERE dept_id = 1;

연속 쿼리 (Continuous Queries)

SQL은 처음에 스트리밍 의미론을 염두에 두고 설계되지는 않았지만, 연속 데이터 파이프라인을 구축하는 강력한 도구예요. Flink SQL이 전통적인 데이터베이스 쿼리와 다른 점은, 행이 도착할 때 지속적으로 소비하고 결과에 대한 업데이트를 생성한다는 것이에요.

연속 쿼리는 결코 종료되지 않으며 결과로 동적 테이블(dynamic table)을 생성해요. 동적 테이블은 스트리밍 데이터에 대한 Flink의 Table API와 SQL 지원의 핵심 개념이에요.

연속 스트림에 대한 집계는 쿼리 실행 중에 집계된 결과를 지속적으로 저장해야 해요. 예를 들어, 들어오는 데이터 스트림에서 각 부서의 직원 수를 세어야 한다고 가정해 보죠. 새 행이 처리될 때 시기적절한 결과를 출력하려면 쿼리가 각 부서에 대한 가장 최신의 카운트를 유지해야 해요.

SELECT 
   dept_id,
   COUNT(*) as emp_count 
FROM employee_information 
GROUP BY dept_id;

이런 쿼리는 상태를 가진(stateful) 것으로 간주돼요. Flink의 고급 결함 허용(fault-tolerance) 메커니즘은 내부 상태와 일관성을 유지하므로, 하드웨어 장애가 발생해도 쿼리는 항상 올바른 결과를 반환해요.

싱크 테이블 (Sink Tables)

이 쿼리를 실행하면 SQL 클라이언트가 실시간으로 하지만 읽기 전용(read-only) 방식으로 출력을 제공해요. 결과를 저장하려면(예: 보고서나 대시보드를 구동하기 위해) 다른 테이블로 작성해야 해요. 이것은 INSERT INTO 문을 사용해 달성할 수 있어요. 이 절에서 참조되는 테이블을 싱크 테이블이라고 해요. INSERT INTO 문은 Flink 클러스터에 분리된(detached) 쿼리로 제출돼요.

먼저 싱크 테이블을 만들어요. 이 예시는 결과를 TaskManager의 로그에 작성하는 print 커넥터를 사용해요:

CREATE TABLE department_counts (
    dept_id INT,
    emp_count BIGINT
) WITH (
    'connector' = 'print'
);

이제 집계된 결과를 이 테이블에 삽입해요:

INSERT INTO department_counts
SELECT
    dept_id,
    COUNT(*) as emp_count
FROM employee_information
GROUP BY dept_id;

제출되면 이것은 백그라운드 작업으로 실행되어 결과를 싱크 테이블에 지속적으로 작성해요. TaskManager 로그에서 출력을 볼 수 있어요:

$ docker compose logs -f taskmanager

운영(production)에서는 일반적으로 JDBC, Kafka 또는 Filesystem 같은 커넥터를 사용해 외부 시스템에 작성해요.

다음 단계 (Next Steps)

이제 Flink SQL을 경험했으니, 학습을 계속하기 위한 몇 가지 경로가 있어요:

Flink SQL 더 깊이 파보기 (Dive Deeper into Flink SQL)

스트리밍 개념 이해하기 (Understand Streaming Concepts)

다른 튜토리얼 시도하기 (Try Other Tutorials)

도움 받기 (Getting Help)

막히면 커뮤니티 지원 리소스를 확인하세요. Apache Flink의 사용자 메일링 리스트는 모든 Apache 프로젝트 중 가장 활발한 리스트 중 하나이며, 빠르게 도움을 받는 좋은 방법이에요.

더 알아보기 (Learn more)