Apache Flink 1.11.2:flink-s3-fs-hadoop插件位置及IDE使用咨询
Hey,我来帮你梳理下Flink 1.11.2对接S3的这些问题,都是实际开发中常遇到的点:
生产环境获取
flink-s3-fs-hadoop-1.11.2.jar的路径 你有几个可靠的途径获取这个插件jar:
- Flink官方二进制包:下载Flink 1.11.2的官方压缩包(比如
flink-1.11.2-bin-scala_2.11.tgz),解压后进入opt目录,里面应该包含flink-s3-fs-hadoop-1.11.2.jar。如果你是通过包管理器(比如yum、apt)安装的Flink,可能需要单独下载对应版本的插件包。 - Maven中央仓库:直接搜索
flink-s3-fs-hadoop的1.11.2版本,下载对应的jar文件。这是最快捷的方式,不用下载整个Flink包。 - 源码编译(不推荐生产用):如果上述途径都不行,可以克隆Flink 1.11.2的源码,编译
s3-fs-hadoop模块,但预编译的官方jar更稳定,适合生产环境。
拿到jar后,把它复制到Flink集群的plugin/s3-fs-hadoop目录(需要手动创建这个子目录),重启Flink节点即可生效。
IDE中使用这些插件:依赖配置
没错,只需在pom.xml中添加对应依赖即可,但要注意scope的设置:
- 如果你需要在IDE中本地运行测试Flink作业,把scope设为
compile,这样依赖会被加载到classpath; - 如果只是打包作业提交到已配置好插件的集群,把scope设为
provided,避免依赖冲突(因为集群的plugin目录已经有这些jar了)。
示例依赖配置:
<!-- S3 Presto 插件(用于Checkpoint) --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-s3-fs-presto</artifactId> <version>1.11.2</version> <scope>compile</scope> <!-- 本地测试用compile,集群提交换为provided --> </dependency> <!-- S3 Hadoop 插件(用于管道数据处理) --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-s3-fs-hadoop</artifactId> <version>1.11.2</version> <scope>compile</scope> </dependency>
IDE中传递S3凭证的几种方式
根据测试场景,你可以选择以下几种方式:
- IDE运行配置添加系统属性:在IDE的Run/Debug Configuration中,添加VM参数:
如果用的是Presto插件,参数是-Dfs.s3a.access.key=你的访问密钥ID -Dfs.s3a.secret.key=你的私有访问密钥-Ds3.access-key和-Ds3.secret-key。 - 配置文件方式:在项目的
src/main/resources目录下创建core-site.xml,添加S3凭证配置:<configuration> <property> <name>fs.s3a.access.key</name> <value>你的访问密钥ID</value> </property> <property> <name>fs.s3a.secret.key</name> <value>你的私有访问密钥</value> </property> </configuration> - 代码中显式配置(仅测试用):不推荐生产环境硬编码凭证,但本地测试可以临时用:
import org.apache.flink.configuration.Configuration; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; public class S3TestJob { public static void main(String[] args) throws Exception { Configuration conf = new Configuration(); // Hadoop插件的配置 conf.setString("fs.s3a.access.key", "你的访问密钥ID"); conf.setString("fs.s3a.secret.key", "你的私有访问密钥"); // Presto插件的配置(如果用它做Checkpoint) conf.setString("s3.access-key", "你的访问密钥ID"); conf.setString("s3.secret-key", "你的私有访问密钥"); StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(conf); // 后续作业逻辑... } } - AWS CLI凭证文件:如果你本地配置了AWS CLI(
~/.aws/credentials),Flink会自动读取这个文件中的凭证,无需额外配置,这是最安全的方式之一。
内容的提问来源于stack exchange,提问作者JanOels
相关产品推荐
相关产品推荐

