Flink 连接器与格式依赖配置完全指南:精简 JAR 与 uber JAR 的选择与实践
Flink 连接器与格式依赖配置完全指南精简 JAR 与 uber JAR 的选择与实践【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flinkFlink 应用需要通过连接器Connector读写 Kafka、JDBC、Filesystem 等外部系统并通过格式Format完成数据编解码。本文围绕 连接器和格式 这一官方配置文档系统讲解flink-connector-NAME精简 JAR 与flink-sql-connector-NAMEuber JAR 两类组件的区别、引入方式与取舍策略并结合本仓库的模块结构与构建配置给出可落地的 Maven 实操方案。连接器和格式Flink 与外部系统的桥梁Flink 应用程序通过连接器读写各种外部系统并通过格式对数据进行编码与解码使其匹配 Flink 内部的数据结构。这一能力对两大 API 同时开放DataStream API参见 DataStream Connectors 概览社区随 Flink 工程一起维护的连接器包括 Apache Kafkasource/sink、Apache Cassandrasource/sink、Amazon DynamoDBsink、Amazon Kinesis Data Streamssource/sink、Amazon Kinesis Data Firehosesink、DataGensource、Elasticsearchsink、Opensearchsink、FileSystemsink、RabbitMQsource/sink、Google PubSubsource/sink、Hybrid Sourcesource、Apache Pulsarsource、JDBCsink、MongoDBsource/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 上发布了两类组件组件类型命名模式内容典型场景精简 JARflink-connector-NAME仅包含连接器代码本身不包含最终第三方依赖项DataStream/Table 作业中按需引入配合构建工具解析传递依赖uber JARflink-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、categoryformat 或 connector、mavenMaven 模块名与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 groupIdorg.apache.flink/groupId artifactIdflink-connector-kafka/artifactId version!-- 替换为你使用的 Flink 版本 --/version /dependency /dependencies然后执行mvn install或mvn clean package即可完成依赖解析。若使用 SQL 客户端或在发行版中全局生效则改为引入flink-sql-connector-kafka并把生成的 JAR 放入 Flink 的/lib目录。需要特别强调的是依赖生效范围scopeFlink 核心依赖应设置为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 的 SPIService 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 groupIdorg.apache.flink/groupId artifactIdflink-sql-connector-hive-3.1.3_2.12/artifactId version!-- 替换为你使用的 Flink 版本 --/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-parquet/artifactId version!-- 替换为你使用的 Flink 版本 --/version /dependency /dependencies build plugins plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-shade-plugin/artifactId executions execution idshade/id phasepackage/phase goals goalshade/goal /goals configuration transformers combine.childrenappend !-- 合并 META-INF/services 文件保证多个连接器/格式工厂可同时被加载 -- transformer implementationorg.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),仅供参考
上一篇/下一篇内容由系统自动关联
返回资讯列表 →