大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载在 Flink DataStream 编程中几乎每一个数据转换算子map、filter、reduce等都需要一个用户自定义函数User-Defined FunctionUDF来定义具体的处理逻辑。本文基于当前仓库中的官方文档docs/content/docs/dev/datastream/user_defined_functions.md系统讲解 Java 与 Scala 两种 API 下定义 UDF 的全部方式——接口实现、匿名类、Lambda 表达式与 Rich Function并深入剖析 Flink 内置的 Accumulator累加器体系从IntCounter、Histogram等内置实现到注册、使用、结果回传的完整流程再到自定义 Accumulator 的接口设计。读完本文你将掌握在实际 Flink 作业中编写可运行 UDF、并借助累加器在作业结束后获取全局统计信息的完整实战方案。UDF 的定义方式大多数 DataStream 算子如map、filter、reduce都要求传入一个用户自定义函数。本节介绍 Java API 与 Scala API 中所有可用的函数指定方式。Java API实现函数接口最基本的方式是实现 Flink 提供的函数接口。例如实现MapFunctionString, Integer将字符串解析为整数class MyMapFunction implements MapFunctionString, Integer { public Integer map(String value) { return Integer.parseInt(value); } } data.map(new MyMapFunction());MapFunction接口位于 flink-core/src/main/java/org/apache/flink/api/common/functions/MapFunction.java其唯一的抽象方法map(IN value)定义输入元素到输出元素的转换逻辑。Flink 为每个常用算子都提供了对应的函数接口如FilterFunction、FlatMapFunction、ReduceFunction等它们共同组成了 DataStream API 的函数类型体系。Java API匿名类如果不希望单独声明一个类可以直接以匿名类的方式传入函数data.map(new MapFunctionString, Integer () { public Integer map(String value) { return Integer.parseInt(value); } });匿名类适用于逻辑简单、只在单个算子中使用的场景省去了定义独立类文件的样板代码。Java APIJava 8 Lambda 表达式Flink 的 Java API 原生支持 Java 8 的 Lambda 表达式可以大幅简化代码。对于只有一个抽象方法的函数式接口如MapFunction、FilterFunction、ReduceFunction编译器会自动推断出对应的函数类型data.filter(s - s.startsWith(http://));data.reduce((i1,i2) - i1 i2);需要注意的是当使用 Lambda 表达式时Flink 依赖类型推断来确认输入输出类型。在某些复杂场景下如泛型嵌套、或需要为结果指定特定类型时可能需要显式提供类型信息关于这一点可以参考仓库中 docs/content/docs/dev/datastream/java_lambdas.md 的专门讨论。Java APIRich Function富函数所有要求传入 UDF 的转换算子都可以改传一个rich富函数。富函数在普通函数接口的基础上额外获得了生命周期回调与运行时上下文访问能力。例如将上面普通的MyMapFunction改为富函数版本class MyMapFunction extends RichMapFunctionString, Integer { public Integer map(String value) { return Integer.parseInt(value); } }然后照常将它传给map转换data.map(new MyMapFunction());富函数同样可以定义为匿名类data.map (new RichMapFunctionString, Integer() { public Integer map(String value) { return Integer.parseInt(value); } });从源码看RichMapFunction定义于 flink-core/src/main/java/org/apache/flink/api/common/functions/RichMapFunction.java它继承自AbstractRichFunction并实现MapFunctionIN, OUT因此既保留了map的核心逻辑又获得了富函数的所有能力。所有Rich*FunctionRichFilterFunction、RichFlatMapFunction、RichReduceFunction等的根接口是 flink-core/src/main/java/org/apache/flink/api/common/functions/RichFunction.java它定义了富函数的核心能力生命周期方法open(...)在函数真正开始处理数据如map、join之前调用一次适合做初始化工作比如连接外部系统、加载配置、注册累加器close()在最后一次数据处理之后调用适合做清理工作。运行时上下文getRuntimeContext()返回 RuntimeContext可获取子任务索引、任务名称、并行度等信息并访问累加器与分布式缓存getIterationRuntimeContext()则在函数参与迭代计算时提供额外的迭代上下文。一个关于open方法的重要版本说明当前仓库中open(Configuration parameters)已标记为Deprecated自 Flink 1.19 起废弃官方推荐实现open(OpenContext openContext)重载并把open(Configuration parameters)留空若只实现了旧的open(Configuration)则默认的open(OpenContext)实现会自动转发调用它参见 RichFunction.java 中的注释说明。Scala APILambda 函数与 Java 类似Scala API 的所有算子都接受 Lambda 函数来描述操作val data: DataStream[String] // [...] data.filter { _.startsWith(http://) }val data: DataStream[Int] // [...] data.reduce { (i1,i2) i1 i2 } // 或者更简洁的写法 data.reduce { _ _ }Scala APIRich FunctionScala 中所有接受 Lambda 函数的转换算子同样可以改传富函数。例如把下面的 Lambda 写法data.map { x x.toInt }改写为富函数版本class MyMapFunction extends RichMapFunction[String, Int] { def map(in: String): Int in.toInt }并照常传给map转换data.map(new MyMapFunction())富函数同样可以定义为匿名类data.map (new RichMapFunction[String, Int] { def map(in: String): Int in.toInt })Accumulators 与 Counters作业级统计利器Accumulator累加器是一种简单的构造它包含一个add 操作和一个最终累加结果最终结果在作业结束后可用。最直观的累加器是counter计数器你可以通过Accumulator.add(V value)方法对其递增。作业结束时Flink 会把所有并行分片的局部结果合并merge起来并把最终结果发送回客户端。累加器在调试阶段或希望快速了解数据分布时非常有用。内置累加器Flink 目前内置了以下累加器它们全部实现了Accumulator接口接口定义见 flink-core/src/main/java/org/apache/flink/api/common/accumulators/Accumulator.javaIntCounter、LongCounter、DoubleCounter整数/长整型/双精度浮点计数器用于统计计数类指标。三者分别位于 IntCounter.java、LongCounter.java 和 DoubleCounter.java下面给出使用计数器的完整示例。Histogram针对离散数值分桶的直方图实现内部本质上是一个Integer - Integer的映射具体实现为TreeMapInteger, Integer可用于计算值的分布。例如统计一个 WordCount 程序中每行单词数的分布情况。实现见 Histogram.java。从 Accumulator.java 的源码注释可以看到Flink 的累加器设计受 Hadoop/MapReduce 的计数器启发每个并行实例创建并更新自己的累加器对象系统在作业结束时将这些并行实例合并最终结果既可以从作业执行结果中获取也可以在 Web 运行时监控界面查看。累加器的使用步骤第一步创建累加器对象。在你希望使用累加器的用户自定义转换函数中创建累加器对象这里以计数器为例private IntCounter numLines new IntCounter();第二步注册累加器。通常在富函数的open()方法中注册累加器并在此处为它指定名称getRuntimeContext().addAccumulator(num-lines, this.numLines);addAccumulator(String name, AccumulatorV, A accumulator)是 RuntimeContext 提供的接口方法对应的getAccumulator(String name)可在函数内部按名称取出累加器对象。第三步使用累加器。注册之后你可以在算子函数的任何位置使用它包括open()和close()方法中this.numLines.add(1);第四步获取最终结果。累加器的整体结果会存放在JobExecutionResult对象中该对象由执行环境的execute()方法返回目前这仅在执行等待作业完成时才有效myJobExecutionResult.getAccumulatorResult(num-lines);getAccumulatorResult(String accumulatorName)定义于 flink-core/src/main/java/org/apache/flink/api/common/JobExecutionResult.java它按名称返回累加器值如果作业中没有产生名为该名称的累加器则返回null。同文件中的getAllAccumulatorResults()可以一次性取出作业产生的所有累加器结果而JobExecutionResult.toString()在打印时会自动附上格式化的 Accumulator Results 摘要便于直接查看。累加器的工作机制与注意事项全局命名空间所有累加器在单个作业内共享同一个命名空间。因此你可以在作业的不同算子函数中使用同名累加器Flink 会在内部将所有同名累加器合并。这一合并行为正是由Accumulator接口的merge(AccumulatorV, R other)方法驱动的。迭代计算限制目前累加器的结果只有在整个作业结束后才可用。官方计划在未来让上一轮迭代的结果在下一轮迭代中可用。在需要按迭代轮次计算统计信息、并基于这些统计信息决定迭代终止条件时可以使用Aggregators见 flink-java/src/main/java/org/apache/flink/api/java/operators/IterativeDataSet.java 中相关方法。自定义累加器要实现自己的累加器只需实现Accumulator接口。你有两种选择实现 Accumulator.java即AccumulatorV, R最灵活的方式。它定义了两个类型参数——V表示要累加的值类型R表示最终结果类型。例如对于直方图V是数值R是整个直方图。实现 SimpleAccumulator.java即SimpleAccumulatorT适用于添加类型与结果类型相同的场景例如计数器IntCounter等内置实现正是如此。从 Accumulator.java 的源码可以看到实现该接口需要提供 5 个方法方法作用void add(V value)向累加器添加一个值R getLocalValue()获取当前 UDF 上下文中的局部值void resetLocal()重置局部值仅影响当前 UDF 上下文void merge(AccumulatorV, R other)供系统在作业结束时合并各分片结果AccumulatorV, R clone()复制累加器所有子类都必须正确实现克隆且不得抛出CloneNotSupportedException以内置的 Histogram.java 为例它是一个典型自定义累加器的参考实现add方法对每个值在TreeMap中的计数加一merge方法把另一个累加器的每个键值对合并进当前映射值相同则相加clone方法复制一份新的TreeMap。如果你认为自己实现的自定义累加器有通用价值可以考虑向 Flink 社区提交 Pull Request 使其随 Flink 一起发布。源码与测试佐证仓库中的测试用例可以直接验证上述 API 的实际用法flink-streaming-java/src/test/java/org/apache/flink/streaming/api/operators/StreamMapTest.java 中的TestOpenCloseMapFunction演示了继承RichMapFunctionString, String并覆写open/close生命周期方法的完整写法flink-streaming-java/src/test/java/org/apache/flink/streaming/api/functions/async/RichAsyncFunctionTest.java 展示了通过runtimeContext.addAccumulator(...)注册累加器的调用方式flink-streaming-java/src/test/java/org/apache/flink/streaming/runtime/tasks/OneInputStreamTaskTest.java 同样包含基于RichMapFunction的TestOpenCloseMapFunction实现可用于对照理解富函数在真实算子链中的行为。小结本文完整覆盖了 Flink DataStream API 中定义用户自定义函数的全部方式Java 侧的接口实现、匿名类、Java 8 Lambda 与富函数Scala 侧的 Lambda 与富函数并系统讲解了 Accumulator 累加器体系——从IntCounter、LongCounter、DoubleCounter、Histogram等内置实现到创建 → 注册 → 使用 → 获取结果的完整四步流程再到基于Accumulator与SimpleAccumulator接口的自定义扩展。其中富函数的open/close生命周期与getRuntimeContext()上下文能力是编写真实 Flink 作业如连接外部资源、注册累加器的关键建议在实际开发中优先使用富函数形式并结合累加器在作业结束后获取有价值的全局统计信息。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Flink Python Table API 用户自定义函数完全指南UDF、UDTF、UDAGG 与 UDTAGG 实战Flink Python Table API 用户自定义函数完全指南UDF、UDTF、UDAGG 与 UDTAGG 实战 本指南围绕 Flink PyFlin大数据流处理批处理数据工程终极开源字体解决方案PingFangSC苹方字体跨平台部署完整指南终极开源字体解决方案PingFangSC苹方字体跨平台部署完整指南 在当今多平台数字产品开发中字体一致性是设计师和开发者面临的核心挑战。PingFangSC前端3分钟快速掌握波形可视化免费交互式声波教学项目完全指南3分钟快速掌握波形可视化免费交互式声波教学项目完全指南 你是否曾好奇声音是如何传播的音乐中的和谐音是如何产生的波形可视化项目Waveforms为你提供上一篇window_manager高级教程自定义无边框窗口与标题栏的完整指南下一篇Telegraf 贡献指南Issue 与 PR 工作流、make 本地验证清单与 Execd 外部插件开发创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
企业数字化 ERP 产品动态
相关推荐
u盘安装系统全对比:3种主流方案完整示例,告别教程依赖症 u盘安装系统全对比:3种主流方案完整示例,告别教程依赖症 看了一堆教程还是不会写项目?手里拿着U盘对着电脑发呆,重启几次还是进不了系统?别急,问题不在你笨,在于那些教程只讲原理不给 完整示例… · 2026/9/23 9:19:26
alphago柯洁速查手册:3步搞定代码跑不通 alphago柯洁速查手册:3步搞定代码跑不通 代码复制完直接报错,是不是瞬间就懵了?别慌,这行干久了都知道,坑就在细节里。 很多人对着屏幕抓头发,其实缺的就是一本 速查手册 式的排错思路。 今天咱们不聊虚的,直接用 Python 拆解… · 2026/9/23 9:19:26
搞定编制军衔源码:3个完整示例彻底解决Stacktrace报错 搞定编制军衔源码:3个完整示例彻底解决Stacktrace报错 报错堆栈一屏红,StackTrace 看得人头皮发麻?别慌,这不是你代码写得烂,是“编制军衔”这块硬骨头没啃透。很多转岗做后端或系统架构的同事,一碰到这种涉及状态机、权限校验和… · 2026/9/23 10:17:17
米安考证选型指南:3个维度对比出最佳实践,避开版本坑 米安考证选型指南:3个维度对比出最佳实践,避开版本坑 版本升级后 API 全变了,是不是让你抓狂? 别再死磕旧教程了,米安体系的 最佳实践 正在重构。 今天把报名、薪资、技术栈一次讲透,让你少走三年弯路。 1.… · 2026/9/23 10:17:11
用 dnx + nuget 分发 MCP 服务:TaoToken 统一 Key 接入 .NET 10 工具链 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/23 10:17:11
健身器材出口欧盟CE认证全解析:EN 20957标准与合规要点 1. 健身器材出口欧盟的合规门槛去年帮一家国内健身器材厂搞定CE认证时,老板拿着厚达两百多页的EN 20957标准文档问我:"这些条款到底在说什么?我们生产的跑步机到底要改哪些地方?"这场景让我意识到,太多厂商困… · 2026/9/23 10:17:11
调拨单模板优化避坑指南:3个技巧提升10倍效率 调拨单模板优化避坑指南:3个技巧提升10倍效率 官方文档太长抓不住重点,导致很多开发在实现“调拨单模板”功能时,往往陷入重复造轮子的困境。别急,这份 避坑指南 直接给你可落地的代码方案,省掉你翻文档两小时的时间。… · 2026/9/23 10:17:11
3招搞定手机怎么下载微信面试难题实战项目解析 3招搞定手机怎么下载微信面试难题实战项目解析 面试被问“手机怎么下载微信”背后的原理,90%的人答不上来。别笑,这看似弱智的问题,实则是考察你对移动应用分发机制、安全校验及网络协议理解的试金石。我带过不少校招新人,他们背了八股文,却连一个A… · 2026/9/23 0:00:03
你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 面试被问“高并发下如何保证消息不丢失”,你张口就是“用Redis”,结果面试官追问“如果Redis宕机了怎么办”,你瞬间卡壳。这种场景太常见了,很多新手在背八股文时,只记住了技术名词… · 2026/9/23 0:00:29