如何在EMR环境的Flink应用中使用s3a文件系统?
EMR 6.1环境Flink 1.11.0 s3a文件系统适配解决方案
协议识别异常修复
- 首先删除已手动安装的Flink s3文件系统插件,避免依赖冲突:默认插件路径为
$FLINK_HOME/plugins/s3-fs-hadoop、$FLINK_HOME/plugins/s3-fs-presto,直接删除对应目录即可。 - 提交Flink作业前先执行环境变量配置命令,加载EMR预安装的Hadoop运行时依赖与配置:
完成该配置后,s3://、s3a://两种协议均可被Flink正常识别,无需额外引入s3相关依赖包。export HADOOP_CLASSPATH=`hadoop classpath`
自定义s3凭证提供器配置生效方法
- 全局生效配置:在
$FLINK_HOME/conf/flink-conf.yaml中添加带统一前缀的配置项:# 配置自定义凭证提供器全类名 fs.s3a.aws.credentials.provider: <自定义凭证提供器全类名> # 自定义提供器所需的其他参数均添加fs.s3a.前缀配置即可 fs.s3a.<自定义参数名>: <参数值> - 单表独立配置(Table API场景):在CREATE TABLE的WITH参数中直接添加配置,优先级高于全局配置:
CREATE TABLE s3_test_table ( -- 表字段定义 ) WITH ( 'connector' = 'filesystem', 'path' = 's3a://<bucket名称>/<存储路径>', 'format' = 'csv', 's3a.aws.credentials.provider' = '<自定义凭证提供器全类名>', 's3a.<自定义参数名>' = '<参数值>' ) - 注意:自定义凭证提供器类需要打包进作业Jar包,或者放入
$FLINK_HOME/lib目录下,保证运行时可被加载。
内容的提问来源于stack exchange,提问作者Ryan Whitcomb
相关产品推荐
相关产品推荐

