当前位置: 首页 > article >正文

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

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这种基础案例建议新手先理解数据流动的完整路径尝试修改并行度观察变化逐步添加窗口、状态等复杂功能

相关文章:

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

1. 环境准备:双语言开发环境搭建 第一次接触Flink时,最让人头疼的就是环境配置。记得2018年我刚从Hadoop转向Flink时,光环境搭建就折腾了两天。现在回想起来,其实只要掌握几个关键点,10分钟就能搞定一个可用的开发环境…...

Windows下用mitmweb抓包实战:从安装证书到过滤百度请求的完整流程

Windows下mitmweb抓包实战:从证书安装到精准流量过滤 引言 在Web开发和测试领域,流量监控与分析是不可或缺的技能。对于Windows平台用户而言,寻找一款高效、易用的抓包工具往往面临诸多挑战。mitmproxy作为业界知名的中间人代理工具&#x…...

AIVideo视频水印技术:基于神经网络的隐形水印方案

AIVideo视频水印技术:基于神经网络的隐形水印方案 1. 引言 视频内容保护一直是创作者们头疼的问题。传统的可见水印影响观看体验,而简单的隐形水印又容易被去除。今天要介绍的AIVideo基于神经网络开发的隐形水印技术,可以说是给视频版权保护…...

Dify前端DIY指南:从修改样式到Docker部署的完整避坑手册

Dify前端DIY指南:从修改样式到Docker部署的完整避坑手册 当你需要为企业内部系统打造独特的品牌界面,或是为教学演示环境定制专属交互体验时,Dify的前端定制能力就显得尤为重要。不同于简单的主题切换,深度定制Dify前端需要掌握从…...

别再手动写CRUD了!用RuoYi代码生成器5分钟搞定MinIO素材管理模块

5分钟极速构建MinIO素材管理系统:RuoYi代码生成器实战指南 每次接到"三天内上线内容管理后台"的需求时,你是否还在重复着建表→写Controller→写Service→调试接口的机械劳动?作为经历过十几个企业级内容平台开发的架构师&#xff…...

Linux下Synopsys2020安装全攻略:从SCL配置到License生成避坑指南

Linux下Synopsys工具链部署实战:从权限管理到License优化的全流程解析 在芯片设计领域,Synopsys工具链的稳定运行直接关系到研发效率。不同于简单的软件安装,EDA工具的部署涉及复杂的权限管理、环境配置和License验证体系。本文将基于真实服务…...

LeetCode 3643.子矩阵垂直翻转算法解析

LeetCode 3643.子矩阵垂直翻转算法解析 题目描述 给定一个二维矩阵 grid 和四个参数 (x, y, k),实现一个函数,将矩阵中以 (x, y) 为左上角、边长为 k 的正方形子矩阵进行上下翻转(垂直镜像翻转)。 算法思路 本题的核心是实现子矩阵…...

Ollama+granite-4.0-h-350m:开源轻量模型在学生编程作业辅导中的应用

Ollamagranite-4.0-h-350m:开源轻量模型在学生编程作业辅导中的应用 1. 为什么需要轻量级编程辅导助手? 作为一名计算机专业的学生,我经常遇到这样的困境:深夜调试代码时遇到问题,找不到人请教;想要理解一…...

基于Ubuntu 24.04与Zabbix 7.0构建云服务器监控体系

1. 环境准备与基础配置 在阿里云ECS上部署Zabbix监控系统前,需要做好充分的环境准备。我建议选择4核8G配置的实例作为Zabbix Server主机,这个配置可以轻松应对中小规模集群的监控需求。实测下来,100G的系统盘空间完全够用,还能保留…...

2024年还用Windows XP?VMware17虚拟化实战:从系统封装到快照管理

2024年企业级Windows XP虚拟化实战:VMware17高级运维指南 在工业控制、金融终端等关键领域,仍有大量关键业务系统依赖Windows XP环境运行。根据行业调研数据显示,全球范围内仍有约3%的企业设备运行这一经典系统,其中银行ATM机和数…...

HY-Motion 1.0参数怎么调?采样步数、动作时长设置全解析

HY-Motion 1.0参数怎么调?采样步数、动作时长设置全解析 [【免费上手链接】HY-Motion 1.0 腾讯混元3D数字人团队开源动作生成模型,十亿参数级文生动作系统,支持一键可视化操作,让文字自然转化为电影级3D律动 镜像地址&#xff1…...

DeepSeek-R1-Distill-Qwen-7B数学推理能力实测:AIME竞赛题解题分析

DeepSeek-R1-Distill-Qwen-7B数学推理能力实测:AIME竞赛题解题分析 1. 引言 如果你关注过最近的大模型进展,应该听说过DeepSeek-R1这个名字。这个系列模型在数学推理能力上表现相当亮眼,特别是那个671B参数的版本,在AIME竞赛题上…...

RevokeMsgPatcher完整指南:让微信/QQ/TIM消息不再消失的终极方案

RevokeMsgPatcher完整指南:让微信/QQ/TIM消息不再消失的终极方案 【免费下载链接】RevokeMsgPatcher :trollface: A hex editor for WeChat/QQ/TIM - PC版微信/QQ/TIM防撤回补丁(我已经看到了,撤回也没用了) 项目地址: https://…...

PSpice AD仿真避坑指南:为什么你的新器件模型导入后无法运行?

PSpice AD仿真避坑指南:为什么你的新器件模型导入后无法运行? 作为一名长期使用PSpice AD进行电路仿真的工程师,我深知导入新器件模型时可能遇到的各种"坑"。本文将结合我的实战经验,系统梳理7个最常见的问题根源及对应…...

SparkFun MicroPressure库解析:MPR微压传感器嵌入式驱动设计

1. SparkFun MicroPressure 库深度解析:Honeywell MPR 系列微压传感器的嵌入式驱动实现1.1 库定位与工程价值SparkFun MicroPressure Library 是一个专为 Honeywell MPR 系列微压传感器(MPR121、MPR031、MPR032 等)设计的轻量级嵌入式 C/C 驱…...

2026程序员破局指南:大模型技能是未来,收藏这份转型路线图

引言 曾几何时,程序员被誉为“21世纪最高薪的职业之一”,是无数人向往的“金饭碗”。然而,步入2026年,这个曾经风光无限的职业似乎正经历一场前所未有的“寒冬”。裁员潮、降薪、AI冲击……种种挑战接踵而至,让许多程…...

避坑指南:SAP供应商付款时F-51的这两种清账方式千万别乱选

SAP供应商付款清账实战避坑:F-51操作中的关键决策逻辑 刚接手SAP财务模块的新人,往往会在供应商付款环节踩坑——尤其是面对F-51事务码中的部分清账与剩余清账选项时。这两个看似简单的功能选择,实际上会直接影响后续对账效率、账龄分析准确性…...

DAMOYOLO-S快速部署:7860端口Web服务访问与HTTPS反向代理配置

DAMOYOLO-S快速部署:7860端口Web服务访问与HTTPS反向代理配置 你是不是也想快速搭建一个属于自己的目标检测服务,但被复杂的模型部署和环境配置劝退?今天,我们就来聊聊如何一键部署DAMOYOLO-S这个高性能的通用检测模型&#xff0…...

Windows10下ENSP云服务配置:解决以太网卡缺失的3种实用方法

Windows 10环境下ENSP云服务配置:全面解决以太网卡缺失问题 在数字化转型浪潮中,网络工程师和IT专业人员经常需要借助仿真工具进行网络拓扑设计和测试。华为ENSP作为业界广泛使用的网络模拟平台,其云服务功能对物理网卡的依赖却让许多使用轻薄…...

VibeVoice实时语音合成系统效果测评:流式播放与长文本支持实测

VibeVoice实时语音合成系统效果测评:流式播放与长文本支持实测 1. 测试环境与准备 1.1 硬件配置 本次测试使用的硬件平台为: GPU:NVIDIA RTX 4090(24GB显存)CPU:AMD Ryzen 9 7950X内存:64GB…...

信号完整性(SIPI)实战解析:高速设计中的串扰抑制策略

1. 串扰的本质与高速设计的挑战 当你把两根电线靠得太近时,总会听到"滋滋"的干扰声——这就是生活中最常见的串扰现象。在高速PCB设计中,这种干扰被放大了无数倍。我最近调试一块DDR5内存板时就深有体会:当数据速率冲到6400Mbps时&…...

学霸同款! 降AIGC网站 千笔·降AIGC助手 VS 知文AI,开源免费首选

在AI技术迅猛发展的今天,越来越多的学生和研究者开始借助AI工具提升论文写作效率。然而,随着学术审查标准的不断升级,AI生成内容的痕迹愈发明显,论文中的AIGC率和重复率问题成为困扰无数人的“隐形炸弹”。面对查重系统日益严格的…...

Python音频处理实战:用wave和numpy生成自定义WAV音效(附完整代码)

Python音频处理实战:用wave和numpy生成自定义WAV音效 1. 音频合成基础与核心概念 音频合成是现代数字音频处理的基础技术之一。想象一下,你正在为一个独立游戏开发音效系统,或者为某个艺术装置设计交互式声音反馈,Python的wave和…...

从汽车NVH到风电监测:阶次跟踪技术的5个跨界应用案例解析

从汽车NVH到风电监测:阶次跟踪技术的5个跨界应用案例解析 阶次跟踪(Order Tracking)技术正悄然改变着工业领域的故障诊断与性能优化方式。这项基于旋转机械转速同步采样的分析方法,已从传统的发动机测试领域,逐步渗透到…...

YOLO标注文件可视化保姆级教程:用Python+OpenCV把txt里的数字变成图像上的框

YOLO标注文件可视化实战指南:从原理到批量处理的完整解决方案 当你第一次拿到YOLO格式的数据集时,面对那些充满数字的txt文件,是否感到无从下手?本文将带你深入理解YOLO标注格式的本质,并手把手教你用Python和OpenCV将…...

vLLM部署千问72B大模型实战:从Docker镜像到API调用的完整避坑指南

vLLM实战:千问72B大模型高效部署与API服务优化指南 在人工智能技术快速迭代的今天,百亿参数级别的大模型已成为企业智能化转型的核心竞争力。如何高效部署这些"庞然大物",使其在实际业务中发挥价值,是每个技术团队面临的…...

MATLAB新手也能搞定!鼠笼式电机矢量控制仿真全流程(附源码)

MATLAB新手也能搞定!鼠笼式电机矢量控制仿真全流程(附源码) 鼠笼式三相交流异步电动机在工业领域应用广泛,而矢量控制技术则是实现其高性能调速的关键。对于电气工程或自动化专业的学生和工程师来说,掌握MATLAB/SIMUL…...

CAN总线信号示波器测试全流程指南

1. CAN总线信号测试的工程实践方法CAN(Controller Area Network)总线自1986年由Bosch公司提出以来,已成为车载电子系统中事实上的通信标准。其差分传输机制、非破坏性仲裁、高抗干扰能力及完善的错误检测机制,使其在汽车动力总成、…...

保姆级教程:用STM32的TIM3测PWM频率和占空比(附完整代码)

STM32实战指南:TIM3精准捕获PWM频率与占空比全解析 在嵌入式开发中,精确测量外部PWM信号的频率和占空比是常见需求。无论是电机控制、传感器数据采集还是通信协议解析,这项技能都至关重要。本文将带您从零开始,使用STM32的TIM3定时…...

xv6 Lab6 COW Fork避坑实录:从引用计数到usertrap,手把手教你搞定MIT操作系统实验

MIT 6.S081 Lab6 COW Fork全攻略:从引用计数陷阱到usertrap实战解析 在操作系统课程中,MIT 6.S081的Lab6堪称一道分水岭——它要求学生在xv6内核中实现Copy-on-Write Fork机制。这个实验不仅考验对虚拟内存系统的理解深度,更需要处理引用计数…...