Flink 连接器与格式依赖配置完全指南:精简 JAR 与 uber JAR 的选择与实践
【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink
Flink 应用需要通过连接器(Connector)读写 Kafka、JDBC、Filesystem 等外部系统,并通过格式(Format)完成数据编解码。本文围绕 连接器和格式 这一官方配置文档,系统讲解flink-connector-<NAME>精简 JAR 与flink-sql-connector-<NAME>uber JAR 两类组件的区别、引入方式与取舍策略,并结合本仓库的模块结构与构建配置给出可落地的 Maven 实操方案。
连接器和格式:Flink 与外部系统的桥梁
Flink 应用程序通过连接器读写各种外部系统,并通过格式对数据进行编码与解码,使其匹配 Flink 内部的数据结构。这一能力对两大 API 同时开放:
- DataStream API:参见 DataStream Connectors 概览,社区随 Flink 工程一起维护的连接器包括 Apache Kafka(source/sink)、Apache Cassandra(source/sink)、Amazon DynamoDB(sink)、Amazon Kinesis Data Streams(source/sink)、Amazon Kinesis Data Firehose(sink)、DataGen(source)、Elasticsearch(sink)、Opensearch(sink)、FileSystem(sink)、RabbitMQ(source/sink)、Google PubSub(source/sink)、Hybrid Source(source)、Apache Pulsar(source)、JDBC(sink)、MongoDB(source/sink)等;
- Table API/SQL:参见 Table & SQL Connectors 概览,原生支持 Filesystem、Elasticsearch、Opensearch、Apache Kafka、Amazon DynamoDB、Amazon Kinesis Data Streams、Amazon Kinesis Data Firehose、JDBC、Apache HBase、Apache Hive、MongoDB 等,并通过
CREATE TABLE ... WITH (...)语句声明连接目标与对应格式。
值得注意的是,上述连接器是 Flink 工程的一部分、包含在发布的源码中,但并不包含在二进制发行版中,需要开发者自行把对应组件引入作业。这正是本文要解决的配置问题。
可用的组件:精简 JAR 与 uber JAR
为了让 Flink 能够访问实现连接器和格式功能的组件,对于 Flink 社区支持的每个连接器,官方在 Maven Central 上发布了两类组件:
| 组件类型 | 命名模式 | 内容 | 典型场景 |
|---|---|---|---|
| 精简 JAR | flink-connector-<NAME> | 仅包含连接器代码本身,不包含最终第三方依赖项 | DataStream/Table 作业中按需引入,配合构建工具解析传递依赖 |
| uber JAR | flink-sql-connector-<NAME> | 连接器代码 + 其全部第三方依赖项的 fat/uber JAR | 主要配合 SQL 客户端使用,也可用于任何 DataStream/Table 应用 |
格式(Format)组件同样适用这一规律,例如flink-json、flink-csv、flink-avro、flink-parquet、flink-orc与对应的flink-sql-*组件。另外请注意:某些连接器没有对应的flink-sql-connector-<NAME>组件,因为它们本身不依赖第三方依赖项,无需打 uber JAR。
本仓库 docs/data/sql_connectors.yml 正是这张组件清单的数据源——它为每个连接器/格式记录了name、category(format 或 connector)、maven(Maven 模块名)与sql_url(对应 uber JAR 的下载地址),并支持通过versions小节声明同一连接器的多版本。例如 Kafka 在清单中的maven为flink-connector-kafka,其 uber JAR 对应flink-sql-connector-kafka;而 JDBC 连接器则直接指向flink-connector-jdbc。下载页与文档中的组件表格均由该文件生成,可作为查询"某个连接器到底该用哪个组件"的第一手依据。
仓库中的 uber JAR 实例:flink-sql-connector-hive
以本仓库 flink-connectors/flink-sql-connector-hive-3.1.3/pom.xml 为例,可以直观看到 uber JAR 的组装方式:它依赖精简连接器模块flink-connector-hive_${scala.binary.version},同时直接引入第三方依赖hive-exec:3.1.3,并显式排除log4j、slf4j-log4j12、guava、avro、reload4j等可能与 Flink 运行时冲突的传递依赖,最终由 shade 插件打成包含 Hive 全部必要依赖的单一 JAR。这与文档中"flink-sql-connector-<NAME>是包含连接器第三方依赖项的 uber JAR"的描述完全吻合。
使用组件:三种引入方式
要把连接器/格式模块引入运行环境,官方文档给出了三种方式:
- 把精简 JAR 及其传递依赖项打包进您的作业 JAR:适用于想自行控制依赖解析的情况,用 Maven 或 Gradle 在构建期完成打包;
- 把 uber JAR 打包进您的作业 JAR:直接以单一依赖形式把
flink-sql-connector-<NAME>塞进作业包,省去处理传递依赖的麻烦; - 把 uber JAR 直接复制到 Flink 发行版的
/lib文件夹内:对集群全局生效,所有提交到该发行版的作业共享同一份连接器。
其中关于打包依赖项的具体操作,请参考 Maven 指南 与 Gradle 指南;关于 Flink 发行版中哪些依赖由运行时提供、哪些需要作业自带的边界,请参考 Flink 依赖剖析 一节——Java API、DataStream Scala API 及运行时模块已由 Flink 本身提供,不应打进作业 uber JAR,而连接器、格式与自定义第三方依赖则应打包进作业 JAR。
Maven 实操:添加一个连接器依赖
以 DataStream/Table 作业引入 Kafka 连接器为例,在项目pom.xml的<dependencies>内添加:
<dependencies> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka</artifactId> <version><!-- 替换为你使用的 Flink 版本 --></version> </dependency> </dependencies>然后执行mvn install(或mvn clean package)即可完成依赖解析。若使用 SQL 客户端或在发行版中全局生效,则改为引入flink-sql-connector-kafka并把生成的 JAR 放入 Flink 的/lib目录。
需要特别强调的是依赖生效范围(scope):Flink 核心依赖应设置为provided(编译需要、但不应打进作业 JAR,避免 JAR 膨胀与版本冲突);而连接器、格式等作业真正需要的第三方依赖应保持compile/runtime范围,使其被正确打包进应用程序 JAR。对于非模板创建的 Maven 项目,建议使用maven-shade-plugin将所有必需依赖合并为 uber/fat JAR,详见 Maven 指南。
三种方式如何选:控制权与运维的权衡
选择 uber JAR、精简 JAR 还是直接内嵌到发行版/lib,取决于你的使用场景,官方文档给出的决策依据如下:
- 使用 uber JAR:将连接器及其第三方依赖作为一个整体打进作业,对作业里的依赖项版本有更多控制权——连接器与其传递依赖的版本绑定关系由你决定;
- 使用精简 JAR:只把连接器代码打入作业,第三方依赖由构建工具按坐标单独解析。由于可以在不更换连接器版本的情况下单独升级某个传递依赖(只要保持二进制兼容),对传递依赖项有更多控制权;
- 把 uber JAR 内嵌到 Flink 发行版的
/lib:连接器对发行版内所有作业全局可见,可以在一处控制所有作业的连接器版本,便于集中运维与升级,代价是不同作业之间无法独立选用不同版本。
简言之:追求作业内版本可控选 uber JAR,追求传递依赖精细管理选精简 JAR,追求集群级统一管控则使用/lib目录方案。
打包实战进阶:多组件 uber JAR 的 SPI 合并问题
当你的项目同时使用多个表连接器/格式并打成 uber JAR 时,会遇到一个隐蔽的坑。Flink 通过 Java 的 SPI(Service Provider Interface) 按工厂标识符(如kafka、json)加载 Table 连接器/格式工厂,而每个连接器/格式的 SPI 资源文件都叫META-INF/services/org.apache.flink.table.factories.Factory、位于同一目录下——在合并 uber JAR 时这些资源文件会互相覆盖,导致 Flink 无法加载工厂。
解决办法是在maven-shade-plugin中配置ServicesResourceTransformer,把META-INF/services下的资源文件合并而非覆盖。以下是一个同时使用flink-sql-connector-hive-3.1.3与flink-parquet的项目示例(该示例同样出现在 Table & SQL Connectors 概览 中):
<dependencies> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-sql-connector-hive-3.1.3_2.12</artifactId> <version><!-- 替换为你使用的 Flink 版本 --></version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-parquet</artifactId> <version><!-- 替换为你使用的 Flink 版本 --></version> </dependency> </dependencies> <build> <plugins> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-shade-plugin</artifactId> <executions> <execution> <id>shade</id> <phase>package</phase> <goals> <goal>shade</goal> </goals> <configuration> <transformers combine.children="append"> <!-- 合并 META-INF/services 文件,保证多个连接器/格式工厂可同时被加载 --> <transformer implementation="org.apache.maven.plugins.shade.resource.ServicesResourceTransformer"/> <!-- ... --> </transformers> </configuration> </execution> </executions> </plugin> </plugins> </build>配置完成后,META-INF/services下的连接器/格式资源文件在构建 uber JAR 时会按行合并,工厂才能被全部正确发现。如果你的作业 JAR 提交后出现"找不到 connector/format 工厂"的异常,优先检查这一步是否遗漏。
小结
连接器与格式的引入是 Flink 作业对接外部系统的第一步,核心决策集中在"精简 JAR vs uber JAR vs/lib内嵌"三选一:flink-connector-<NAME>提供最小依赖、flink-sql-connector-<NAME>提供开箱即用的 fat JAR,而复制到/lib则面向集群统一管控。构建 uber JAR 时还需留意 SPI 工厂文件合并,避免多组件共存场景下的加载失败。结合 Maven 指南、Gradle 指南 与 项目配置概览 一起阅读,即可为你的 Flink 作业搭建出稳定、可控的依赖体系。
【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考