使用Flink SQL将Hudi表写入Minio S3桶失败求助
Flink SQL写入Hudi到Minio S3失败问题
问题描述
尝试用Flink SQL将数据写入存储在Minio S3上的Hudi表时失败,Hudi表虽已创建,但仅生成了.hoodie元数据目录,结构如下:
myminio/flink-hudi └─ t1 └─ .hoodie ├─ .aux │ ├─ .bootstrap │ │ ├─ .fileids │ │ └─ .partitions │ └─ ckp_meta ├─ .schema ├─ .temp └─ archived
复现步骤
- 创建Flink Hudi表
CREATE TABLE t1( uuid VARCHAR(20) PRIMARY KEY NOT ENFORCED, name VARCHAR(10), age INT, ts TIMESTAMP(3), `partition` VARCHAR(20) ) PARTITIONED BY (`partition`) WITH ( 'connector' = 'hudi', 'path' = 's3a://flink-hudi/t1', 'table.type' = 'MERGE_ON_READ' );
- 插入数据到Hudi表
INSERT INTO t1 VALUES ('id1','Danny',23,TIMESTAMP '1970-01-01 00:00:01','par1');
环境信息
- Hudi版本:0.12.0
- Hadoop版本:3.2.4
- Flink版本:1.15.2
- 存储:Minio S3
- 部署环境:非Docker
配置与依赖
已添加的依赖包
- hadoop-aws-3.2.4.jar
- aws-java-sdk-bundle-1.11.901.jar
- flink-s3-fs-hadoop-1.15.2.jar
Hadoop core-site.xml配置
<property> <name>fs.s3a.access.key</name> <value>xxx</value> </property> <property> <name>fs.s3a.secret.key</name> <value>xxx</value> </property> <property> <name>fs.s3a.endpoint</name> <value>xxx</value> </property> <property> <name>fs.s3a.path.style.access</name> <value>true</value> </property> <property> <name>fs.s3.impl</name> <value>org.apache.hadoop.fs.s3a.S3AFileSystem</value> </property>
Flink flink-conf.yaml配置
taskmanager.numberOfTaskSlots: 4 s3a.endpoint: xxx s3a.access-key: xxx s3a.secret-key: xxx s3a.path.style.access: true fs.hdfs.hadoopconf: /export/servers/hadoop-3.2.4/etc/hadoop state.backend: rocksdb state.backend.incremental: true state.checkpoints.dir: s3a://flink-state/checkpoint execution.checkpointing.interval: 30000 classloader.check-leaked-classloader: false
Flink启动命令
export HADOOP_CLASSPATH=`$HADOOP_HOME/bin/hadoop classpath` ./bin/start-cluster.sh ./bin/sql-client.sh embedded -j /opt/flink/jars/hudi-flink1.15-bundle-0.12.0.jar shell
错误堆栈
org.apache.hudi.exception.HoodieException: Exception while scanning the checkpoint meta files under path: s3a://flink-hudi/t1/.hoodie/.aux/ckp_meta at org.apache.hudi.sink.meta.CkpMetadata.load(CkpMetadata.java:169) at org.apache.hudi.sink.meta.CkpMetadata.lastPendingInstant(CkpMetadata.java:175) at org.apache.hudi.sink.common.AbstractStreamWriteFunction.lastPendingInstant(AbstractStreamWriteFunction.java:243) at org.apache.hudi.sink.common.AbstractStreamWriteFunction.initializeState(AbstractStreamWriteFunction.java:151) at org.apache.flink.streaming.util.functions.StreamingFunctionUtils.tryRestoreFunction(StreamingFunctionUtils.java:189) at org.apache.flink.streaming.util.functions.StreamingFunctionUtils.restoreFunctionState(StreamingFunctionUtils.java:171) at org.apache.flink.streaming.api.operators.AbstractUdfStreamOperator.initializeState(AbstractUdfStreamOperator.java:94) at org.apache.flink.streaming.api.operators.StreamOperatorStateHandler.initializeOperatorState(StreamOperatorStateHandler.java:122) at org.apache.flink.streaming.api.operators.AbstractStreamOperator.initializeState(AbstractStreamOperator.java:286) at org.apache.flink.streaming.runtime.tasks.RegularOperatorChain.initializeStateAndOpenOperators(RegularOperatorChain.java:106) at org.apache.flink.streaming.runtime.tasks.StreamTask.restoreGates(StreamTask.java:700) at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.call(StreamTaskActionExecutor.java:55) at org.apache.flink.streaming.runtime.tasks.StreamTask.restoreInternal(StreamTask.java:676) at org.apache.flink.streaming.runtime.tasks.StreamTask.restore(StreamTask.java:643) at org.apache.flink.runtime.taskmanager.Task.runWithSystemExitMonitoring(Task.java:948) at org.apache.flink.runtime.taskmanager.Task.restoreAndInvoke(Task.java:917) at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:741) at org.apache.flink.runtime.taskmanager.Task.run(Task.java:563) at java.lang.Thread.run(Thread.java:748) Caused by: java.io.FileNotFoundException: No such file or directory: s3a://flink-hudi/t1/.hoodie/.aux/ckp_meta at org.apache.hadoop.fs.s3a.S3AFileSystem.s3GetFileStatus(S3AFileSystem.java:2344) at org.apache.hadoop.fs.s3a.S3AFileSystem.innerGetFileStatus(S3AFileSystem.java:2226) at org.apache.hadoop.fs.s3a.S3AFileSystem.getFileStatus(S3AFileSystem.java:2160) at org.apache.hadoop.fs.s3a.S3AFileSystem.innerListStatus(S3AFileSystem.java:1961) at org.apache.hadoop.fs.s3a.S3AFileSystem.lambda$listStatus$9(S3AFileSystem.java:1940) at org.apache.hadoop.fs.s3a.Invoker.once(Invoker.java:109) at org.apache.hadoop.fs.s3a.S3AFileSystem.listStatus(S3AFileSystem.java:1940) at org.apache.hudi.common.fs.HoodieWrapperFileSystem.lambda$listStatus$15(HoodieWrapperFileSystem.java:365) at org.apache.hudi.common.fs.HoodieWrapperFileSystem.executeFuncWithTimeMetrics(HoodieWrapperFileSystem.java:106) at org.apache.hudi.common.fs.HoodieWrapperFileSystem.listStatus(HoodieWrapperFileSystem.java:364) at org.apache.hudi.sink.meta.CkpMetadata.scanCkpMetadata(CkpMetadata.java:216) at org.apache.hudi.sink.meta.CkpMetadata.load(CkpMetadata.java:167) ... 18 more
排查与解决方案
从错误堆栈看,核心问题是Hudi扫描ckp_meta目录时找不到该路径,结合Minio是对象存储的特性,给出以下排查方向和解决办法:
1. 验证Minio权限与路径
- 用Minio控制台或
mc工具检查使用的access key/secret key是否拥有flink-hudibucket的读写、创建对象/目录权限 - 确认
fs.s3a.endpoint配置正确,能正常访问Minio服务;检查路径s3a://flink-hudi/t1是否存在,手动创建该路径下的ckp_meta目录后重试插入操作
2. 调整Hadoop S3A客户端配置
由于S3是对象存储,没有真正的目录结构,Hadoop S3A客户端需要通过标记文件模拟目录,在core-site.xml中添加以下配置:
<property> <name>fs.s3a.directory.marker.retention</name> <value>keep</value> </property> <property> <name>fs.s3a.directory.marker.create</name> <value>true</value> </property> <!-- 如果Minio未开启SSL,添加下面配置 --> <property> <name>fs.s3a.connection.ssl.enabled</name> <value>false</value> </property>
3. 清理旧元数据后重试
删除Minio中flink-hudi/t1下的所有内容,重新创建Hudi表并执行插入操作,避免旧的不完整元数据干扰
4. 检查依赖兼容性与冲突
- 确认
hadoop-aws、aws-java-sdk-bundle版本与Hadoop 3.2.4完全匹配,避免版本不一致导致的API调用错误 - 检查Flink lib目录下是否存在重复的aws-sdk或hadoop-aws包,如有则删除多余包
5. 测试关闭Checkpoint
临时关闭Flink的checkpoint(注释execution.checkpointing.interval配置),执行插入操作看是否能成功写入。如果成功,再逐步调整checkpoint配置,排查是否是checkpoint与S3交互的问题
内容的提问来源于stack exchange,提问作者he wang
相关产品推荐
相关产品推荐

