Skip to content

FlinkSQL 集成 Hive 与 FileSystem:流批读写、维表 Join 与小文件合并

本文基于 Flink 1.12+。Flink 在 1.12 版本中完善了Hive和FileSystem connector,增加了很多新特性以满足批和流的需求,本文主要介绍1.12版本下hive和FileSystem connector的相关特性和使用,并简要对比FileSystem connector和旧Hdfs connector的相关区别。

读Hive表

建完catalog后,就可以用类似:catalogName.databaseName.tableName来引用表,SQL如下:

sql

create table pt (
    `timestamp` BIGINT,
    `time` STRING,
    id BIGINT,
    product STRING,
    price DOUBLE,
    canSell STRING,
    selledNum BIGINT,
    dt STRING,    -- partition Column
    `hour` STRING,  -- partition Column
    `min` STRING  -- partition Column
) with (
    'connector'='print'
);

insert into pt
select * from 
hive_catalog.database_name.table_name;

对于Hive表的读取,目前Flink支持流模式和批模式,下面分别进行介绍。

a. 流模式

在流模式下,Flink会去监听表中是否有新数据,有可见的新数据就会通知消费,该功能通过参数streaming-source.enable = true 来启用。使用Catalog后不再需要写DDL,那么建表的相关参数就需要动态指定,Flink提供了SQL Hints,示例如下:

ada
SELECT * FROM hive_table 
/*+ OPTIONS('streaming-source.enable'='true', 'streaming-source.consume-start-offset'='2020-05-20') */;

流模式下还有几个与消费起点相关的参数(streaming-source.consume-start-offsetstreaming-source.partition-order 等),用途见官方文档的 Hive Read 一节。

需要注意的是,如果新建的Hive表没有分区的话,需要将streaming-source.partition-order设置为create-time。

b. 批模式

Flink目前也支持通过批的方式去读取Hive表,开启方法如下:

sql
SET execution.runtime-mode = batch;   -- 执行模式,默认 streaming
SET catalog = hive_catalog;

insert into pt
select * from 
hive_catalog.database_name.table_name;

因为streaming-source.enable的值默认是false,所以不需要额外指定。Batch模式会有一些额外的优化,并且不会执行checkpoint,避免不必要的性能损耗。在批模式下,Hive source的并发默认会根据文件和block的数量来动态生成,如果要显式指定,可以用 table.exec.hive.infer-source-parallelism 一族参数覆盖。

2. 写Hive表

有了Catalog后,写hive表同样变得非常简单,只需书写DML即可。示例如下:

sql
SET execution.runtime-mode = batch;   -- 或 streaming,默认 streaming
SET catalog = hive_catalog;

CREATE TABLE source_table (
    name STRING,
    id STRING,
    `timestamp` BIGINT
) with (
    'connector'='datagen',
    'rows-per-second'='1',
    'number-of-rows'='2000'
);

INSERT INTO hive_catalog.demo_db.test_hive_table
SELECT * FROM source_table;

上述例子中的hive表是一个简单的非分区表,对于分区表而言,在写入时需要额外写入分区字段,使得该条数据能够正确找到分区,或者在DML中直接指定固定的分区。

sql
从数据中提取分区信息:

CREATE TABLE Kafka_source(
    name STRING,
    id STRING,
    `timestamp` BIGINT,     -- 13位unix时间戳
) with (
    ......
);
SET catalog = hive_catalog;

INSERT INTO hive_catalog.demo_db.test_hive_table
SELECT *, FROM_UNIXTIME(`timestamp`, 'yyyy-MM-dd') from Kafka_source;
--------------------------------------------------------------------------------
指定固定的分区:
INSERT INTO hive_catalog.demo_db.test_hive_table PARTITION (dt = '2021-06-21')
SELECT * from Kafka_source;
通过该方式指定的分区需要是确定值,暂不支持计算。
a. 流模式

通过上述示例即可在流模式下写数据到Hive,需要注意的是,如果选择parquet这类块存储作为序列化类型,那么文件的切分建议oncheckpoint来执行,否则会有丢数据的风险。如果频繁的进行checkpoint,会导致小文件过多,建议作业的checkpoint间隔不要小于5分钟,如果作业数据量不大,可以调小并发,并增大checkpoint间隔至15分钟左右。在运行时,Flink的每一个task都会产生一个独立的文件,如果并发较大,也会受到小文件的影响,该问题目前可以通过Flink的小文件自动合并来解决,这部分内容和更多的Hive Streaming Sink的参数可以参考下文FileSystem Sink部分,因为其底层实际借助于FileSystem Sink来实现,参数也是通用的。

b. 批模式

在批模式下,只有当Flink作业结束,数据才会可见。批模式下还会额外支持**INSERT**`` OVERWRITE 语句。

FlinkSQL使用Hive表来进行维表join

原先,FlinkSQL支持的Hive维表Join场景比较单一,没有考虑该表是否是分区表,只是当检测到更新时进行全量的加载。1.12版本的Hive维表Join丰富了分区表的Join场景,支持Join Hive表的最新分区,同时也兼容以往Join latest Hive表的场景。维表Join的写法较之前并未发生变化,变化的只是Hive表的配置,在实际使用时可以使用SQL hints进行添加。下面给出官方文档上Join最新分区表的示例,更多内容可以参考下官方文档

sql
-- Assume the data in hive table is updated per day, every day contains the latest and complete dimension data
CREATE TABLE dimension_table (
  product_id STRING,
  product_name STRING,
  unit_price DECIMAL(10, 4),
  pv_count BIGINT,
  like_count BIGINT,
  comment_count BIGINT,
  update_time TIMESTAMP(3),
  update_user STRING,
  ...) PARTITIONED BY (pt_year STRING, pt_month STRING, pt_day STRING) TBLPROPERTIES (
  -- using default partition-name order to load the latest partition every 12h (the most recommended and convenient way)
  'streaming-source.enable' = 'true',
  'streaming-source.partition.include' = 'latest',
  'streaming-source.monitor-interval' = '12 h',
  'streaming-source.partition-order' = 'partition-name',  -- option with default value, can be ignored.

  -- using partition file create-time order to load the latest partition every 12h
  'streaming-source.enable' = 'true',
  'streaming-source.partition.include' = 'latest',
  'streaming-source.partition-order' = 'create-time',
  'streaming-source.monitor-interval' = '12 h'

  -- using partition-time order to load the latest partition every 12h
  'streaming-source.enable' = 'true',
  'streaming-source.partition.include' = 'latest',
  'streaming-source.monitor-interval' = '12 h',
  'streaming-source.partition-order' = 'partition-time',
  'partition.time-extractor.kind' = 'default',
  'partition.time-extractor.timestamp-pattern' = '$pt_year-$pt_month-$pt_day 00:00:00' );

上面给出的是创建Hive表的DDL,同时给出了三种配置,分别是按分区名、分区文件创建时间、分区时间来加载最新的分区。由于事先在数据工厂建好了表,所以此处可以省略。
--------------------------------------------------------------------------------
CREATE TABLE orders_table (
  order_id STRING,
  order_amount DOUBLE,
  product_id STRING,
  log_ts TIMESTAMP(3),
  proctime as PROCTIME()) WITH (...);

-- streaming sql, kafka temporal join a hive dimension table. Flink will automatically reload data from the configured latest partition in the interval of 'streaming-source.monitor-interval'.
SELECT * FROM orders_table AS order 
JOIN dimension_table /*+ OPTIONS('streaming-source.enable' = 'true') */
FOR SYSTEM_TIME AS OF order.proctime AS dim ON order.product_id = dim.product_id;

--------------------------------------------------------------------------------
如果Hive表是每天进行全量的更新,可以用一下参数来开启全量的加载
SELECT * FROM orders_table AS order 
JOIN dimension_table /*+ OPTIONS('streaming-source.enable' = 'false', 'streaming-source.partition.include' = 'all', 'lookup.join.cache.ttl' = '12 h') */
FOR SYSTEM_TIME AS OF order.proctime AS dim ON order.product_id = dim.product_id;

目前Hive维表Join实现逻辑也是通过将Hive表的所有数据加载到内存中来实现快速Join,使用时请确保TaskManager拥有足够的内存。

FlinkSQL通过FileSystem读写文件

上文主要介绍了FlinkSQL如何读写Hive,和Hive connector的一些新特性,由于Hive connector实际对文件的读写操作是通过FileSystem connector来实现的,本章简要介绍下FileSystem connector。很早之前,Flink就引入了FileSystem connector来读写Hdfs文件,在经过不断地迭代之后,自1.11版本开始,Filesystem connector被放在了Flink-table-runtime-blink项目下,使用时不再需要引入额外的connector包。FileSystem connector可以独立进行分布式文件系统的读写,不需要与Hive metastore进行通信,同时该connector也支持读写本地文件系统。由于不与Hive metastore进行通信,那么读写文件时需要先行建表,同时因为不能获取文件的schema,在建表时也需要指定format(目前支持CSV、JSON、Avro、Parqurt、Orc、Canal-JSON等主要的format)。下面给出一个示例:

sql
create table fs_sink (
    `timestamp` BIGINT,
    `time` STRING,
    id BIGINT,
    product STRING,
    price DOUBLE,
    canSell STRING,
    selledNum BIGINT,
    dt STRING,
    `hour` STRING,
    `min` STRING
) PARTITIONED BY (dt, `hour`, `min`) with (
    'connector'='filesystem',
    'path'='hdfs://nameservice1/user/XXXX/XXXXX',
    'format'='parquet',
    'sink.rolling-policy.rollover-interval'='30 s',
    'sink.partition-commit.policy.kind'='success-file'
);

insert into fs_sink
select * from source_table;

由DDL可见,Filesystem目前也支持指定分区。需要注意的是,目前FileSystem作为Source表暂不支持streaming场景下的增量读,仅支持全量读。 在FileSystem作为Sink表的场景下,情况就要复杂得多。目前FileSystem作为Sink表的配置主要有三个方面,分别是rolling策略,小文件合并以及分区提交,分区提交相关配置不再赘述,可以参考官方文档

a. Rolling Policy

在streaming模式下,需要定期对文件进行切分,将处在 in-progress状态的文件close,以保证数据的近实时可见,但如果切分过于频繁,又会导致小文件过多,对Hdfs非常不友好,这就需要额外的措施来进行平衡。FileSystem connector设计了下面展示的三个参数来配置切分策略。

Flink将FileSystem支持的format分为了两类,分别是Row Format和Bulk Format。 对于是否rolling on checkpoint,Flink通过 isBulkFormat || isOpenCompaction这两个参数来进行判断。 对于Row Format,例如:CSV、JSON,只要不配置小文件合并,默认可以不在checkpoint的时候进行切分,那么就可以保证每次切分大小基本相同。 但对于Bulk Format,例如:parquet、orc,不管开不开启小文件合并,默认都是rolling on checkpoint的。为什么Bulk Format默认需要rolling on checkpoint,以parquet举例。parquet写的是块文件,默认会一个bucket去写一次,在写入前这个块的数据会被缓存在内存中,默认大小是128M。因为parquet并未提供flush接口,如果在checkpoint的时候没有将文件落盘,数据还停留在内存中,在checkpoint完成后,作业失败会导致内存中的数据丢失。所以目前使用parquet需要去rolling on checkpoint。

b. File Compaction

上文提到Bulk Format需要rolling on checkpoint,那么checkpoint和sink.rolling-policy.file-size满足其1都会去切分文件,这导致部分文件没有达到预期的大小就会被切出,Flink又引入了小文件合并来解决这个问题。

通过上述配置可以开启小文件合并功能,但需要注意,合并的文件是一个checkpoint周期内所有task产生的文件,不同checkpoint 周期的文件并不会被合并。在添加完小文件合并的相关参数后,可以得到如下的 DAG(作业并发为 3):

FileSystem Sink 开启小文件合并后的 DAG

由DAG图可以看到,Stream-writer算子逻辑为向Hdfs文件写数据,三个并发会向指定目录写三个inprogress文件,当做完checkpoint后会告诉下游compact-coordinator节点,coordinator节点在收到上游所有节点的checkpoint完成通知后,会筛选出需要compact的文件发往下游compact-operator进行合并,合并完后直接commit。 从上述流程不难看出,目前小文件合并功能只作用于同一个checkpoint周期,如果作业并发较小且checkpoint周期较短,也很难合并出理想大小的文件。


本文是「我的数据空间」实时计算实践笔记的一篇。平台把这类链路做成了托管作业(Flink SQL 流作业、实时集成、Iceberg 表管理与维护),支持私有化部署与 OEM 合作 —— 见核心能力总览技术博客,交流合作 QQ:1559851993。