You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

复现步骤

  1. 创建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'
);
  1. 插入数据到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>
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-hudi bucket的读写、创建对象/目录权限
  • 确认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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.11 18:45:31