Flink环境与部署-元一软件
flink是一款开源的大数据流式处理框架他可以同时批处理和流处理具有容错性、高吞吐、低延迟等优势本文简述flink在windows和linux中安装步骤和示例程序的运行包括本地调试环境集群环境。另外介绍Flink的开发工程的构建。首先要想运行Flink我们需要下载并解压Flink的二进制包下载地址如下https://flink.apache.org/down…我们可以选择Flink与Scala结合版本这里我们选择最新的1.9版本Apache Flink 1.9.0 for Scala 2.12进行下载。下载成功后在windows系统中可以通过Windows的bat文件或者Cygwin来运行Flink。在linux系统中分为单机集群和Hadoop等多种情况。通过Windows的bat文件运行首先启动cmd命令行窗口进入flink文件夹运行bin目录下的start-cluster.bat注意运行flink需要java环境请确保系统已经配置java环境变量。$cdflink $cdbin $ start-cluster.bat Starting alocalcluster with one JobManager process and one TaskManager process. You can terminate the processes via CTRL-Cinthe spawned shell windows. Web interface by default on http://localhost:8081/.显示启动成功后我们在浏览器访问 http://localhost:8081/可以看到flink的管理页面。通过Cygwin运行Cygwin是一个在windows平台上运行的类UNIX模拟环境官网下载http://cygwin.com/install.html安装成功后启动Cygwin终端运行start-cluster.sh脚本。$cdflink $ bin/start-cluster.sh Starting cluster.显示启动成功后我们在浏览器访问 http://localhost:8081/可以看到flink的管理页面。Linux系统上安装flink单节点安装在Linux上单节点安装方式与cygwin一样下载Apache Flink 1.9.0 for Scala 2.12然后解压后只需要启动start-cluster.sh。集群安装集群安装分为以下几步1、在每台机器上复制解压出来的flink目录。2、选择一个作为master节点然后修改所有机器conf/flink-conf.yamljobmanager.rpc.addressmaster主机名3、修改conf/slaves,将所有work节点写入work01 work024、在master上启动集群bin/start-cluster.sh安装在Hadoop我们可以选择让Flink运行在Yarn集群上。下载Flink for Hadoop的包保证 HADOOP_HOME已经正确设置即可启动 bin/yarn-session.sh运行flink示例程序批处理示例提交flink的批处理examples程序bin/flink run examples/batch/WordCount.jar这是flink提供的examples下的批处理例子程序统计单词个数。$ bin/flink run examples/batch/WordCount.jar Starting execution of program Executing WordCount example with default input data set. Use--inputto specifyfileinput. Printing result to stdout. Use--outputto specify output path.(a,5)(action,1)(after,1)(against,1)(all,2)(and,12)(arms,1)(arrows,1)(awry,1)(ay,1)得到结果这里统计的是默认的数据集可以通过–input --output指定输入输出。我们可以在页面中查看运行的情况流处理示例启动nc服务器nc-l9000提交flink的批处理examples程序bin/flink run examples/streaming/SocketWindowWordCount.jar--port9000这是flink提供的examples下的流处理例子程序接收socket数据传入统计单词个数。在nc端写入单词$nc-l9000lorem ipsum ipsum ipsum ipsum bye输出在日志中$tail-flog/flink-*-taskexecutor-*.out lorem:1bye:1ipsum:4停止flink$ ./bin/stop-cluster.sh在安装好Flink以后只要快速构建Flink工程并完成相关代码开发就可以轻松入手Flink。构建工具Flink项目可以使用不同的构建工具进行构建。为了能够快速入门Flink 为以下构建工具提供了项目模版MavenGradle这些模版可以帮助你搭建项目结构并创建初始构建文件。Maven环境要求唯一的要求是使用Maven 3.0.4或更高版本和安装Java 8.x。创建项目使用以下命令之一来创建项目使用Maven archetypes$ mvn archetype:generate\-DarchetypeGroupIdorg.apache.flink\-DarchetypeArtifactIdflink-quickstart-java\-DarchetypeVersion1.9.0运行quickstart脚本curlhttps://flink.apache.org/q/quickstart.sh|bash-s1.9.0下载完成后查看项目目录结构tree quickstart/ quickstart/ ├── pom.xml └── src └── main ├──java│ └── org │ └── myorg │ └── quickstart │ ├── BatchJob.java │ └── StreamingJob.java └── resources └── log4j.properties示例项目是一个Maven project它包含了两个类StreamingJob和BatchJob分别是DataStreamandDataSet程序的基础骨架程序。main方法是程序的入口既可用于IDE测试/执行也可用于部署。我们建议你将此项目导入IDE来开发和测试它。IntelliJ IDEA 支持 Maven 项目开箱即用。如果你使用的是 Eclipse使用m2e 插件 可以导入 Maven 项目。一些 Eclipse 捆绑包默认包含该插件其他情况需要你手动安装。请注意对 Flink 来说默认的 JVM 堆内存可能太小你应当手动增加堆内存。在 Eclipse 中选择Run Configurations - Arguments并在VM Arguments对应的输入框中写入-Xmx800m。在 IntelliJ IDEA 中推荐从菜单Help | Edit Custom VM Options来修改 JVM 选项。构建项目如果你想要构建/打包你的项目请在项目目录下运行 ‘mvn clean package’ 命令。命令执行后你将找到一个JAR文件里面包含了你的应用程序以及已作为依赖项添加到应用程序的连接器和库target/-.jar。注意如果你使用其他类而不是StreamingJob作为应用程序的主类/入口我们建议你相应地修改pom.xml文件中的mainClass配置。这样Flink 可以从 JAR 文件运行应用程序而无需另外指定主类。Gradle环境要求唯一的要求是使用Gradle 3.x(或更高版本) 和安装Java 8.x。创建项目使用以下命令之一来创建项目Gradle示例build.gradlebuildscript{repositories{jcenter()// this applies only to the GradleShadowplugin}dependencies{classpathcom.github.jengelman.gradle.plugins:shadow:2.0.4}}plugins{idjavaidapplication// shadow plugin to produce fat JARsidcom.github.johnrengelman.shadowversion2.0.4}// artifact properties grouporg.myorg.quickstartversion0.1-SNAPSHOTmainClassNameorg.myorg.quickstart.StreamingJobdescriptionFlink Quickstart Job ext{javaVersion1.8flinkVersion1.9.0scalaBinaryVersion2.11slf4jVersion1.7.7log4jVersion1.2.17}sourceCompatibilityjavaVersion targetCompatibilityjavaVersion tasks.withType(JavaCompile){options.encodingUTF-8}applicationDefaultJvmArgs[-Dlog4j.configurationlog4j.properties]task wrapper(type: Wrapper){gradleVersion3.1}//declarewhere tofindthe dependencies of your project repositories{mavenCentral()maven{urlhttps://repository.apache.org/content/repositories/snapshots/}}// 注意我们不能使用compileOnly或者shadow配置这会使我们无法在 IDE 中或通过使用gradle run命令运行代码。 // 我们也不能从 shadowJar 中排除传递依赖请查看 https://github.com/johnrengelman/shadow/issues/159)。 // -显式定义我们想要包含在flinkShadowJar配置中的类库!configurations{flinkShadowJar // dependencieswhichgo into the shadowJar // 总是排除这些依赖也来自传递依赖因为 Flink 会提供这些依赖。 flinkShadowJar.exclude group:org.apache.flink, module:force-shadingflinkShadowJar.exclude group:com.google.code.findbugs, module:jsr305flinkShadowJar.exclude group:org.slf4jflinkShadowJar.exclude group:log4j}//declarethe dependenciesforyour production andtestcode dependencies{// -------------------------------------------------------------- // 编译时依赖不应该包含在 shadow jar 中 // 这些依赖会在 Flink 的 lib 目录中提供。 // -------------------------------------------------------------- compileorg.apache.flink:flink-java:${flinkVersion}compileorg.apache.flink:flink-streaming-java_${scalaBinaryVersion}:${flinkVersion}// -------------------------------------------------------------- // 应该包含在 shadow jar 中的依赖例如连接器。 // 它们必须在 flinkShadowJar 的配置中 // -------------------------------------------------------------- //flinkShadowJarorg.apache.flink:flink-connector-kafka-0.11_${scalaBinaryVersion}:${flinkVersion}compilelog4j:log4j:${log4jVersion}compileorg.slf4j:slf4j-log4j12:${slf4jVersion}// Addtestdependencies here. // testCompilejunit:junit:4.12}//makecompileOnly dependencies availablefortests: sourceSets{main.compileClasspathconfigurations.flinkShadowJar main.runtimeClasspathconfigurations.flinkShadowJar test.compileClasspathconfigurations.flinkShadowJar test.runtimeClasspathconfigurations.flinkShadowJar javadoc.classpathconfigurations.flinkShadowJar}run.classpathsourceSets.main.runtimeClasspath jar{manifest{attributesBuilt-By:System.getProperty(user.name),Build-Jdk:System.getProperty(java.version)}}shadowJar{configurations[project.configurations.flinkShadowJar]}setting.gradlerootProject.namequickstart或者运行quickstart脚本bash-c$(curlhttps://flink.apache.org/q/gradle-quickstart.sh)--1.9.02.11查看目录结构tree quickstart/ quickstart/ ├── README ├── build.gradle ├── settings.gradle └── src └── main ├──java│ └── org │ └── myorg │ └── quickstart │ ├── BatchJob.java │ └── StreamingJob.java └── resources └── log4j.properties示例项目是一个Gradle 项目它包含了两个类StreamingJob和BatchJob是DataStream和DataSet程序的基础骨架程序。main方法是程序的入口即可用于IDE测试/执行也可用于部署。我们建议你将此项目导入你的 IDE来开发和测试它。IntelliJ IDEA 在安装Gradle插件后支持 Gradle 项目。Eclipse 则通过 Eclipse Buildship 插件支持 Gradle 项目鉴于shadow插件对 Gradle 版本有要求请确保在导入向导的最后一步指定 Gradle 版本 3.0。你也可以使用 Gradle’s IDE integration 从 Gradle 创建项目文件。构建项目如果你想要构建/打包项目请在项目目录下运行 ‘gradle clean shadowJar’ 命令。命令执行后你将找到一个 JAR 文件里面包含了你的应用程序以及已作为依赖项添加到应用程序的连接器和库build/libs/--all.jar。注意如果你使用其他类而不是StreamingJob作为应用程序的主类/入口我们建议你相应地修改build.gradle文件中的mainClassName配置。这样Flink 可以从 JAR 文件运行应用程序而无需另外指定主类。

相关新闻

Java实习面试高频考点解析与实战技巧

Java实习面试高频考点解析与实战技巧

1. Java实习面试通关指南:那些被问烂的题目与实战解法刚结束三个月的地狱式刷题,终于拿下了某大厂的Java实习Offer。作为面过15公司的"老油条",我发现80%的面试问题都来自那几个固定题库。今天就把这些高频考点掰开揉碎&#xff0c…

2026/7/21 3:29:01 阅读更多 →
AI与自动化本质区别:决策机制、学习能力与技术范式辨析

AI与自动化本质区别:决策机制、学习能力与技术范式辨析

1. 这不是AI,是自动化——一个被严重误用的词正在拖垮整个行业的认知基础我第一次在客户现场听到“我们上线了AI客服系统”这句话时,正蹲在机房角落调试一台老旧的票据扫描仪。客户CTO拍着我的肩膀,语气里带着一种近乎虔诚的兴奋,…

2026/7/21 3:27:52 阅读更多 →
职场高效学习系统:破除学习幻觉的实战方法论

职场高效学习系统:破除学习幻觉的实战方法论

1. 为什么我们总在"假装学习"?上周整理书架时翻出五本塑封完好的"年度必读书",健身App里存着三个月没打开的课程,收藏夹里吃灰的"Python入门教程"已经积了厚厚一层电子尘埃——这大概就是当代职场人最熟悉的学…

2026/7/21 3:27:52 阅读更多 →

最新新闻

19-插件系统入门-安全安装与核心管理

19-插件系统入门-安全安装与核心管理

19 插件系统入门:安全安装与核心管理 ——少即是多,稳才是真 小杨第一次接触Obsidian时,被社区里琳琅满目的插件惊呆了。主题美化、日历视图、思维导图、自动补全……每一个看起来都"刚需"。他一口气装了50多个插件,满心期待Obsidian变身为超级生产力工具。结…

2026/7/21 14:10:02 阅读更多 →
GalTransl终极指南:零基础掌握AI游戏汉化全流程

GalTransl终极指南:零基础掌握AI游戏汉化全流程

GalTransl终极指南:零基础掌握AI游戏汉化全流程 【免费下载链接】GalTransl 支持GPT-4/Claude/Deepseek/Sakura等大语言模型的Galgame自动化翻译解决方案 Automated translation solution for visual novels supporting GPT-4/Claude/Deepseek/Sakura 项目地址: h…

2026/7/21 14:09:55 阅读更多 →
TI PRU中断控制器INTC:三级映射与优先级仲裁实战解析

TI PRU中断控制器INTC:三级映射与优先级仲裁实战解析

1. PRU中断控制器:实时系统的神经中枢如果你在搞基于TI Sitara系列处理器的嵌入式实时项目,比如用AM335x做电机控制,或者用AM437x做高速数据采集,那你肯定绕不开PRU(Programmable Real-Time Unit)。这俩200…

2026/7/21 14:09:48 阅读更多 →
TI处理器PRCM模块PLL寄存器配置实战:从原理到稳定时钟输出

TI处理器PRCM模块PLL寄存器配置实战:从原理到稳定时钟输出

1. 项目概述:从寄存器手册到实战配置如果你正在开发基于德州仪器(TI)处理器(比如AM335x, AM437x, AM57xx系列)的嵌入式系统,并且已经翻开了那本动辄几千页的技术参考手册(TRM)&#…

2026/7/21 14:09:43 阅读更多 →
084、色彩校正矩阵:CCM设计与色彩管理系统的协同

084、色彩校正矩阵:CCM设计与色彩管理系统的协同

084、色彩校正矩阵:CCM设计与色彩管理系统的协同 一个让我失眠三天的色偏问题 去年夏天,某款旗舰手机的主摄模组在量产前夜,QA报告了一个诡异的问题:同一批次的Sensor,在D65光源下拍摄灰卡,R/G/B通道的响应曲线居然有0.8%的差异。生产线的兄弟急得跳脚,说再不解决就要延…

2026/7/21 14:09:40 阅读更多 →
提升开发效率的10个CLI工具:Awesome Vibe Coding终端助手精选

提升开发效率的10个CLI工具:Awesome Vibe Coding终端助手精选

提升开发效率的10个CLI工具:Awesome Vibe Coding终端助手精选 【免费下载链接】awesome-vibe-coding A hand-picked collection of tools and resources for Vibe Coding 项目地址: https://gitcode.com/gh_mirrors/aweso/awesome-vibe-coding 在AI辅助开发的…

2026/7/21 14:08:22 阅读更多 →

日新闻

Octane Render与C4D汉化版安装与优化指南

Octane Render与C4D汉化版安装与优化指南

1. Octane Render与C4D的黄金组合:为什么选择这个方案?在三维创作领域,渲染器的选择往往决定了作品的最终呈现质量和工作效率。作为Cinema 4D(C4D)用户,Octane Render的GPU加速特性与实时预览功能&#xff…

2026/7/21 0:00:19 阅读更多 →
GPMC接口设计:异步/同步模式与多路复用配置实战

GPMC接口设计:异步/同步模式与多路复用配置实战

1. GPMC接口设计:从硬件连接到软件配置的全局视角在嵌入式系统开发中,尤其是基于TI Sitara系列如AM263x这类高性能微控制器的项目里,外部存储器的扩展几乎是绕不开的一环。无论是存放大量非易失性代码的NOR Flash,还是作为高速数据…

2026/7/21 0:00:19 阅读更多 →
UE5 GAS框架下RPG被动技能系统:从核心原理到实战实现

UE5 GAS框架下RPG被动技能系统:从核心原理到实战实现

1. 项目概述:UE5 GAS RPG被动技能的核心价值在UE5里用GAS(Gameplay Ability System)做RPG游戏,主动技能像是你手里的武器,按一下打一下,逻辑直接,反馈也快。但被动技能,它更像是你身…

2026/7/21 0:00:19 阅读更多 →

周新闻

Go语言静态资源打包方案对比与实践指南

Go语言静态资源打包方案对比与实践指南

1. 项目背景与核心需求在Go语言开发中,我们经常需要处理静态资源文件的打包问题。无论是Web应用的模板文件、前端资源,还是配置文件、证书等,都需要随程序一起分发。传统做法是将这些文件与编译后的二进制文件放在同一目录下,但这…

2026/7/21 8:48:31 阅读更多 →
Go语言实现高性能LDAP认证服务的架构与实践

Go语言实现高性能LDAP认证服务的架构与实践

1. 项目背景与核心价值LDAP(轻量级目录访问协议)作为企业级身份认证的黄金标准,已经服务了超过80%的财富500强公司。我在金融科技领域实施统一认证体系时,发现传统Java方案存在启动慢、内存占用高等痛点。而Go语言凭借其协程并发模…

2026/7/21 5:34:47 阅读更多 →
【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

更多请点击: https://intelliparadigm.com 第一章:AI面试官实战指南的核心价值与适用场景 AI面试官并非替代人类HR的“黑箱工具”,而是以可解释、可审计、可迭代的方式,赋能招聘全链路的关键基础设施。其核心价值在于将主观经验沉…

2026/7/21 8:25:39 阅读更多 →

月新闻