Flink DataStreamAPI实战指南——从环境搭建到WordCount(Java/Scala双语言版)

news2026/3/22 9:57:17
1. 环境准备双语言开发环境搭建第一次接触Flink时最让人头疼的就是环境配置。记得2018年我刚从Hadoop转向Flink时光环境搭建就折腾了两天。现在回想起来其实只要掌握几个关键点10分钟就能搞定一个可用的开发环境。1.1 JDK版本选择Flink 1.16.x对JDK的要求比较灵活支持JDK 8和JDK 11。但这里有个坑需要注意如果你后续需要整合Hive等组件建议选择JDK 8。我去年在一个金融项目中就踩过这个坑当时用JDK 11跑得好好的一整合Hive 3.1.3就各种报错。安装完JDK后记得检查环境变量java -version javac -version这两个命令应该显示相同的版本号否则后续Maven编译会出问题。1.2 IDE配置技巧IntelliJ IDEA确实是Flink开发的首选特别是它的Scala插件非常智能。有个小技巧分享安装插件时直接搜索Scala可能会找到多个版本建议选择JetBrains官方维护的那个。安装完成后创建一个空项目时我习惯这样组织模块结构MyFlinkProject ├── java-module (Java代码) └── scala-module (Scala代码)这种结构比单独建两个项目更便于管理共用依赖。1.3 Maven配置实战Maven的settings.xml配置直接影响依赖下载速度。建议配置阿里云镜像mirror idaliyunmaven/id mirrorOf*/mirrorOf name阿里云公共仓库/name urlhttps://maven.aliyun.com/repository/public/url /mirror对于Flink 1.16.0核心依赖配置要注意Scala版本后缀的变化。从1.15开始如果你只用Java API可以不用带scala后缀的包。但用Scala API时必须匹配Scala 2.12版本。2. 项目初始化双语言工程搭建2.1 Java模块创建创建Java模块时我推荐使用Maven的原型(archetype)maven-archetype-quickstart这样生成的pom.xml比较干净方便后续添加Flink专属依赖。关键依赖配置dependency groupIdorg.apache.flink/groupId artifactIdflink-clients/artifactId version${flink.version}/version /dependency这个依赖已经包含了DataStream API的核心功能。2.2 Scala模块特殊配置Scala模块创建后需要特别注意两点在Project Structure中添加Scala SDK右键模块选择Add Framework Support添加Scala支持pom.xml中必须明确指定Scala版本properties scala.version2.12.15/scala.version scala.binary.version2.12/scala.binary.version /propertiesScala的依赖比Java复杂需要三个核心包dependency groupIdorg.apache.flink/groupId artifactIdflink-scala_${scala.binary.version}/artifactId version${flink.version}/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-scala_${scala.binary.version}/artifactId version${flink.version}/version /dependency2.3 日志配置技巧很多新手运行程序时看不到日志输出这是因为缺少log4j配置。建议在resources目录下创建log4j.propertieslog4j.rootLoggerERROR, console log4j.appender.consoleorg.apache.log4j.ConsoleAppender log4j.appender.console.layoutorg.apache.log4j.PatternLayout log4j.appender.console.layout.ConversionPattern%d{HH:mm:ss} %p %c{2}: %m%n对应的Maven依赖要特别注意版本兼容性dependency groupIdorg.slf4j/groupId artifactIdslf4j-log4j12/artifactId version1.7.36/version /dependency3. 批处理WordCount实现3.1 Java批处理实现Java版的批处理WordCount有几个关键点需要注意ExecutionEnvironment是批处理的入口flatMap操作需要指定returns类型Tuple2的泛型信息会被擦除需要显式声明完整代码示例public class BatchWordCountJava { public static void main(String[] args) throws Exception { ExecutionEnvironment env ExecutionEnvironment.getExecutionEnvironment(); DataSourceString lines env.readTextFile(input.txt); FlatMapOperatorString, String words lines.flatMap((String line, CollectorString out) - { for (String word : line.split( )) { out.collect(word); } }).returns(Types.STRING); MapOperatorString, Tuple2String, Integer wordPairs words.map(word - new Tuple2(word, 1)) .returns(Types.TUPLE(Types.STRING, Types.INT)); AggregateOperatorTuple2String, Integer counts wordPairs.groupBy(0).sum(1); counts.print(); } }3.2 Scala批处理实现Scala版代码更简洁但要注意隐式转换的导入import org.apache.flink.api.scala._ object BatchWordCountScala { def main(args: Array[String]): Unit { val env ExecutionEnvironment.getExecutionEnvironment val counts env.readTextFile(input.txt) .flatMap(_.split( )) .map((_, 1)) .groupBy(0) .sum(1) counts.print() } }这里有个性能优化技巧对于小数据集可以在print()前加上.counts.setParallelism(1)这样输出不会乱序。4. 流处理WordCount实现4.1 Java流处理实现流处理与批处理的主要区别使用StreamExecutionEnvironment需要调用execute()触发任务keyBy代替groupBy典型实现public class StreamingWordCountJava { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); DataStreamString lines env.readTextFile(input.txt); DataStreamTuple2String, Integer counts lines .flatMap((String line, CollectorTuple2String, Integer out) - { for (String word : line.split( )) { out.collect(new Tuple2(word, 1)); } }) .returns(Types.TUPLE(Types.STRING, Types.INT)) .keyBy(0) .sum(1); counts.print(); env.execute(Streaming WordCount); } }4.2 Scala流处理实现Scala流处理代码的链式调用非常优雅import org.apache.flink.streaming.api.scala._ object StreamingWordCountScala { def main(args: Array[String]): Unit { val env StreamExecutionEnvironment.getExecutionEnvironment val counts env.readTextFile(input.txt) .flatMap(_.split( )) .map((_, 1)) .keyBy(_._1) .sum(1) counts.print() env.execute() } }4.3 执行模式的选择从Flink 1.12开始可以通过setRuntimeMode统一批流处理env.setRuntimeMode(RuntimeExecutionMode.BATCH);三种模式的区别BATCH优化批处理执行计划STREAMING纯流模式AUTOMATIC根据数据源自动判断实际项目中建议在提交任务时指定模式flink run -Dexecution.runtime-modeBATCH -c MainClass app.jar5. 核心概念解析5.1 DataStream API设计思想Flink的DataStream API采用了惰性求值设计只有在调用execute()时才会真正执行。这种设计使得Flink可以优化整个执行计划。我经常用这个类比来解释就像写SQL一样前面的操作只是定义了一个查询计划最后执行时才真正运行。5.2 类型系统处理Java的类型擦除是个大问题Flink提供了两种解决方案通过returns()显式指定类型使用TypeHint保留泛型信息Scala由于有更丰富的类型信息通常不需要特别处理。5.3 并行度设置技巧并行度设置直接影响性能有几个经验值本地开发时设为1方便调试生产环境通常设为CPU核数的2-3倍可以通过env.setParallelism()全局设置查看并行度的好方法System.out.println(当前并行度 env.getParallelism());6. 常见问题排查6.1 类型推断失败这是Java版最常见的问题错误信息通常包含TypeInformation could not be created。解决方案检查所有lambda表达式是否都加了returns复杂类型建议实现ResultTypeQueryable接口6.2 依赖冲突特别是当整合Hadoop生态时容易发生jar包冲突。建议使用mvn dependency:tree查看依赖树用排除冲突包6.3 本地执行问题如果在IDEA中运行报错可以尝试添加scope为provided的依赖设置env.enableCheckpointing(1000)7. 性能优化建议7.1 批处理优化对于批处理作业合理设置批量大小使用sortPartition预排序考虑使用DataSet API(虽然已标记为Legacy)7.2 流处理优化流处理优化点设置合适的checkpoint间隔使用增量检查点配置合理的状态后端7.3 内存配置通过conf/flink-conf.yaml调整taskmanager.memory.process.size: 4096m taskmanager.memory.task.heap.size: 2048m8. 扩展应用场景8.1 对接Kafka实际项目中数据源通常是Kafkaenv.addSource(new FlinkKafkaConsumer(topic, new SimpleStringSchema(), props))8.2 使用状态有状态的流处理示例keyedStream.process(new KeyedProcessFunctionString, Tuple2String, Integer, String() { private ValueStateInteger state; Override public void open(Configuration parameters) { state getRuntimeContext().getState(new ValueStateDescriptor(count, Integer.class)); } Override public void processElement(Tuple2String, Integer value, Context ctx, CollectorString out) { int current state.value() null ? 0 : state.value(); current value.f1; state.update(current); out.collect(value.f0 : current); } });8.3 窗口计算滚动窗口示例stream.keyBy(_._1) .window(TumblingProcessingTimeWindows.of(Time.seconds(5))) .sum(1)9. 测试与调试9.1 单元测试Flink提供了专门的测试工具Test public void testPipeline() throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getTestEnvironment(); // 测试代码 }9.2 日志调试建议在开发时添加dataStream.map(value - { System.out.println(处理元素 value); return value; });9.3 Web UI使用启动本地环境后访问http://localhost:8081可以查看作业执行计划各个算子的吞吐量背压情况10. 生产环境建议10.1 资源配置根据数据量合理配置TaskManager数量每个TaskManager的slot数网络缓冲区大小10.2 监控方案推荐组合Prometheus Grafana监控指标ELK收集日志自定义告警规则10.3 升级策略Flink版本升级注意先在小规模测试环境验证检查API变更特别注意状态兼容性11. 最佳实践总结经过多个项目的实践我总结了这些经验开发环境尽量和生产环境保持一致重要作业要设置重启策略合理利用savepoint进行版本管理监控指标要包含延迟和吞吐量对于WordCount这种基础案例建议新手先理解数据流动的完整路径尝试修改并行度观察变化逐步添加窗口、状态等复杂功能

本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若转载,请注明出处:http://www.coloradmin.cn/o/2436562.html

如若内容造成侵权/违法违规/事实不符,请联系多彩编程网进行投诉反馈,一经查实,立即删除!

相关文章

SpringBoot-17-MyBatis动态SQL标签之常用标签

文章目录 1 代码1.1 实体User.java1.2 接口UserMapper.java1.3 映射UserMapper.xml1.3.1 标签if1.3.2 标签if和where1.3.3 标签choose和when和otherwise1.4 UserController.java2 常用动态SQL标签2.1 标签set2.1.1 UserMapper.java2.1.2 UserMapper.xml2.1.3 UserController.ja…

wordpress后台更新后 前端没变化的解决方法

使用siteground主机的wordpress网站,会出现更新了网站内容和修改了php模板文件、js文件、css文件、图片文件后,网站没有变化的情况。 不熟悉siteground主机的新手,遇到这个问题,就很抓狂,明明是哪都没操作错误&#x…

网络编程(Modbus进阶)

思维导图 Modbus RTU(先学一点理论) 概念 Modbus RTU 是工业自动化领域 最广泛应用的串行通信协议,由 Modicon 公司(现施耐德电气)于 1979 年推出。它以 高效率、强健性、易实现的特点成为工业控制系统的通信标准。 包…

UE5 学习系列(二)用户操作界面及介绍

这篇博客是 UE5 学习系列博客的第二篇,在第一篇的基础上展开这篇内容。博客参考的 B 站视频资料和第一篇的链接如下: 【Note】:如果你已经完成安装等操作,可以只执行第一篇博客中 2. 新建一个空白游戏项目 章节操作,重…

IDEA运行Tomcat出现乱码问题解决汇总

最近正值期末周,有很多同学在写期末Java web作业时,运行tomcat出现乱码问题,经过多次解决与研究,我做了如下整理: 原因: IDEA本身编码与tomcat的编码与Windows编码不同导致,Windows 系统控制台…

利用最小二乘法找圆心和半径

#include <iostream> #include <vector> #include <cmath> #include <Eigen/Dense> // 需安装Eigen库用于矩阵运算 // 定义点结构 struct Point { double x, y; Point(double x_, double y_) : x(x_), y(y_) {} }; // 最小二乘法求圆心和半径 …

使用docker在3台服务器上搭建基于redis 6.x的一主两从三台均是哨兵模式

一、环境及版本说明 如果服务器已经安装了docker,则忽略此步骤,如果没有安装,则可以按照一下方式安装: 1. 在线安装(有互联网环境): 请看我这篇文章 传送阵>> 点我查看 2. 离线安装(内网环境):请看我这篇文章 传送阵>> 点我查看 说明&#xff1a;假设每台服务器已…

XML Group端口详解

在XML数据映射过程中&#xff0c;经常需要对数据进行分组聚合操作。例如&#xff0c;当处理包含多个物料明细的XML文件时&#xff0c;可能需要将相同物料号的明细归为一组&#xff0c;或对相同物料号的数量进行求和计算。传统实现方式通常需要编写脚本代码&#xff0c;增加了开…

LBE-LEX系列工业语音播放器|预警播报器|喇叭蜂鸣器的上位机配置操作说明

LBE-LEX系列工业语音播放器|预警播报器|喇叭蜂鸣器专为工业环境精心打造&#xff0c;完美适配AGV和无人叉车。同时&#xff0c;集成以太网与语音合成技术&#xff0c;为各类高级系统&#xff08;如MES、调度系统、库位管理、立库等&#xff09;提供高效便捷的语音交互体验。 L…

(LeetCode 每日一题) 3442. 奇偶频次间的最大差值 I (哈希、字符串)

题目&#xff1a;3442. 奇偶频次间的最大差值 I 思路 &#xff1a;哈希&#xff0c;时间复杂度0(n)。 用哈希表来记录每个字符串中字符的分布情况&#xff0c;哈希表这里用数组即可实现。 C版本&#xff1a; class Solution { public:int maxDifference(string s) {int a[26]…

【大模型RAG】拍照搜题技术架构速览:三层管道、两级检索、兜底大模型

摘要 拍照搜题系统采用“三层管道&#xff08;多模态 OCR → 语义检索 → 答案渲染&#xff09;、两级检索&#xff08;倒排 BM25 向量 HNSW&#xff09;并以大语言模型兜底”的整体框架&#xff1a; 多模态 OCR 层 将题目图片经过超分、去噪、倾斜校正后&#xff0c;分别用…

【Axure高保真原型】引导弹窗

今天和大家中分享引导弹窗的原型模板&#xff0c;载入页面后&#xff0c;会显示引导弹窗&#xff0c;适用于引导用户使用页面&#xff0c;点击完成后&#xff0c;会显示下一个引导弹窗&#xff0c;直至最后一个引导弹窗完成后进入首页。具体效果可以点击下方视频观看或打开下方…

接口测试中缓存处理策略

在接口测试中&#xff0c;缓存处理策略是一个关键环节&#xff0c;直接影响测试结果的准确性和可靠性。合理的缓存处理策略能够确保测试环境的一致性&#xff0c;避免因缓存数据导致的测试偏差。以下是接口测试中常见的缓存处理策略及其详细说明&#xff1a; 一、缓存处理的核…

龙虎榜——20250610

上证指数放量收阴线&#xff0c;个股多数下跌&#xff0c;盘中受消息影响大幅波动。 深证指数放量收阴线形成顶分型&#xff0c;指数短线有调整的需求&#xff0c;大概需要一两天。 2025年6月10日龙虎榜行业方向分析 1. 金融科技 代表标的&#xff1a;御银股份、雄帝科技 驱动…

观成科技:隐蔽隧道工具Ligolo-ng加密流量分析

1.工具介绍 Ligolo-ng是一款由go编写的高效隧道工具&#xff0c;该工具基于TUN接口实现其功能&#xff0c;利用反向TCP/TLS连接建立一条隐蔽的通信信道&#xff0c;支持使用Let’s Encrypt自动生成证书。Ligolo-ng的通信隐蔽性体现在其支持多种连接方式&#xff0c;适应复杂网…

铭豹扩展坞 USB转网口 突然无法识别解决方法

当 USB 转网口扩展坞在一台笔记本上无法识别,但在其他电脑上正常工作时,问题通常出在笔记本自身或其与扩展坞的兼容性上。以下是系统化的定位思路和排查步骤,帮助你快速找到故障原因: 背景: 一个M-pard(铭豹)扩展坞的网卡突然无法识别了,扩展出来的三个USB接口正常。…

未来机器人的大脑:如何用神经网络模拟器实现更智能的决策?

编辑&#xff1a;陈萍萍的公主一点人工一点智能 未来机器人的大脑&#xff1a;如何用神经网络模拟器实现更智能的决策&#xff1f;RWM通过双自回归机制有效解决了复合误差、部分可观测性和随机动力学等关键挑战&#xff0c;在不依赖领域特定归纳偏见的条件下实现了卓越的预测准…

Linux应用开发之网络套接字编程(实例篇)

服务端与客户端单连接 服务端代码 #include <sys/socket.h> #include <sys/types.h> #include <netinet/in.h> #include <stdio.h> #include <stdlib.h> #include <string.h> #include <arpa/inet.h> #include <pthread.h> …

华为云AI开发平台ModelArts

华为云ModelArts&#xff1a;重塑AI开发流程的“智能引擎”与“创新加速器”&#xff01; 在人工智能浪潮席卷全球的2025年&#xff0c;企业拥抱AI的意愿空前高涨&#xff0c;但技术门槛高、流程复杂、资源投入巨大的现实&#xff0c;却让许多创新构想止步于实验室。数据科学家…

深度学习在微纳光子学中的应用

深度学习在微纳光子学中的主要应用方向 深度学习与微纳光子学的结合主要集中在以下几个方向&#xff1a; 逆向设计 通过神经网络快速预测微纳结构的光学响应&#xff0c;替代传统耗时的数值模拟方法。例如设计超表面、光子晶体等结构。 特征提取与优化 从复杂的光学数据中自…