FileSystem SQL 커넥터

FileSystem SQL 커넥터 (FileSystem SQL Connector)

이 커넥터는 Flink FileSystem 추상화가 지원하는 파일 시스템의 파티션 파일에 대한 접근을 제공합니다.

파일 시스템 커넥터 자체는 Flink에 포함되어 있어 추가 의존성이 필요하지 않습니다. 해당 jar는 Flink 배포의 /lib 디렉터리 안에서 찾을 수 있습니다. 파일 시스템에서 행을 읽고 쓰기 위해 해당 포맷을 지정해야 합니다.

파일 시스템 커넥터는 로컬 또는 분산 파일 시스템에서 읽고 쓸 수 있게 해줍니다. 파일 시스템 테이블은 다음과 같이 정의할 수 있습니다:

CREATE TABLE MyUserTable (
  column_name1 INT,
  column_name2 STRING,
  ...
  part_name1 INT,
  part_name2 STRING
) PARTITIONED BY (part_name1, part_name2) WITH (
  'connector' = 'filesystem',           -- required: specify the connector
  'path' = 'file:///path/to/whatever',  -- required: path to a directory
  'format' = '...',                     -- required: file system connector requires to specify a format,
                                        -- Please refer to Table Formats
                                        -- section for more details
  'partition.default-name' = '...',     -- optional: default partition name in case the dynamic partition
                                        -- column value is null/empty string
  'source.path.regex-pattern' = '...',  -- optional: regex pattern to filter files to read under the
                                        -- directory of `path` option. This regex pattern should be
                                        -- matched with the absolute file path. If this option is set,
                                        -- the connector  will recursive all files under the directory
                                        -- of `path` option

  -- optional: the option to enable shuffle data by dynamic partition fields in sink phase, this can greatly
  -- reduce the number of file for filesystem sink but may lead data skew, the default value is false.
  'sink.shuffle-by-partition.enable' = '...',
  ...
)

Flink File System 특화 의존성을 포함해야 합니다.

파일 시스템 커넥터의 동작은 이전 레거시 파일 시스템 커넥터(previous legacy filesystem connector)와 매우 다릅니다: path 파라미터는 파일이 아니라 디렉터리를 지정하며, 선언한 경로에서 사람이 읽을 수 있는 파일을 얻을 수 없습니다.

출처: 문서

본문

파티션 파일 (Partition Files)

Flink의 파일 시스템 파티션 지원은 표준 hive 포맷을 사용합니다. 그러나 파티션이 테이블 카탈로그에 미리 등록될 필요는 없습니다. 파티션은 디렉터리 구조에 기반해 발견되고 추론됩니다. 예를 들어 아래 디렉터리에 기반한 파티션 테이블은 datetimehour 파티션을 포함하는 것으로 추론됩니다.

path
└── datetime=2019-08-25
    └── hour=11
        ├── part-0.parquet
        ├── part-1.parquet
    └── hour=12
        ├── part-0.parquet
└── datetime=2019-08-26
    └── hour=6
        ├── part-0.parquet

파일 시스템 테이블은 파티션 삽입과 덮어쓰기 삽입(overwrite inserting)을 모두 지원합니다. INSERT 문을 참고하세요. 파티션 테이블에 insert overwrite하면 전체 테이블이 아니라 해당 파티션만 덮어씁니다.

파일 포맷 (File Formats)

파일 시스템 커넥터는 여러 포맷을 지원합니다:

  • CSV: RFC-4180.
  • JSON: 파일 시스템 커넥터의 JSON 포맷은 일반적인 JSON 파일이 아니라 줄바꿈으로 구분된 JSON(newline delimited JSON)임에 유의하세요.
  • Avro: Apache Avro. avro.codec 구성으로 압축을 지원합니다.
  • Parquet: Apache Parquet. Hive와 호환.
  • Orc: Apache Orc. Hive와 호환.
  • Debezium-JSON: debezium-json.
  • Canal-JSON: canal-json.
  • Raw: raw.

소스 (Source)

파일 시스템 커넥터는 단일 파일 또는 전체 디렉터리를 단일 테이블로 읽는 데 사용될 수 있습니다.

디렉터리를 소스 경로로 사용할 때 디렉터리 안 파일의 수집 순서는 정의되어 있지 않습니다.

디렉터리 감시 (Directory watching)

기본적으로 파일 시스템 커넥터는 유한(bounded)이며, 구성된 경로를 한 번 스캔한 후 스스로 닫힙니다.

source.monitor-interval 옵션을 구성해 연속 디렉터리 감시를 활성화할 수 있습니다:

기본값 타입 설명
source.monitor-interval (none) Duration 소스가 새 파일을 확인하는 간격. 간격은 0보다 커야 합니다. 각 파일은 경로로 고유하게 식별되며, 발견되는 즉시 한 번 처리됩니다. 이미 처리된 파일 집합은 소스의 수명 전체 동안 상태에 유지되므로, 체크포인트와 세이브포인트에 소스 상태와 함께 영속됩니다. 더 짧은 간격은 파일을 더 빨리 발견함을 의미하지만, 파일 시스템/객체 저장소의 더 빈번한 목록화나 디렉터리 순회를 의미하기도 합니다. 이 구성 옵션이 설정되지 않으면 제공된 경로가 한 번 스캔되므로 소스는 유한합니다.

사용 가능한 메타데이터 (Available Metadata)

다음 커넥터 메타데이터는 테이블 정의에서 메타데이터 컬럼으로 접근할 수 있습니다. 모든 메타데이터는 읽기 전용입니다.

데이터 타입 설명
file.path STRING NOT NULL 입력 파일의 전체 경로.
file.name STRING NOT NULL 파일 이름, 즉 파일 경로의 루트에서 가장 먼 요소.
file.size BIGINT NOT NULL 파일의 바이트 수.
file.modification-time TIMESTAMP_LTZ(3) NOT NULL 파일의 수정 시간.

확장된 CREATE TABLE 예시는 이러한 메타데이터 필드를 노출하는 문법을 보여줍니다:

CREATE TABLE MyUserTableWithFilepath (
  column_name1 INT,
  column_name2 STRING,
  `file.path` STRING NOT NULL METADATA
) WITH (
  'connector' = 'filesystem',
  'path' = 'file:///path/to/whatever',
  'format' = 'json'
)

스트리밍 싱크 (Streaming Sink)

파일 시스템 커넥터는 Flink의 FileSystem에 기반한 스트리밍 쓰기를 지원해 레코드를 파일에 씁니다. 행 인코딩(Row-encoded) 포맷은 CSV와 JSON입니다. 벌크 인코딩(Bulk-encoded) 포맷은 Parquet, ORC, Avro입니다.

SQL을 직접 작성해 파티션되지 않은 테이블에 스트림 데이터를 삽입할 수 있습니다. 파티션 테이블이면 파티션 관련 연산을 구성할 수 있습니다. 자세한 내용은 Partition Commit을 참고하세요.

롤링 정책 (Rolling Policy)

파티션 디렉터리 안의 데이터는 part 파일로 나뉩니다. 각 파티션은 해당 파티션에 데이터를 받은 싱크의 각 서브태스크에 대해 최소 하나의 part 파일을 포함합니다. 진행 중(in-progress) part 파일이 닫히고 구성 가능한 롤링 정책에 따라 추가 part 파일이 생성됩니다. 이 정책은 크기와 파일이 열려 있을 수 있는 최대 기간을 지정하는 시간 초과에 따라 part 파일을 롤링합니다.

옵션 필수 Forwarded 기본값 타입 설명
sink.rolling-policy.file-size optional yes 128MB MemorySize 롤링 전 최대 part 파일 크기.
sink.rolling-policy.rollover-interval optional yes 30 min Duration 롤링 전 part 파일이 열려 있을 수 있는 최대 시간(기본 30분으로 너무 많은 작은 파일을 방지). 이것이 확인되는 빈도는 'sink.rolling-policy.check-interval' 옵션으로 제어됩니다.
sink.rolling-policy.check-interval optional yes 1 min Duration 시간 기반 롤링 정책을 확인하는 간격. 'sink.rolling-policy.rollover-interval'에 따라 part 파일이 롤오버되어야 하는지 확인하는 빈도를 제어합니다.

참고: 벌크 포맷(parquet, orc, avro)의 경우 롤링 정책과 체크포인트 간격(pending 파일은 다음 체크포인트에서 완료됨)이 이러한 part 파일들의 크기와 수를 제어합니다.

참고: 행 포맷(csv, json)의 경우, 파일 시스템에 데이터가 존재함을 관찰하기 전에 오래 기다리고 싶지 않다면 커넥터 속성의 sink.rolling-policy.file-size 또는 sink.rolling-policy.rollover-interval 파라미터와 Flink 구성 파일의 execution.checkpointing.interval 파라미터를 함께 설정할 수 있습니다. 다른 포맷(avro, orc)의 경우 Flink 구성 파일의 execution.checkpointing.interval 파라미터를 설정하기만 하면 됩니다.

파일 압축 (File Compaction)

파일 싱크는 파일 압축을 지원하여, 애플리케이션은 많은 수의 파일을 생성하지 않고도 더 작은 체크포인트 간격을 가질 수 있습니다.

옵션 필수 Forwarded 기본값 타입 설명
auto-compaction optional no false Boolean 스트리밍 싱크에서 자동 압축을 활성화할지 여부. 데이터는 임시 파일에 기록됩니다. 체크포인트가 완료된 후, 한 체크포인트로 생성된 임시 파일이 압축됩니다. 임시 파일은 압축 전에는 보이지 않습니다.
compaction.file-size optional yes (none) MemorySize 압축 대상 파일 크기. 기본값은 롤링 파일 크기.

활성화되면 파일 압축은 대상 파일 크기에 따라 여러 작은 파일을 더 큰 파일로 병합합니다. 프로덕션에서 파일 압축을 실행할 때는 다음을 알고 있어야 합니다:

  • 단일 체크포인트의 파일만 압축됩니다. 즉, 체크포인트 수만큼의 파일 수는 최소한 생성됩니다.
  • 병합 전의 파일은 보이지 않으므로 파일의 가시성은 체크포인트 간격 + 압축 시간일 수 있습니다.
  • 압축이 너무 오래 걸리면 작업에 백프레셔를 가하고 체크포인트 기간을 늘립니다.

파티션 커밋 (Partition Commit)

파티션을 쓴 뒤 다운스트림 애플리케이션에 알리는 것이 종종 필요합니다. 예를 들어 Hive metastore에 파티션을 추가하거나 디렉터리에 _SUCCESS 파일을 쓰는 것입니다. 파일 시스템 싱크는 커스텀 정책을 구성할 수 있는 파티션 커밋 기능을 포함합니다. 커밋 동작은 트리거(triggers)정책(policies)의 조합에 기반합니다.

  • 트리거 (Trigger): 파티션 커밋의 시점은 파티션에서 추출한 시간 또는 처리 시간으로 결정할 수 있습니다.
  • 정책 (Policy): 파티션을 커밋하는 방법. 내장 정책은 success 파일과 metastore의 커밋을 지원하며, hive의 분석을 트리거해 통계를 생성하거나 작은 파일을 병합하는 등 자신만의 정책을 구현할 수도 있습니다.

참고: 파티션 커밋은 동적 파티션 삽입에서만 동작합니다.

파티션 커밋 트리거 (Partition commit trigger)

언제 파티션을 커밋할지 정의하려면 파티션 커밋 트리거를 제공하세요:

옵션 필수 Forwarded 기본값 타입 설명
sink.partition-commit.trigger optional yes process-time String 파티션 커밋의 트리거 타입: 'process-time': 머신의 시간에 기반하며 파티션 시간 추출이나 워터마크 생성이 필요하지 않습니다. '현재 시스템 시간'이 '파티션 생성 시스템 시간'에 'delay'를 더한 것을 지나면 파티션을 커밋합니다. 'partition-time': 파티션 값에서 추출한 시간에 기반하며 워터마크 생성이 필요합니다. '워터마크'가 '파티션 값에서 추출한 시간'에 'delay'를 더한 것을 지나면 파티션을 커밋합니다.
sink.partition-commit.delay optional yes 0 s Duration 파티션은 지연 시간이 지나기 전에는 커밋되지 않습니다. 일일 파티션이면 '1 d'여야 하고, 시간별 파티션이면 '1 h'여야 합니다.
sink.partition-commit.watermark-time-zone optional yes UTC String 긴 워터마크 값을 TIMESTAMP 값으로 파싱하기 위한 시간대. 파싱된 워터마크 타임스탬프는 파티션 시간과 비교되어 파티션을 커밋할지 결정하는 데 사용됩니다. 이 옵션은 sink.partition-commit.trigger가 'partition-time'으로 설정된 경우에만 적용됩니다. 이 옵션이 올바르게 구성되지 않으면(예: 소스 rowtime이 TIMESTAMP_LTZ 컬럼에 정의되었는데 이 구성이 없으면) 사용자는 몇 시간 후에 파티션이 커밋되는 것을 볼 수 있습니다. 기본값은 'UTC'이며, 워터마크가 TIMESTAMP 컬럼에 정의되었거나 정의되지 않았음을 의미합니다. 워터마크가 TIMESTAMP_LTZ 컬럼에 정의되면 워터마크의 시간대는 세션 시간대입니다. 옵션 값은 'America/Los_Angeles' 같은 전체 이름 또는 'GMT-08:00' 같은 커스텀 시간대 id입니다.

트리거에는 두 가지 유형이 있습니다:

  • 첫 번째는 파티션 처리 시간입니다. 파티션 시간 추출이나 워터마크 생성이 필요하지 않습니다. 파티션 생성 시간과 현재 시스템 시간에 따라 파티션 커밋을 트리거합니다. 이 트리거는 더 범용적이지만 정확하지는 않습니다. 예를 들어 데이터 지연이나 페일오버는 조기 파티션 커밋으로 이어질 수 있습니다.
  • 두 번째는 파티션 값에서 추출한 시간과 워터마크에 따른 파티션 커밋 트리거입니다. 이는 작업에 워터마크 생성이 필요하고 파티션이 시간별 파티션이나 일일 파티션처럼 시간에 따라 나뉘어야 합니다.

데이터가 완전한지 여부와 무관하게 다운스트림이 파티션을 가능한 한 빨리 보도록 하려면:

  • 'sink.partition-commit.trigger'='process-time' (기본값)
  • 'sink.partition-commit.delay'='0s' (기본값) 파티션에 데이터가 생기는 즉시 커밋합니다. 참고: 파티션은 여러 번 커밋될 수 있습니다.

데이터가 완전할 때만 다운스트림이 파티션을 보도록 하고, 작업에 워터마크 생성이 있으며 파티션 값에서 시간을 추출할 수 있다면:

  • 'sink.partition-commit.trigger'='partition-time'
  • 'sink.partition-commit.delay'='1h' (파티션이 시간별이면 '1h', 파티션 타입에 따라 다름) 이것이 파티션 커밋의 가장 정확한 방법이며, 커밋된 파티션이 가능한 한 데이터 완전하도록 보장하려 합니다.

데이터가 완전할 때만 다운스트림이 파티션을 보도록 하지만 워터마크가 없거나 파티션 값에서 시간을 추출할 수 없다면:

  • 'sink.partition-commit.trigger'='process-time' (기본값)
  • 'sink.partition-commit.delay'='1h' (파티션이 시간별이면 '1h', 파티션 타입에 따라 다름) 파티션 커밋을 정확히 시도하지만 데이터 지연이나 페일오버는 조기 파티션 커밋으로 이어집니다.

늦은 데이터 처리: 레코드가 이미 커밋된 파티션에 쓰여져야 한다고 판단되면 레코드가 해당 파티션에 기록되고, 이후 이 파티션의 커밋이 다시 트리거됩니다.

파티션 시간 추출기 (Partition Time Extractor)

시간 추출기는 파티션 값에서 시간을 추출하는 것을 정의합니다.

옵션 필수 Forwarded 기본값 타입 설명
partition.time-extractor.kind optional no default String 파티션 값에서 시간을 추출하는 시간 추출기. default와 custom을 지원합니다. default는 타임스탬프 패턴/포맷터를 구성할 수 있습니다. custom은 추출기 클래스를 구성해야 합니다.
partition.time-extractor.class optional no (none) String PartitionTimeExtractor 인터페이스를 구현하는 추출기 클래스.
partition.time-extractor.timestamp-pattern optional no (none) String 'default' 구성 방식은 사용자가 파티션 필드를 사용해 합법적인 타임스탬프 패턴을 얻을 수 있게 해줍니다. 기본값은 첫 번째 필드의 'yyyy-MM-dd hh:mm:ss'를 지원합니다. 단일 파티션 필드 'dt'에서 타임스탬프를 추출해야 하면 '$dt'를 구성할 수 있습니다. 'year', 'month', 'day', 'hour' 같은 여러 파티션 필드에서 추출해야 하면 '$year-$month-$day $hour:00:00'을 구성할 수 있습니다. 'dt'와 'hour' 두 파티션 필드에서 추출해야 하면 '$dt $hour:00:00'을 구성할 수 있습니다.
partition.time-extractor.timestamp-formatter optional no yyyy-MM-dd HH:mm:ss String 파티션 타임스탬프 문자열 값을 타임스탬프로 포맷하는 포맷터. 파티션 타임스탬프 문자열 값은 'partition.time-extractor.timestamp-pattern'으로 표현됩니다. 예를 들어 'year', 'month', 'day'라고 하는 여러 파티션 필드에서 파티션 타임스탬프가 추출되면 'partition.time-extractor.timestamp-pattern'을 '$year$month$day'로, partition.time-extractor.timestamp-formatter를 'yyyyMMdd'로 구성할 수 있습니다. 기본 포맷터는 'yyyy-MM-dd HH:mm:ss'입니다.

timestamp-formatter는 Java의 DateTimeFormatter와 호환됩니다.

기본 추출기는 파티션 필드로 구성된 타임스탬프 패턴에 기반합니다. PartitionTimeExtractor 인터페이스에 기반해 완전히 커스텀한 파티션 추출 구현을 지정할 수도 있습니다.

public class HourPartTimeExtractor implements PartitionTimeExtractor {
    @Override
    public LocalDateTime extract(List<String> keys, List<String> values) {
        String dt = values.get(0);
        String hour = values.get(1);
		return Timestamp.valueOf(dt + " " + hour + ":00:00").toLocalDateTime();
	}
}
파티션 커밋 정책 (Partition Commit Policy)

파티션 커밋 정책은 파티션이 커밋될 때 어떤 동작이 취해지는지 정의합니다.

  • 첫 번째는 metastore입니다. metastore 정책은 hive 테이블만 지원합니다. 파일 시스템은 디렉터리 구조를 통해 파티션을 관리합니다.
  • 두 번째는 success 파일입니다. 파티션에 해당하는 디렉터리에 빈 파일을 씁니다.
옵션 필수 Forwarded 기본값 타입 설명
sink.partition-commit.policy.kind optional yes (none) String 파티션을 커밋하는 정책은 다운스트림 애플리케이션에 파티션 쓰기가 끝났음을 알리는 것입니다. 파티션은 읽을 준비가 되었습니다. metastore: metastore에 파티션 추가. metastore 정책은 hive 테이블만 지원하며 파일 시스템은 디렉터리 구조를 통해 파티션을 관리합니다. success-file: 디렉터리에 '_success' 파일 추가. 둘 다 동시에 구성할 수 있습니다: 'metastore,success-file'. custom: 정책 클래스를 사용해 커밋 정책 생성. 여러 정책 구성 지원: 'metastore,success-file'.
sink.partition-commit.policy.class optional yes (none) String PartitionCommitPolicy 인터페이스를 구현하는 파티션 커밋 정책 클래스. 커스텀 커밋 정책에서만 동작합니다.
sink.partition-commit.policy.class.parameters optional yes (none) String 커스텀 커밋 정책의 생성자에 전달되는 파라미터. 여러 파라미터는 세미콜론으로 구분됩니다(예: 'param1;param2'). 구성 값은 목록(['param1', 'param2'])으로 분할되어 커스텀 커밋 정책 클래스의 생성자에 전달됩니다. 이 옵션은 선택이며, 구성되지 않으면 기본 생성자가 사용됩니다.
sink.partition-commit.success-file.name optional yes _SUCCESS String success-file 파티션 커밋 정책의 파일 이름. 기본값은 '_SUCCESS'.

커밋 정책 구현을 확장할 수 있습니다. 커스텀 커밋 정책 구현은 다음과 같습니다:

public class AnalysisCommitPolicy implements PartitionCommitPolicy {
    private HiveShell hiveShell;

    @Override
	public void commit(Context context) throws Exception {
	    if (hiveShell == null) {
	        hiveShell = createHiveShell(context.catalogName());
	    }

        hiveShell.execute(String.format(
            "ALTER TABLE %s ADD IF NOT EXISTS PARTITION (%s = '%s') location '%s'",
	        context.tableName(),
	        context.partitionKeys().get(0),
	        context.partitionValues().get(0),
	        context.partitionPath()));
	    hiveShell.execute(String.format(
	        "ANALYZE TABLE %s PARTITION (%s = '%s') COMPUTE STATISTICS FOR COLUMNS",
	        context.tableName(),
	        context.partitionKeys().get(0),
	        context.partitionValues().get(0)));
	}
}

싱크 병렬도 (Sink Parallelism)

외부 파일 시스템( Hive 포함)에 파일을 쓰는 병렬도는 해당 테이블 옵션으로 구성할 수 있으며, 스트리밍 모드와 배치 모드 모두에서 지원됩니다. 기본적으로 병렬도는 마지막 업스트림 체인 연산자와 동일하게 구성됩니다. 업스트림 병렬도와 다른 병렬도가 구성되면 파일 쓰기 연산자와 파일 압축 연산자(사용된 경우)가 그 병렬도를 적용합니다.

옵션 필수 Forwarded 기본값 타입 설명
sink.parallelism optional no (none) Integer 외부 파일 시스템에 파일을 쓰는 병렬도. 값은 0보다 커야 하며 그렇지 않으면 예외가 발생합니다.

참고: 현재 싱크 병렬도 구성은 업스트림의 changelog 모드가 INSERT-ONLY인 경우에만 지원됩니다. 그렇지 않으면 예외가 발생합니다.

전체 예시 (Full Example)

아래 예시는 파일 시스템 커넥터를 사용해 스트리밍 쿼리가 Kafka에서 파일 시스템으로 데이터를 쓰고, 배치 쿼리가 그 데이터를 다시 읽는 방법을 보여줍니다.

CREATE TABLE kafka_table (
  user_id STRING,
  order_amount DOUBLE,
  log_ts TIMESTAMP(3),
  WATERMARK FOR log_ts AS log_ts - INTERVAL '5' SECOND
) WITH (...);

CREATE TABLE fs_table (
  user_id STRING,
  order_amount DOUBLE,
  dt STRING,
  `hour` STRING
) PARTITIONED BY (dt, `hour`) WITH (
  'connector'='filesystem',
  'path'='...',
  'format'='parquet',
  'sink.partition-commit.delay'='1 h',
  'sink.partition-commit.policy.kind'='success-file'
);

-- streaming sql, insert into file system table
INSERT INTO fs_table
SELECT
    user_id,
    order_amount,
    DATE_FORMAT(log_ts, 'yyyy-MM-dd'),
    DATE_FORMAT(log_ts, 'HH')
FROM kafka_table;

-- batch sql, select with partition pruning
SELECT * FROM fs_table WHERE dt='2020-05-20' and `hour`='12';

워터마크가 TIMESTAMP_LTZ 컬럼에 정의되고 partition-time으로 커밋하는 경우, 세션 시간대로 sink.partition-commit.watermark-time-zone을 설정해야 합니다. 그렇지 않으면 파티션 커밋이 몇 시간 후에 발생할 수 있습니다.

CREATE TABLE kafka_table (
  user_id STRING,
  order_amount DOUBLE,
  ts BIGINT, -- time in epoch milliseconds
  ts_ltz AS TO_TIMESTAMP_LTZ(ts, 3),
  WATERMARK FOR ts_ltz AS ts_ltz - INTERVAL '5' SECOND -- Define watermark on TIMESTAMP_LTZ column
) WITH (...);

CREATE TABLE fs_table (
  user_id STRING,
  order_amount DOUBLE,
  dt STRING,
  `hour` STRING
) PARTITIONED BY (dt, `hour`) WITH (
  'connector'='filesystem',
  'path'='...',
  'format'='parquet',
  'partition.time-extractor.timestamp-pattern'='$dt $hour:00:00',
  'sink.partition-commit.delay'='1 h',
  'sink.partition-commit.trigger'='partition-time',
  'sink.partition-commit.watermark-time-zone'='Asia/Shanghai', -- Assume user configured time zone is 'Asia/Shanghai'
  'sink.partition-commit.policy.kind'='success-file'
);

-- streaming sql, insert into file system table
INSERT INTO fs_table
SELECT
    user_id,
    order_amount,
    DATE_FORMAT(ts_ltz, 'yyyy-MM-dd'),
    DATE_FORMAT(ts_ltz, 'HH')
FROM kafka_table;

-- batch sql, select with partition pruning
SELECT * FROM fs_table WHERE dt='2020-05-20' and `hour`='12';

더 알아보기 (Learn more)