
大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载导读本文是 Flink 项目工程化配置的入门核心指南涵盖使用 Maven/Gradle 创建 Flink 项目、选择并声明正确的 API 与连接器依赖、打包 fat/uber JAR 提交作业以及测试依赖与高级配置Flink 依赖剖析、Scala 版本、Table 依赖、Hadoop 依赖等完整链路。读完本文你将能独立搭建一个可编译、可测试、可打包、可提交运行的标准 Flink 作业工程并理解 Flink 发行版中/lib、/opt目录与类加载机制对依赖声明方式的深层影响。开始创建 Flink 项目每个 Flink 应用程序都依赖于一组 Flink 库。最低限度下应用依赖 Flink API此外通常还需要某些连接器库如 Kafka、Cassandra以及用户自定义数据处理逻辑所需的第三方依赖。本节展示如何用 Maven 与 Gradle 快速创建项目骨架。提示所有 Flink Scala API 均已弃用deprecated并将在未来的 Flink 版本中移除。你仍然可以用 Scala 构建应用但建议迁移到 DataStream 和/或 Table API 的 Java 版本参见 FLIP-265。使用 Maven 创建项目Maven 用户可通过官方 archetype 原型 创建项目$ mvn archetype:generate \ -DarchetypeGroupIdorg.apache.flink \ -DarchetypeArtifactIdflink-quickstart-java \ -DarchetypeVersionflink-version执行后Maven 会交互式地询问 groupId、artifactId 与 package 名称。该原型对应的源码位于仓库 flink-quickstart/flink-quickstart-java其pom.xml使用maven-archetype打包方式pom.xml即 Flink 官方脚手架即由该模块生成。也可以直接使用快速启动脚本一步到位$ curl https://flink.apache.org/q/quickstart.sh | bash -s flink-version使用 Gradle 创建项目Gradle 用户可创建一个空项目手动创建src/main/java与src/main/resources目录后开始编写类或直接使用下面的构建脚本与快速启动脚本获得功能完整的启动工程。请在脚本所在目录执行gradle命令。build.gradle官方模板可直接复制使用plugins { id java id application // shadow plugin to produce fat JARs id com.github.johnrengelman.shadow version 7.1.2 } // artifact properties group org.quickstart version 0.1-SNAPSHOT mainClassName org.quickstart.DataStreamJob description Flink Quickstart Job ext { javaVersion 1.8 flinkVersion flink-version scalaBinaryVersion scala-binary-version slf4jVersion 1.7.36 log4jVersion 2.17.1 } sourceCompatibility javaVersion targetCompatibility javaVersion tasks.withType(JavaCompile) { options.encoding UTF-8 } applicationDefaultJvmArgs [-Dlog4j.configurationFilelog4j2.properties] // declare where to find the dependencies of your project repositories { mavenCentral() maven { url https://repository.apache.org/content/repositories/snapshots mavenContent { snapshotsOnly() } } } // NOTE: We cannot use compileOnly or shadow configurations since then we could not run code // in the IDE or with gradle run. We also cannot exclude transitive dependencies from the // shadowJar yet. // - Explicitly define the libraries we want to be included in the flinkShadowJar configuration! configurations { flinkShadowJar // dependencies which go into the shadowJar // always exclude these (also from transitive dependencies) since they are provided by Flink flinkShadowJar.exclude group: org.apache.flink, module: force-shading flinkShadowJar.exclude group: com.google.code.findbugs, module: jsr305 flinkShadowJar.exclude group: org.slf4j flinkShadowJar.exclude group: org.apache.logging.log4j } // declare the dependencies for your production and test code dependencies { // -------------------------------------------------------------- // Compile-time dependencies that should NOT be part of the // shadow (uber) jar and are provided in the lib folder of Flink // -------------------------------------------------------------- implementation org.apache.flink:flink-streaming-java:${flinkVersion} implementation org.apache.flink:flink-clients:${flinkVersion} // -------------------------------------------------------------- // Dependencies that should be part of the shadow jar, e.g. // connectors. These must be in the flinkShadowJar configuration! // -------------------------------------------------------------- //flinkShadowJar org.apache.flink:flink-connector-kafka:${flinkVersion} runtimeOnly org.apache.logging.log4j:log4j-slf4j-impl:${log4jVersion} runtimeOnly org.apache.logging.log4j:log4j-api:${log4jVersion} runtimeOnly org.apache.logging.log4j:log4j-core:${log4jVersion} // Add test dependencies here. // testCompile junit:junit:4.12 } // make compileOnly dependencies available for tests: sourceSets { main.compileClasspath configurations.flinkShadowJar main.runtimeClasspath configurations.flinkShadowJar test.compileClasspath configurations.flinkShadowJar test.runtimeClasspath configurations.flinkShadowJar javadoc.classpath configurations.flinkShadowJar } run.classpath sourceSets.main.runtimeClasspath jar { manifest { attributes Built-By: System.getProperty(user.name), Build-Jdk: System.getProperty(java.version) } } shadowJar { configurations [project.configurations.flinkShadowJar] }settings.gradlerootProject.name quickstart也可使用 Gradle 快速启动脚本需要传递 Flink 版本与 Scala 二进制版本两个参数bash -c $(curl https://flink.apache.org/q/gradle-quickstart.sh) -- flink-version scala-binary-version需要哪些依赖项要开始一个 Flink 作业通常需要如下依赖Flink API用于开发作业连接器和格式用于将作业与外部系统集成测试实用程序用于测试作业。除此之外开发自定义功能时还需按需添加第三方依赖。Flink API 依赖对照Flink 提供两大 API——DataStream API 与 Table API SQL二者可单独使用也可混合使用。按需选择并加入构建描述符即可你要使用的 API需要添加的依赖项DataStreamflink-streaming-javaDataStream Scala 版flink-streaming-scala_scala-binary-versionTable APIflink-table-api-javaTable API Scala 版flink-table-api-scala_scala-binary-versionTable API DataStreamflink-table-api-java-bridgeTable API DataStream Scala 版flink-table-api-scala-bridge_scala-binary-version从仓库的模块结构可以看到这些 API 模块的真实布局例如flink-datastream-api、flink-streaming-java、flink-table-api-java、flink-table-api-java-bridge均以独立模块形式存在于 flink-streaming-java 与 flink-table 下。Scala 相关的 API 模块则带_2.12等二进制版本后缀这正是下一节Scala 版本所强调的二进制不兼容问题的来源。运行和打包理解 fat/uber JAR如果你想通过简单地执行主类来运行作业classpath 中需要包含flink-clients对于 Table API 程序还需要flink-table-runtime与flink-table-planner-loader。经验法则建议将应用程序代码及其全部必需依赖连接器、格式、第三方依赖打包进一个 fat/uber JAR。此规则不适用于 Java API、DataStream Scala API 以及上述运行时模块——它们已由 Flink 发行版自带不应打入作业 uber JAR。打包后的作业 JAR 可以直接提交给已运行的 Flink 集群或轻松加入 Flink 应用容器镜像而无需改动发行版。当前仓库的版本为2.0-SNAPSHOT见根 pom.xml在实际使用上述命令与坐标时请将flink-version替换为你使用的具体发布版本号。连接器与格式的依赖管理Flink 应用通过连接器读写各种外部系统并通过格式对数据进行编解码。Flink 社区为每个连接器在 Maven Central 发布两类组件详细指南flink-connector-NAME精简 JAR只包含连接器自身代码不含最终第三方依赖flink-sql-connector-NAME包含连接器及其第三方依赖的 uber JAR。格式组件同理。部分连接器如文件系统连接器因不需要第三方依赖而没有对应的flink-sql-connector-NAME组件。仓库中可以看到典型示例flink-formats下同时存在flink-json与flink-sql-json两个模块flink-connectors下的连接器模块也遵循这一命名约定。三种使用方式使用连接器/格式模块可选把精简 JAR 及其传递依赖打包进作业 JAR把 uber JAR 打包进作业 JAR把 uber JAR 直接复制到 Flink 发行版的/lib文件夹。三者取舍使用uber JAR对作业里的依赖版本有更多控制权使用精简 JAR可在不更换连接器版本的情况下单独升级传递依赖需保证二进制兼容直接内嵌到发行版/lib可在一处统一控制所有作业的连接器版本。uber/fat JAR 主要配合 SQL 客户端 使用但同样可用于任何 DataStream/Table 应用程序。使用 Maven 配置项目Maven 是 Apache 软件基金会开源的自动化构建工具可用于构建、发布与部署项目。要求Maven 3.8.6Java 8已弃用或 Java 11。导入 IDE创建好项目后建议导入 IDE 开发测试。IntelliJ IDEA 开箱即用地支持 Maven 项目Eclipse 通过 m2e 插件导入。堆内存Java 默认 JVM 堆大小对 Flink 而言可能过小建议手动调大。Eclipse 中在Run Configurations - Arguments的VM Arguments填-Xmx800mIntelliJ IDEA 推荐通过Help | Edit Custom VM Options修改 JVM 属性。IntelliJ 注意运行配置需勾选Include dependencies with Provided scope若选项不可用较旧版本 IDEA可创建调用应用main()方法的测试用例来运行。构建与添加依赖在项目目录执行mvn clean package会在target/artifact-id-version.jar生成包含应用及连接器依赖的 JAR。若主类不是DataStreamJob建议同步修改pom.xml中的mainClassName使 Flink 可直接通过 JAR 运行而无需额外指定主类。在pom.xml的dependencies内添加依赖例如 Kafka 连接器dependencies dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-kafka/artifactId versionflink-version/version /dependency /dependencies随后执行mvn install。由官方模板创建的项目mvn clean package会自动将依赖打入应用 JAR非模板项目则建议使用 Maven Shade 插件。重要提示所有核心 Flink 依赖的生效范围应设为provided——即编译期可用但不打包进应用 JAR。否则最好的情况是 JAR 因包含全部 Flink 核心依赖而过大最坏情况是打入的 Flink 核心依赖与你的其他依赖发生版本冲突通常通过反向类加载规避。而需要打进应用 JAR 的依赖如连接器生效范围必须为compile。打包与提交仅使用 Flink 自带依赖如 JSON 格式的文件系统连接器时无需构建 uber/fat JAR使用发行版未内置的外部依赖时可将其加入发行版类路径或打入 uber/fat 应用 JAR。提交 uber/fat JAR 到本地或远程集群bin/flink run -c org.example.MyJob myFatJar.jar更多部署细节见 部署 CLI 指南。maven-shade-plugin 模板构建包含全部连接器与库依赖的应用 JAR可使用如下 shade 插件定义build plugins plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-shade-plugin/artifactId version3.1.1/version executions execution phasepackage/phase goals goalshade/goal /goals configuration artifactSet excludes excludecom.google.code.findbugs:jsr305/exclude /excludes /artifactSet filters filter !-- Do not copy the signatures in the META-INF folder. Otherwise, this might cause SecurityExceptions when using the JAR. -- artifact*:*/artifact excludes excludeMETA-INF/*.SF/exclude excludeMETA-INF/*.DSA/exclude excludeMETA-INF/*.RSA/exclude /excludes /filter /filters transformers transformer implementationorg.apache.maven.plugins.shade.resource.ManifestResourceTransformer !-- Replace this with the main class of your job -- mainClassmy.programs.main.clazz/mainClass /transformer transformer implementationorg.apache.maven.plugins.shade.resource.ServicesResourceTransformer/ /transformers /configuration /execution /executions /plugin /plugins /build注意两个细节剔除META-INF/*.SF、*.DSA、*.RSA签名文件可避免使用 JAR 时触发 SecurityExceptionServicesResourceTransformer用于合并 META-INF/services 下的服务描述文件。Maven Shade 插件默认包含生效范围为runtime或compile的全部依赖。使用 Gradle 配置项目Gradle 是开源通用构建工具可用于自动化开发任务。要求Gradle 7.xJava 8已弃用或 Java 11。导入 IDEIntelliJ IDEA 通过 Gradle 插件支持 Gradle 项目Eclipse 通过 Eclipse Buildship 插件导入导入向导最后一步需指定 Gradle 版本 ≥ 3.0shadow 插件会用到。JVM 堆与Include dependencies with Provided scope的注意事项与 Maven 章节相同。构建与添加依赖执行gradle clean shadowJar会在build/libs/project-name-version-all.jar生成包含应用及连接器依赖的 JAR。若主类不是StreamingJob请同步修改build.gradle中的mainClassName。在build.gradle的 dependencies 块中配置依赖例如使用模板工程时添加 Kafka 连接器dependencies { ... flinkShadowJar org.apache.flink:flink-connector-kafka:${flinkVersion} ... }注意需要打入 fat JAR 的连接器依赖必须放在flinkShadowJar配置中而由 Flink 发行版提供的核心依赖flink-streaming-java、flink-clients放在implementation中详见上文的完整模板。依赖生效范围的原则与 Maven 章节一致核心依赖应为 provided应用依赖应为 compile。打包与提交无外部依赖时gradle clean installDistGradle Wrapper 则用./gradlew clean installDist有外部依赖时gradle clean installShadowDist在/build/install/yourProject/lib生成 fat JARWrapper 则用./gradlew clean installShadowDist也可将依赖加入发行版类路径。提交 uber/fat JARbin/flink run -c org.example.MyJob myFatJar.jar用于测试的依赖项Flink 提供测试作业的实用程序可直接添加为依赖测试依赖指南。DataStream API 测试为 DataStream 作业开发测试用例需添加flink-test-utils依赖test生效范围。该模块提供了MiniCluster——一个可配置的轻量级 Flink 集群可在 JUnit 测试中运行并直接执行作业。从源码看仓库中 MiniClusterWithClientResource 封装了运行时测试资源其内部使用MiniClusterClient与MiniClusterResourceConfiguration是编写集成测试的常用入口。Table API 测试要在 IDE 本地测试 Table API/SQL 程序除flink-test-utils外还需添加flink-table-test-utilstest生效范围。它位于 flink-table/flink-table-test-utils会自动引入查询计划器与运行时分别用于查询的计划与执行。该模块自 Flink 1.15 引入目前视为实验性模块。高级配置主题Flink 依赖剖析Flink 自身由一组核心类与依赖构成运行时核心在应用启动时必须存在提供通信协调、网络管理、检查点、容错、API、算子如窗口、资源管理等服务。这些核心依赖打包在flink-dist.jar位于发行版/lib目录也是 Flink 容器镜像的基础可近似理解为 Java 核心库之于 JVM。为保持核心依赖精简并避免冲突Flink Core Dependencies不包含任何连接器或库如 CEP、SQL、ML。发行版/lib还包含常用模块如执行 Table 作业的必需模块、一组连接器和格式默认自动加载若需禁止加载直接从 classpath 的/lib删除对应 JAR 即可。/opt目录下的额外可选依赖通过移动 JAR 到/lib来启用。类加载机制细节见 Flink 类加载。Scala 版本不同 Scala 版本的二进制不兼容所有传递地依赖 Scala 的 Flink 依赖项都以构建时的 Scala 版本为后缀如flink-streaming-scala_2.12。仅使用 Java API 时可用任意 Scala 版本使用 Scala API 时则需选择与应用匹配的 Scala 版本。注意2.12.8 之后的 Scala 版本与之前 2.12.x 二进制不兼容本地为更高 Scala 版本构建 Flink 时需添加-Djapicmp.skip跳过二进制兼容性检查。Table 依赖剖析Flink 发行版默认在/lib包含执行 SQL 任务所需的三个 JARflink-table-api-java-uber-version.jar包含全部 Java APIflink-table-runtime-version.jarTable 运行时flink-table-planner-loader-version.jar查询计划器。自 Flink 1.15 起原先打包为一个flink-table.jar的依赖被拆分为上述三个 JAR允许以flink-table-planner-loader充当内部flink-table-planner。Table Java API 内置于发行版但默认不包含 Table Scala API使用 Scala API 的格式与连接器时需手动下载 JAR 放入/lib推荐或打入 uber/fat JAR。Table Planner 与 Planner Loader自 Flink 1.15 起发行版包含两个 plannerflink-table-planner_scala-version-version.jar位于/opt包含查询计划器需使用与发行版相同版本的 Scalaflink-table-planner-loader-version.jar位于/lib默认加载计划器被隐藏在独立 classpath 中无法直接使用io.apache.flink.table.planner包因 Scala 已打包其中无需考虑 Scala 版本。两个 JAR 代码功能相同、打包方式不同。切勿将二者同时放入 classpath否则 Table 任务将失败。Flink 计划停止在发行版中发布flink-table-planner_scala-version组件强烈建议迁移作业/自定义连接器/格式以使用前述 API 模块而不依赖内部 planner。Hadoop 依赖一般规则无需直接向应用添加 Hadoop 依赖。要与 Hadoop 集成应让 Flink 系统本身携带 Hadoop 依赖而非将其作为用户代码依赖。通过环境变量指定export HADOOP_CLASSPATHhadoop classpath设计原因有二其一部分 Hadoop 交互发生在用户应用启动之前如为检查点配置 HDFS、通过 Hadoop Kerberos 令牌认证、在 YARN 上部署其二Flink 的反向类加载方式在核心依赖中隐藏了许多传递依赖含 Hadoop 依赖使应用可使用相同依赖的不同版本而不冲突这对大型依赖树尤为重要。若仅在 IDE 开发/测试时需要 Hadoop 依赖如 HDFS 访问应限定其生效范围为test或provided。下一步是什么开始开发作业查看 DataStream API 与 Table API SQL按构建工具深入打包细节Maven 指南、Gradle 指南更多项目配置高级主题高级配置。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Flink CDC 3.6 DataStream API 作业打包指南基于 MySQL CDC 的 Maven 工程搭建与 JAR 发布实战Flink CDC 3.6 DataStream API 作业打包指南基于 MySQL CDC 的 Maven 工程搭建与 JAR 发布实战 本篇指南围绕 F后端数据集成大数据流处理变更数据捕获数据同步gl-rs未来路线图OpenGL 4.6支持与Vulkan集成gl rs未来路线图OpenGL 4.6支持与Vulkan集成 引言gl rs项目概述 gl rs是一个专为Rust编程语言设计的 OpenGL函数大数据流处理批处理数据工程Apache Flink CDC PostgreSQL Connector 全指南从建表配置到增量快照与 DataStream 实战Apache Flink CDC PostgreSQL Connector 全指南从建表配置到增量快照与 DataStream 实战 PostgreSQL C后端数据集成大数据流处理变更数据捕获数据同步上一篇OP-TEE OS深入解析可信执行环境的安全基石 - 终极指南下一篇Molemac 清理工具完整指南快速分析磁盘空间、安全清理缓存、卸载应用与监控系统创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考