为何无法在Flink中使用Datagen连接器完成数据插入?
Flink Datagen写入Hudi无数据问题排查
问题现象
通过Flink SQL创建Datagen源表sourceT和Hudi表t2,insert语句提交成功,但查询t2时显示「结果检索已取消」,未获取到数据。
具体操作与执行日志
Flink SQL> CREATE TABLE sourceT ( > uuid varchar(20), > name varchar(10), > age int, > ts timestamp(3), > `partition` varchar(20) > ) WITH ( > 'connector' = 'datagen', > 'rows-per-second' = '1' > ); [INFO] 执行语句成功。 Flink SQL> create table t2( > uuid varchar(20), > name varchar(10), > age int, > ts timestamp(3), > `partition` varchar(20) > ) > with ( > 'connector' = 'hudi', > 'path' = '/tmp/hudi_flink/t2', > 'table.type' = 'MERGE_ON_READ' > ); [INFO] 执行语句成功。 Flink SQL> insert into t2 select * from sourceT; [INFO] 正在向集群提交SQL更新语句... 2022-11-29 11:55:39,776 WARN org.apache.flink.yarn.configuration.YarnLogConfigUtil [] - 配置目录('/opt/module/flink-1.13.6/conf')已包含LOG4J配置文件。若要使用logback,请删除或重命名日志配置文件。 2022-11-29 11:55:40,018 INFO org.apache.hadoop.yarn.client.RMProxy [] - 正在连接至hadoop163/192.168.10.163:8032上的ResourceManager 2022-11-29 11:55:40,151 INFO org.apache.flink.yarn.YarnClusterDescriptor [] - 未传入Flink Jar包路径,将通过org.apache.flink.yarn.YarnClusterDescriptor类的位置定位Jar包 2022-11-29 11:55:40,224 INFO org.apache.flink.yarn.YarnClusterDescriptor [] - 已找到应用程序'application_1669685396475_0006'的Web界面hadoop163:45413。 [INFO] SQL更新语句已成功提交至集群: Job ID: 9f2283c2f6d943c170068abe39747bc0 Flink SQL> select * from t2; 2022-11-29 11:55:57,226 INFO org.apache.hadoop.hdfs.protocol.datatransfer.sasl.SaslDataTransferClient [] - SASL加密信任检查:localHostTrusted = false,remoteHostTrusted = false 2022-11-29 11:55:57,516 INFO org.apache.hadoop.yarn.client.RMProxy [] - 正在连接至hadoop163/192.168.10.163:8032上的ResourceManager 2022-11-29 11:55:57,516 INFO org.apache.flink.yarn.YarnClusterDescriptor [] - 未传入Flink Jar包路径,将通过org.apache.flink.yarn.YarnClusterDescriptor类的位置定位Jar包 2022-11-29 11:55:57,520 INFO org.apache.flink.yarn.YarnClusterDescriptor [] - 已找到应用程序'application_1669685396475_0006'的Web界面hadoop163:45413。 [INFO] 结果检索已取消。
排查与解决步骤
1. 补全Hudi表主键配置
Hudi的MERGE_ON_READ类型表必须指定主键,否则写入逻辑会异常。修改t2表的创建语句,添加主键配置:
create table t2( uuid varchar(20), name varchar(10), age int, ts timestamp(3), `partition` varchar(20) ) with ( 'connector' = 'hudi', 'path' = '/tmp/hudi_flink/t2', 'table.type' = 'MERGE_ON_READ', 'primary.key' = 'uuid' -- 新增主键配置 );
删除原表后重新创建,再提交insert作业。
2. 确认Insert作业运行状态
- 访问Flink Web UI(地址
hadoop163:45413),查看Job ID9f2283c2f6d943c170068abe39747bc0的作业:- 检查作业是否正常运行,无失败、重启情况
- 查看Metrics中的
numRecordsIn/numRecordsOut指标,确认数据是否从Datagen流入并写入Hudi
3. 直接验证HDFS数据文件
执行HDFS命令查看目标路径下是否生成数据文件:
hdfs dfs -ls /tmp/hudi_flink/t2 hdfs dfs -ls /tmp/hudi_flink/t2/base # 查看base目录文件 hdfs dfs -ls /tmp/hudi_flink/t2/log # 查看log目录文件
如果存在文件,说明数据已写入,问题出在查询环节;如果无文件,说明写入逻辑有问题。
4. 调整查询参数
- 尝试带限制的查询:
select * from t2 limit 10;
- 若为流查询,添加流读取配置:
select * from t2 /*+ OPTIONS('streaming'='true', 'read.streaming.check-interval'='1s') */;
5. 查看作业详细日志
通过YARN命令获取作业日志,排查写入异常:
yarn logs -applicationId application_1669685396475_0006 | grep -i hudi
内容的提问来源于stack exchange,提问作者Jiangchao Yang
相关产品推荐
相关产品推荐

