ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Flink 编程模型概念透析:从有状态流处理到 SQL 的四层 API 抽象

Apache Flink 编程模型概念透析:从有状态流处理到 SQL 的四层 API 抽象 大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载Flink 为流式/批式处理应用程序的开发提供了从底层有状态流处理到顶层 SQL 查询的四级编程抽象。本文以 Flink 官方文档概念透析Concepts in Depth概览章节为主线深入剖析每一层抽象的定位、能力边界与适用场景并结合 Flink 仓库源码ProcessFunction.java、Table.java 等验证其底层实现机制。读完本文你将能够根据业务场景在四层抽象之间做出正确选型并理解各层 API 之间如何无缝切换与协同。说明本文对应仓库文档 overview.md属于概念透析Concepts in Depth章节的导览。该章节深入分析 Flink 分布式运行时架构如何实现有状态与及时流处理等核心概念。阅读路线本章在 Flink 概念体系中的位置Flink 官方文档将概念学习拆分为两个层次实践练习Hands-on Training介绍作为 Flink API 根基的有状态实时流处理基本概念并举例说明如何在应用中使用这些机制。对应仓库文档 learn-flink/overview.md 下的三个子章节Data Pipelines ETL 中的Stateful Transformations小节引入了有状态流处理的概念Fault Tolerance 对有状态流处理进行了深入展开检查点、状态恢复等Streaming Analytics 介绍了及时流处理事件时间、水印、窗口的概念。概念透析Concepts in Depth即本章concepts/overview.md深入分析 Flink 分布式运行时架构如何实现上述概念。本文聚焦本概览章节的核心主题——Flink 的四层 API 抽象从最底层到最顶层逐一展开。总览Flink 的四层 API 抽象Flink 为流式/批式处理应用程序的开发提供了不同级别的抽象自底向上依次为有状态实时流处理Stateful Timely Stream Processing——最底层抽象Core APIsDataStream API / DataSet API——通用编程接口层Table API——以表为中心的声明式 DSLSQL——最顶层、与 Table API 语义紧密关联的查询语言。下图为官方文档中展示的编程抽象层级示意图片来源docs/static/fig/levels_of_abstraction.svg另见 docs/static/fig/concepts/levels_of_abstraction.svg抽象层级越高表达越简洁、声明性越强层级越低表达力越强、控制粒度越细。下面逐一深入。第一层有状态实时流处理底层抽象Flink API 最底层的抽象为有状态实时流处理其抽象实现是Process Function并且被 Flink 框架集成到了 DataStream API 中使用。Process Function 的能力从源码看ProcessFunction.java 定义了一个处理流的函数抽象核心回调方法包括processElement(I value, Context ctx, CollectorO out)对输入流中的每一个元素被调用可产出零个或多个输出元素onTimer(long timestamp, OnTimerContext ctx, CollectorO out)当通过TimerService注册的定时器触发时被调用同样可产出输出并继续注册新定时器Context提供timestamp()当前元素/触发定时器的时间戳、timerService()查询时间与注册定时器以及output(OutputTag, value)向侧输出流发射记录等能力OnTimerContext额外提供timeDomain()用于区分定时器所属的时间域事件时间或处理时间。它允许用户自由地处理来自单流或多流的事件数据并提供具有全局一致性和容错保障的状态。此外用户可以在此层抽象中注册**事件时间event time和处理时间processing time**回调方法从而允许程序实现复杂计算。关键实现细节Keyed 状态访问根据 ProcessFunction.java 中的注释只有将ProcessFunction应用在KeyedStream上时才能访问 Keyed State 和定时器两者都作用域于某个 key。Rich 函数特性ProcessFunction继承自AbstractRichFunction因此永远是一个RichFunction可以通过open(OpenContext)/close()生命周期方法获取RuntimeContext见 ProcessFunction.java。Keyed 变体对于需要按 key 处理的场景Flink 提供了KeyedProcessFunctionKeyedProcessFunction.java其processElement中的Context额外暴露getCurrentKey()等方法。使用示意伪代码形态具体 API 详见 DataStream API 文档DataStreamEvent stream env.addSource(...); stream .keyBy(Event::getKey) .process(new KeyedProcessFunctionString, Event, Output() { Override public void processElement(Event value, Context ctx, CollectorOutput out) { // 访问 keyed state、查询时间、注册定时器 ctx.timerService().registerEventTimeTimer(value.getTimestamp() 60000); } Override public void onTimer(long timestamp, OnTimerContext ctx, CollectorOutput out) { // 定时器触发逻辑 } });第二层Core APIsDataStream API 与 DataSet API实际上许多应用程序并不需要用到最底层抽象而是可以直接使用Core APIs进行编程。DataStream APIDataStream API应用于有界/无界数据流场景是这一层的主力。Core APIs 提供的流式 APIFluent API即链式调用风格为数据处理提供了通用的模块组件包括但不限于各种形式的用户自定义转换transformations联接joins聚合aggregations窗口windows**状态state**操作。此层 API 中处理的数据类型在每种编程语言中都有其对应的类例如 Java 中的 POJO、Scala 中的 case class 等。相关的 DataStream API 文档位于 docs/content.zh/docs/dev/datastream/overview.md。底层抽象与 Core API 的集成Process Function这类底层抽象和DataStream API相互集成使得用户可以选择按需使用更底层的抽象 API 来实现自己的需求。典型用法是在 DataStream 上调用.process(...)/.keyBy(...).process(...)从而在普通转换中插入有状态、定时器、侧输出等底层能力。DataSet API有界批处理DataSet API额外提供了一些面向有界数据集的原语比如循环/迭代loop/iteration操作适用于批式迭代计算场景如图算法、机器学习迭代训练等。注在当前的 Flink 版本中批处理能力正在统一收敛到 DataStream API 之上流批一体DataSet API 属于较早期的有界批处理编程接口。第三层Table API声明式 DSLFlink API 第三层抽象是Table API它是以**表Table**为中心的声明式编程DSLAPI。核心特性动态表在流式数据场景下Table 可以表示一张正在动态改变的表dynamic table遵循扩展关系模型表拥有 schema类似于关系型数据库中的 schema关系型操作提供类似于关系模型中的操作如select、project、join、group-by和aggregate等。从源码可以验证这些操作的真实存在Table.java 中定义了Table select(Expression... fields)第 129 行Table filter(Expression predicate)第 199 行GroupedTable groupBy(Expression... fields)第 234 行Table join(Table right)、Table join(Table right, Expression joinPredicate)第 262/284 行Table joinLateral(...)表函数 lateral join第 403/437 行等。声明式与表达力权衡Table API 程序以声明的方式定义应执行的逻辑操作而不是确切地指定程序应该执行的代码——即表达做什么what而非怎么做how。尽管 Table API 使用起来很简洁并且可以由各种类型的用户自定义函数UDF扩展功能但它的表达能力还是比 Core API 差。此外Table API 程序在执行之前会经过优化器应用优化规则对用户编写的表达式进行优化例如谓词下推、列裁剪、算子重排等这显著减少了用户手写优化的工作量。与 DataStream/DataSet 的无缝切换表和DataStream/DataSet可以进行无缝切换Flink 允许用户在编写应用程序时将Table API与DataStream/DataSetAPI 混合使用。典型场景先使用 DataStream 做复杂的底层处理再转成 Table 做关系型分析或者先用 Table 做声明式清洗再转回 DataStream 做精细控制。相关 API 文档位于 docs/content.zh/docs/dev/table/overview.md。第四层SQL最顶层抽象Flink API 最顶层抽象是SQL。这层抽象在语义和程序表达力上都类似于Table API但程序实现全部是SQL 查询表达式SQL 抽象与 Table API 抽象之间的关联非常紧密SQL 查询语句可以直接在Table API中定义的表上执行换句话说Table API 构建的 Table 对象就是 SQL 查询的数据源二者共享同一套关系模型与优化器。SQL 抽象相关的说明位于 docs/content.zh/docs/dev/table/overview.md 的 SQL 小节。使用形态示例在 TableEnvironment 中执行 SQL// 在 Table API 中注册或创建表 tableEnv.createTemporaryView(orders, ordersTable); // 在 Table API 定义的表上执行 SQL 查询 Table result tableEnv.sqlQuery( SELECT user_id, COUNT(*) AS cnt FROM orders GROUP BY user_id);四层抽象的选型建议抽象层级核心形态表达风格表达能力典型适用场景有状态实时流处理Process Function命令式回调最强状态、定时器、侧输出复杂事件处理、精确时间语义控制、底层算子定制Core APIsDataStream / DataSet APIFluent 命令式强通用流/批处理、窗口聚合、连接、状态编程Table API表 声明式 DSL声明式what中等可用 UDF 扩展关系型分析、需要优化器的场景SQLSQL 查询表达式声明式与 Table API 相当数据分析师、标准 SQL 场景、与 BI 工具对接选型思路需要细粒度控制时间、状态与输出→ 下沉到 Process Function 层常规流/批数据处理→ 使用 DataStream API有界/无界统一希望编写更少代码、享受优化器红利→ 使用 Table API面向标准 SQL 或非程序员协作→ 使用 SQL并可与 Table API 自由混用。与运行时架构的关系本概览章节属于概念透析Concepts in Depth系列的开篇其姊妹章节深入剖析了支撑上述 API 的运行时机制均可在 docs/content.zh/docs/concepts 目录下找到flink-architecture.mdFlink 分布式运行时架构stateful-stream-processing.md有状态流处理的底层实现time.md时间概念事件时间/处理时间/水印glossary.md概念术语表。无论使用哪一层 API 编程最终都会被编译为分布在 TaskManager 上执行的有状态算子图operator graph由 JobManager 负责调度与容错协调——这正是概念透析后续章节要回答的架构如何实现概念这一问题的核心。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐XState核心概念深度解析从有限状态机到Actor模型XState核心概念深度解析从有限状态机到Actor模型 本文深入探讨XState状态管理库的核心架构从基础的有限状态机原理出发系统解析层次状态机设计、并前端后端从批处理到实时计算Apache Flink SQL如何重构现代数据处理流程从批处理到实时计算Apache Flink SQL如何重构现代数据处理流程 你是否还在为传统ETL工具的批处理延迟而烦恼当业务需要实时决策时每天凌晨运行的大数据流处理批处理数据工程pytype 抽象值系统解析从概念到实现的终极指南pytype 抽象值系统解析从概念到实现的终极指南 想要深入理解 Python 静态类型检查的核心机制吗pytype 的抽象值系统正是其类型推断能力的静态分析开发工具代码质量创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表