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

本地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环境部署步骤

  1. 自定义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
  1. 配置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
  1. 修改Hudi表路径协议:建议将表路径从s3://改为s3a://,这是Hadoop官方推荐的S3协议,兼容性更强:
'path' = 's3a://path/to/hudi',

内容的提问来源于stack exchange,提问作者Rinze

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 03:27:29