本地Docker部署Flink同步MySQL CDC到S3(Hudi)遇S3协议错误求依赖包
问题
本地通过Docker部署Flink,编写Flink任务使用MySQL CDC将数据同步至S3(以Apache Hudi格式存储),任务代码如下:
t_env = StreamTableEnvironment.create(env, environment_settings=settings) t_env.execute_sql(f""" CREATE TABLE source_table ( xxx ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = 'xxx', 'port' = 'xxx', 'username' = 'xxx', 'password' = 'xxx', 'database-name' = 'xxx', 'table-name' = 'xxx', 'scan.incremental.snapshot.enabled' = 'true', 'scan.incremental.snapshot.chunk.key-column' = 'id' ) """) t_env.execute_sql(f""" CREATE TABLE hudi_table ( xxx, PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'hudi', 'path' = 's3://path/to/hudi', 'table.type' = 'COPY_ON_WRITE', 'write.precombine.field' = 'id', 'write.operation' = 'upsert', 'hoodie.datasource.write.recordkey.field' = 'id', 'hoodie.datasource.write.partitionpath.field' = '', 'write.tasks' = '1', 'compaction.tasks' = '1', 'compaction.async.enabled' = 'true' ) """) t_env.execute_sql(f""" INSERT INTO hudi_table SELECT * FROM source_table """)
提交任务后出现如下错误:
org.apache.flink.table.api.ValidationException: Unable to create a sink for writing table 'default_catalog.default_database.hudi_table'. ... Caused by: java.io.IOException: No FileSystem for scheme: s3 at org.apache.hadoop.fs.FileSystem.getFileSystemClass(FileSystem.java:2799) at org.apache.hadoop.fs.FileSystem.createFileSystem(FileSystem.java:2810) at org.apache.hadoop.fs.FileSystem.access$200(FileSystem.java:100) at org.apache.hadoop.fs.FileSystem$Cache.getInternal(FileSystem.java:2849) at org.apache.hadoop.fs.FileSystem$Cache.get(FileSystem.java:2831) at org.apache.hadoop.fs.FileSystem.get(FileSystem.java:389) at org.apache.hadoop.fs.Path.getFileSystem(Path.java:356) at org.apache.hudi.hadoop.fs.HadoopFSUtils.getFs(HadoopFSUtils.java:126) ... 42 more
需补充哪些JAR包让Flink具备访问S3的能力?
解决方案
要解决S3文件系统未识别的问题,需要补充Hadoop兼容的S3文件系统依赖包,具体如下(版本需匹配你的Flink和Hudi版本):
hadoop-aws:提供S3文件系统的核心实现aws-java-sdk-bundle:整合AWS SDK相关依赖,包含S3客户端核心逻辑
Docker环境部署步骤
- 自定义Flink镜像:基于官方Flink镜像,将上述JAR包复制到Flink的
/opt/flink/lib目录,示例Dockerfile:
FROM flink:1.17.1-scala_2.12 # 下载对应版本的依赖包(示例版本适配Flink 1.17、Hudi 0.14) RUN wget -P /opt/flink/lib/ https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-aws/3.3.4/hadoop-aws-3.3.4.jar RUN wget -P /opt/flink/lib/ https://repo1.maven.org/maven2/com/amazonaws/aws-java-sdk-bundle/1.12.512/aws-java-sdk-bundle-1.12.512.jar
- 配置S3访问凭证:在
flink-conf.yaml中添加AWS凭证(也可通过环境变量传递):
fs.s3a.access.key: your-access-key fs.s3a.secret.key: your-secret-key # 若使用S3兼容存储(如MinIO),额外添加: # fs.s3a.endpoint: http://your-s3-endpoint # fs.s3a.path.style.access: true
- 修改Hudi表路径协议:建议将表路径从
s3://改为s3a://,这是Hadoop官方推荐的S3协议,兼容性更强:
'path' = 's3a://path/to/hudi',
内容的提问来源于stack exchange,提问作者Rinze
相关产品推荐
相关产品推荐

