
# FlinkServer作业对接Hudi表
本章节适用于MRS 3.1.2及之后的版本。
#### 操作场景
MRS大数据平台中，Hudi作为数据湖表格式支持流批一体读写，而Flink是主流的流计算引擎。当用户需要将Kafka实时数据流式写入Hudi表，或从Hudi表流式读取数据进行实时分析时，需要通过FlinkServer的FlinkSQL接口实现对接。本指南通过使用FlinkServer写FlinkSQL对接Hudi，涵盖MOR/COW表读写、Lookup Join维表关联、Hive元数据同步等场景。FlinkSQL读写Hudi时，不支持定义TINYINT、SMALLINT和TIME类型。
Flink对Hudi表的COW表、MOR表类型读写支持详情见[表1]。
FlinkSQL读写Hudi表参数简化可参考[Hudi创建表参数简化指导](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_300859.html)。
 表1Flink对Hudi表的读写支持 
| Flink SQL | COW表 | MOR表 |
|:---|:---|:---|
| 批量写       | 支持   | 支持   |
| 批量读       | 支持   | 支持   |
| 流式写       | 支持   | 支持   |
| 流式读       | 支持   | 支持   |
   
#### 基本概念
- **COW表**（Copy On Write，写时复制）：Hudi的一种表类型，每次更新操作都会将新数据与已有基础文件合并，生成新的Parquet列式存储文件。写入开销较大但读取性能好，适合读多写少的批量场景。
- **MOR表**（Merge On Read，读时合并）：Hudi的一种表类型，写入时将新数据追加到Avro行式日志文件中，读取时将日志文件与Parquet基础文件合并。写入开销小但读取需额外合并，适合流式高频写入场景。
- **CheckPoint**：Flink的周期性状态快照机制，将算子状态持久化到外部存储。Flink写Hudi表时，数据在CheckPoint触发时才执行Hudi Commit操作提交到HDFS，因此CheckPoint是写入Hudi表的必要条件。
 
#### 前提条件
- 集群已安装HDFS、Yarn、Flink和Hudi等服务。
- 包含Hudi服务的客户端已安装，例如安装路径为：/opt/client。
- Flink要求1.12.2及以后版本，Hudi要求0.9.0及以后版本。
- 参考[创建FlinkServer权限角色](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_24049.html)创建一个具有FlinkServer管理员权限的用户用于访问Flink WebUI，如：flink_admin。并且用户需要添加hadoop、hive、kafkaadmin用户组，以及Manager_administrator角色。
 
#### 流式写入和读取Hudi表操作步骤
1. 使用**flink_admin**登录Manager，选择"集群 \> 服务 \> Flink"，在"Flink WebUI"右侧，单击链接，访问Flink的WebUI。
2. 参考[集群连接模式创建Flink SQL作业](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_24024.html)，新建Flink SQL流作业，在作业开发界面进行如下作业配置。然后输入SQL，执行SQL校验通过后，启动作业。如下SQL示例将作为3个作业分别添加，依次运行。
   
   需勾选"基础参数"中的"开启CheckPoint"，"时间间隔（ms）"可设置为"60000"，"模式"可使用默认值。
   ![](https://support.huaweicloud.com/cmpntguide-lts-mrs/public_sys-resources/note_3.0-zh-cn.png)
   - 由于FlinkSQL作业在触发CheckPoint时才会往Hudi表中写数据，所以需要在Flink WebUI界面中开启CheckPoint。CheckPoint间隔根据业务需要调整，建议间隔调大。
   
   - 如果CheckPoint间隔太短，数据来不及刷新会导致作业异常；建议CheckPoint间隔为分钟级。
   
   - FlinkSQL作业写MOR表时需要做异步compaction，控制compaction间隔的参数，见Hudi官网：https://hudi.apache.org/docs/configurations.html
   
   - MRS 3.2.1及以后版本默认Hudi写表是Flink状态索引，如果需要使用bucket索引需要在Hudi写表中添加参数：
     ```
     'index.type'='BUCKET',
     'hoodie.bucket.index.num.buckets'='Hudi表中每个分区划分桶的个数',
     'hoodie.bucket.index.hash.field'='recordkey.field'
     ```
     - hoodie.bucket.index.num.buckets：Hudi表中每个分区划分桶的个数，每个分区内的数据通过Hash方式放入每个桶内。建表或第一次写入数据时设置后不能修改，否则更新数据会存在异常。
     
     - hoodie.bucket.index.hash.field：进行分桶时计算Hash值的字段，必须为主键的子集，默认为Hudi表的主键。该参数不填则默认为recordkey.field。
      
   
   - MRS 3.2.1及以后版本对于同一张Hudi表，可以被Flink、Spark引擎的bucket索引交叉混写。
    
   1. 作业1：FlinkSQL流式写入MOR表。
      ```
      CREATE TABLE stream_mor(
      uuid VARCHAR(20),
      name VARCHAR(10),
      age INT,
      ts INT,
      `p` VARCHAR(20)
      ) PARTITIONED BY (`p`) WITH (
      'connector' = 'hudi',
      'path' = 'hdfs://hacluster/tmp/hudi/stream_mor',
      'table.type' = 'MERGE_ON_READ',
      'hoodie.datasource.write.recordkey.field' = 'uuid',
      'write.precombine.field' = 'ts',
      'write.tasks' = '4'
      );
      CREATE TABLE kafka(
      uuid VARCHAR(20),
      name VARCHAR(10),
      age INT,
      ts INT,
      `p` VARCHAR(20)
      ) WITH (
      'connector' = 'kafka',
      'topic' = 'writehudi',
      'properties.bootstrap.servers' = 'Kafka的Broker实例业务IP:Kafka端口号',
      'properties.group.id' = 'testGroup1',
      'scan.startup.mode' = 'latest-offset',
      'format' = 'json',
      'properties.sasl.kerberos.service.name' = 'kafka',--普通模式集群不需要该参数，同时删除上一行的逗号
      'properties.security.protocol' = 'SASL_PLAINTEXT',--普通模式集群不需要该参数
      'properties.kerberos.domain.name' = 'hadoop.系统域名'--普通模式集群不需要该参数
      );
      insert into
      stream_mor
      select
      *
      from
      kafka;
      ```
      
   
   2. 作业2：FlinkSQL流式写入COW表。
      ```
      CREATE TABLE stream_write_cow(
      uuid VARCHAR(20),
      name VARCHAR(10),
      age INT,
      ts INT,
      `p` VARCHAR(20)
      ) PARTITIONED BY (`p`) WITH (
      'connector' = 'hudi',
      'path' = 'hdfs://hacluster/tmp/hudi/stream_cow',
      'hoodie.datasource.write.recordkey.field' = 'uuid',
      'write.precombine.field' = 'ts',
      'write.tasks' = '4'
      );
      CREATE TABLE kafka(
      uuid VARCHAR(20),
      name VARCHAR(10),
      age INT,
      ts INT,
      `p` VARCHAR(20)
      ) WITH (
      'connector' = 'kafka',
      'topic' = 'writehudi',
      'properties.bootstrap.servers' = 'Kafka的Broker实例业务IP:Kafka端口号',
      'properties.group.id' = 'testGroup1',
      'scan.startup.mode' = 'latest-offset',
      'format' = 'json',
      'properties.sasl.kerberos.service.name' = 'kafka',--普通模式集群不需要该参数，同时删除上一行的逗号
      'properties.security.protocol' = 'SASL_PLAINTEXT',--普通模式集群不需要该参数
      'properties.kerberos.domain.name' = 'hadoop.系统域名'--普通模式集群不需要该参数
      );
      insert into
      stream_write_cow
      select
      *
      from
      kafka;
      ```
      
   
   3. 作业3：FlinkSQL流式读取MOR和COW表并合并数据输出到Kafka。注意作业3需要等待作业1和作业2均启动后，状态显示为"运行中"后再执行SQL校验和启动作业（否则SQL校验可能提示错误，找不到Hudi表目录）。
      ```
      CREATE TABLE stream_mor(
      uuid VARCHAR(20),
      name VARCHAR(10),
      age INT,
      ts INT,
      `p` VARCHAR(20)
      ) PARTITIONED BY (`p`) WITH (
      'connector' = 'hudi',
      'path' = 'hdfs://hacluster/tmp/hudi/stream_mor',
      'table.type' = 'MERGE_ON_READ',
      'hoodie.datasource.write.recordkey.field' = 'uuid',
      'write.precombine.field' = 'ts',
      'read.tasks' = '4',
      'read.streaming.enabled' = 'true',
      'read.streaming.check-interval' = '5',
      'read.streaming.start-commit' = 'earliest'
      );
      CREATE TABLE stream_write_cow(
      uuid VARCHAR(20),
      name VARCHAR(10),
      age INT,
      ts INT,
      `p` VARCHAR(20)
      ) PARTITIONED BY (`p`) WITH (
      'connector' = 'hudi',
      'path' = 'hdfs://hacluster/tmp/hudi/stream_cow',
      'hoodie.datasource.write.recordkey.field' = 'uuid',
      'write.precombine.field' = 'ts',
      'read.tasks' = '4',
      'read.streaming.enabled' = 'true',
      'read.streaming.check-interval' = '5',
      'read.streaming.start-commit' = 'earliest'
      );
      CREATE TABLE kafka(
      uuid VARCHAR(20),
      name VARCHAR(10),
      age INT,
      ts INT,
      `p` VARCHAR(20)
      ) WITH (
      'connector' = 'kafka',
      'topic' = 'readhudi',
      'properties.bootstrap.servers' = 'Kafka的Broker实例业务IP:Kafka端口号',
      'properties.group.id' = 'testGroup1',
      'scan.startup.mode' = 'latest-offset',
      'format' = 'json',
      'properties.sasl.kerberos.service.name' = 'kafka',--普通模式集群不需要该参数，同时删除上一行的逗号
      'properties.security.protocol' = 'SASL_PLAINTEXT',--普通模式集群不需要该参数
      'properties.kerberos.domain.name' = 'hadoop.系统域名'--普通模式集群不需要该参数
      );
      insert into 
      kafka 
      select 
      * 
      from 
      stream_mor union all select * from stream_write_cow;
      ```
      
   
   
   ![](https://support.huaweicloud.com/cmpntguide-lts-mrs/public_sys-resources/note_3.0-zh-cn.png)
   - Kafka Broker实例IP地址及Kafka端口号：
     - 服务的实例IP地址可通过登录FusionInsight Manager后，单击"集群 \> 服务 \> Kafka \> 实例"，在实例列表页面中查询。
     
     - 集群的"认证模式"为"安全模式"时为"sasl.port"的值，默认为"21007"。
     
     - 集群的"认证模式"为"普通模式"时为"port"的值，默认为"9092"。如果配置端口号为9092，则需要配置"allow.everyone.if.no.acl.found"参数为true，具体操作如下： 登录FusionInsight Manager系统，选择"集群 \> 服务 \> Kafka \> 配置 \> 全部配置"，搜索"allow.everyone.if.no.acl.found"配置，修改参数值为true，保存配置即可。
       
      
   
   - 系统域名：可登录FusionInsight Manager，选择"系统 \> 权限 \> 域和互信"，查看"本端域"参数，即为当前系统域名。
    
   
   
3. 参考[使用Kafka生产消费数据](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_0379.html)，向kafka中写入数据。
   
   ```
   sh kafka-console-producer.sh --broker-list Kafka角色实例所在节点的IP地址:Kafka端口号 --topic 主题名称 --producer.config 客户端目录/Kafka/kafka/config/producer.properties
   ```
   例如本示例使用主题名称为"writehudi"：
   ```
   sh kafka-console-producer.sh --broker-list Kafka角色实例所在节点的IP地址:Kafka端口号 --topic writehudi --producer.config /opt/client/Kafka/kafka/config/producer.properties
   ```
   输入消息内容：
   ```
   {"uuid": "1","name":"a01","age":10,"ts":10,"p":"1"}
   {"uuid": "2","name":"a02","age":20,"ts":20,"p":"2"}
   ```
   输入完成后按回车发送消息。
   
   
4. 消费kafka topic的数据，读取Flink流读Hudi表的结果。 
   ```
   sh kafka-console-consumer.sh --bootstrap-server Kafka的Broker实例所在节点的IP地址:Kafka端口号 --topic 主题名称 --consumer.config 客户端目录/Kafka/kafka/config/consumer.properties --from-beginning
   ```
   例如本示例使用主题名称为"readhudi"：
   ```
   sh kafka-console-consumer.sh --bootstrap-server Kafka的Broker实例所在节点的IP地址:Kafka端口号 --topic readhudi --consumer.config /opt/client/Kafka/kafka/config/consumer.properties --from-beginning
   ```
   读取结果如下（顺序不固定）：
   ```
   {"uuid": "1","name":"a01","age":10,"ts":10,"p":"1"}
   {"uuid": "2","name":"a02","age":20,"ts":20,"p":"2"}
   {"uuid": "1","name":"a01","age":10,"ts":10,"p":"1"}
   {"uuid": "2","name":"a02","age":20,"ts":20,"p":"2"}
   ```
   如果Kafka消费端能持续读取到Hudi表流式读取的结果数据，且数据内容与Kafka生产端写入的数据一致，则表示Flink作业对接Hudi表成功。
   
   
 
#### FlinkSQL Lookup Join Hudi使用须知
适用于MRS 3.5.0及以后版本。
- 使用lookup.join.cache.ttl参数来控制维表数据的加载周期，默认值为60min。
- Hudi维表数据会被加载到Flink TaskManager Heap中，所以不推荐大于10万行记录的Hudi表作为维表。
- 维表的新增、更新数据需要等到下一次加载周期后，才能被加载进来参与计算。
SQL示例如下：
```
CREATE TABLE hudimor(
  uuid VARCHAR(20),
  name VARCHAR(10),
  age INT,
  ts INT,
  `p` VARCHAR(20),
  PRIMARY KEY (uuid) NOT ENFORCED
) PARTITIONED BY (`p`) WITH (
  'connector' = 'hudi',
  'path' = 'hdfs://hacluster/tmp/hudimor',
  'table.type' = 'MERGE_ON_READ',
  'hoodie.datasource.write.recordkey.field' = 'uuid',
  'write.precombine.field' = 'ts',
  'lookup.join.cache.ttl' = '60min'
);
CREATE TABLE datagen(uuid varchar(20), proctime as PROCTIME()) WITH (
  'connector' = 'datagen',
  'rows-per-second' = '1'
);
CREATE TABLE blackhole (
  uuid VARCHAR(20),
  name VARCHAR(10),
  age INT,
  ts INT,
  `p` VARCHAR(20)
) WITH ('connector' = 'blackhole');
insert into
  blackhole
select
  t1.uuid as uuid,
  t2.name as name,
  t2.age as age,
  t2.ts as ts,
  t2.p as p
FROM
  datagen AS t1
  left JOIN hudimor FOR SYSTEM_TIME AS OF t1.proctime AS t2 ON t1.uuid = t2.uuid;
```
#### WITH主要参数说明
表2WITH主要参数说明 
| 方式                                                                         | 配置项                         | 是否必选 | 默认值         | 描述                                                                                                                                 |
|:---|:---|:---|:---|:---|
| 读取                       | read.tasks                  | 否    | 1           | 读Hudi表的task并行度，和FlinkServer作业开发界面基础参数配置的并行度一致。                                                                                     |
| 读取                       | read.streaming.enabled      | 否    | false       | 是否开启流读模式。                                                                                                                          |
| 读取                       | read.streaming.start-commit | 否    | 默认从最新commit | Stream和Batch增量消费，指定"yyyyMMddHHmmss"格式时间的开始消费位置（闭区间）。                                                                               |
| 读取                       | read.end-commit             | 否    | 默认到最新commit | Stream和Batch增量消费，指定"yyyyMMddHHmmss"格式时间的结束消费位置（闭区间）。                                                                               |
| 写入       | write.tasks                 | 否    | 1           | 写Hudi表的task并行度，和FlinkServer作业开发界面基础参数配置的并行度一致。                                                                                     |
| 写入       | index.bootstrap.enabled     | 否    | false       | 是否开启索引加载，开启后会将已存表的最新数据一次性加载到state中。 如果有全量数据接增量的需求，且已经有全量的离线Hoodie表，需要接上实时写入，同时保证数据不重复，可以开启索引加载功能。 |
| 写入       | write.index_bootstrap.tasks | 否    | 4           | 如果启动作业时索引加载缓慢，可以调大该值，调大该值后可以加快bootstrap阶段的效率，但bootstrap阶段会阻塞CheckPoint。                                                            |
| 写入       | compaction.async.enabled    | 否    | true        | 是否开启在线压缩。                                                                                                                          |
| 写入       | compaction.schedule.enabled | 否    | true        | 是否阶段性生成压缩plan，即使关闭在线压缩的情况下也建议开启。                                                                                                   |
| 写入       | compaction.tasks            | 否    | 10          | 压缩Hudi表task并行度。                                                                                                                    |
| 写入       | index.state.ttl             | 否    | 7D          | 索引保存的时间，默认为7天（单位：天），小于"0"表示永久保存 索引是判断数据重复的核心数据结构，对于长时间的更新，比如更新一个月前的数据，需要将该值调大。                       |
   
#### Flink On Hudi同步元数据到Hive
启动此特性后，Flink写数据至Hudi表将自动在Hive上创建出Hudi表并同步添加分区，然后供SparkSQL、Hive等服务读取Hudi表数据。
如下是支持的两种同步元数据方式，后续操作步骤以JDBC方式为示例：
适用于MRS 3.2.0及之后版本。
- 使用JDBC方式同步元数据到Hive
  ```
  CREATE TABLE stream_mor(
  uuid VARCHAR(20),
  name VARCHAR(10),
  age INT,
  ts INT,
  `p` VARCHAR(20)
  ) PARTITIONED BY (`p`) WITH (
  'connector' = 'hudi',
  'path' = 'hdfs://hacluster/tmp/hudi/stream_mor',
  'table.type' = 'MERGE_ON_READ',
  'hive_sync.enable' = 'true',
  'hive_sync.table' = '要同步到Hive的表名',
  'hive_sync.db' = '要同步到Hive的数据库名',
  'hive_sync.metastore.uris' = 'Hive客户端hive-site.xml文件中hive.metastore.uris的值',
  'hive_sync.jdbc_url' = 'Hive客户端component_env文件中CLIENT_HIVE_URI的值'
  );
  ```
  ![](https://support.huaweicloud.com/cmpntguide-lts-mrs/public_sys-resources/notice_3.0-zh-cn.png)
  - hive_sync.jdbc_url：Hive客户端component_env文件中CLIENT_HIVE_URI的值，如果该值中存在"\\"需将其删除。
  
  - 如果需要使用Hive风格分区，需同时配置如下参数：
    - 'hoodie.datasource.write.hive_style_partitioning' = 'true'
    
    - 'hive_sync.partition_extractor_class' = 'org.apache.hudi.hive.MultiPartKeysValueExtractor'
     
  
  - Flink on Hudi并同步数据至Hive的任务，因为Hudi对大小写敏感，Hive对大小写不敏感，所以在Hudi表中的字段不建议使用大写字母，否则可能会造成数据无法正常读写。
    
- 使用HMS方式同步元数据到Hive
  ```
  CREATE TABLE stream_mor(
  uuid VARCHAR(20),
  name VARCHAR(10),
  age INT,
  ts INT,
  `p` VARCHAR(20)
  ) PARTITIONED BY (`p`) WITH (
  'connector' = 'hudi',
  'path' = 'hdfs://hacluster/tmp/hudi/stream_mor',
  'table.type' = 'MERGE_ON_READ',
  'hive_sync.enable' = 'true',
  'hive_sync.table' = '要同步到Hive的表名',
  'hive_sync.db' = '要同步到Hive的数据库名',
  'hive_sync.mode' = 'hms',
  'hive_sync.metastore.uris' = 'Hive客户端hive-site.xml文件中hive.metastore.uris的值',
  'properties.hive.metastore.kerberos.principal' = 'Hive客户端hive-site.xml文件中hive.metastore.kerberos.principal的值'
  );
  ```
  
 
使用JDBC方式同步元数据到Hive示例：
1. 使用**flink_admin**登录Manager，选择"集群 \> 服务 \> Flink"，在"Flink WebUI"右侧，单击链接，访问Flink的WebUI。
2. 参考[集群连接模式创建Flink SQL作业](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_24024.html)，新建Flink SQL流作业，在作业开发界面进行如下作业配置。然后填入SQL，执行SQL校验通过后，启动作业。
   
   需勾选"基础参数"中的"开启CheckPoint"，"时间间隔（ms）"可设置为"60000"，"模式"可使用默认值。
   ```
   CREATE TABLE stream_mor2(
   uuid VARCHAR(20),
   name VARCHAR(10),
   age INT,
   ts INT,
   `p` VARCHAR(20)
   ) PARTITIONED BY (`p`) WITH (
   'connector' = 'hudi',
   'path' = 'hdfs://hacluster/tmp/hudi/stream_mor2',
   'table.type' = 'MERGE_ON_READ',
   'hoodie.datasource.write.recordkey.field' = 'uuid',
   'write.precombine.field' = 'ts',
   'write.tasks' = '4',
   'hive_sync.enable' = 'true',
   'hive_sync.table' = '要同步到Hive的表名，如stream_mor2',
   'hive_sync.db' = '要同步到Hive的数据库名，如default',
   'hive_sync.metastore.uris' = 'Hive客户端hive-site.xml文件中hive.metastore.uris的值',
   'hive_sync.jdbc_url' = 'Hive客户端component_env文件中CLIENT_HIVE_URI的值'
   );
   CREATE TABLE datagen (
   uuid varchar(20), name varchar(10), age int, ts INT, p varchar(20)
   ) WITH (
   'connector' = 'datagen',
   'rows-per-second' = '1',
   'fields.p.length' = '1'
   );insert into stream_mor2 select * from datagen;
   ```
   
   
3. 等待Flink作业运行一段时间，将datagen生成的随机测试数据持续写入Hudi表。可通过"更多 \> 作业详情"跳转到Flink作业的原生UI页面，查看Job运行情况。
4. 登录客户端所在节点，加载环境变量，执行beeline命令登录Hive客户端，执行SQL查看是否在Hive上成功创建Hudi Sink表，并且查询表是否可读出数据。 
   ```
   cd /opt/hadoopclient
   source bigdata_env
   beeline
   desc formatted default.stream_mor2;
   select * from default.stream_mor2 limit 5;
   show partitions default.stream_mor2;
   ```
   
   
#### FlinkSQL建Hudi表支持字段注释
通过Hudi Catalog创建的Hudi表，支持同步字段注释信息到Hive。
SQL示例如下：
```
CREATE CATALOG hoodie_catalog
  WITH (
    'type' = 'hudi',
    'default-database' ='default',
    'catalog.path' = 'hdfs:///tmp/',--catalog表默认存储路径
    'mode' = 'hms',
    'cluster.name' ='hudiTest'
  );
use catalog hoodie_catalog;
CREATE TABLE stream_mor(
       id int comment '主键',
       name VARCHAR(20) comment '名字',
       age INT comment '年龄',
       `date` VARCHAR(20)
) PARTITIONED BY (`date`) WITH (
  'connector' = 'hudi',
  'path' = 'hdfs://hacluster/tmp/hudi_mor',--Hudi表存储路径
  'table.type' = 'MERGE_ON_READ',
  'hoodie.datasource.write.recordkey.field' = 'id',
  'write.precombine.field' = 'age',
  'hoodie.datasource.write.hive_style_partitioning' = 'true',
  'hive_sync.partition_extractor_class' = 'org.apache.hudi.hive.MultiPartKeysValueExtractor',
  'index.type' =  'BUCKET',
  'hoodie.bucket.index.num.buckets' = '8',
  'write.tasks' = '8'
);
```
#### 相关文档
- 如何在FlinkServer中创建Flink SQL作业，请参考[集群连接模式创建Flink SQL作业](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_24024.html)。
- 如何创建具有FlinkServer管理员权限的用户，请参考[创建FlinkServer权限角色](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_24049.html)。
- FlinkSQL读写Hudi表参数简化，请参考[Hudi创建表参数简化指导](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_300859.html)。
 
