ARTICLE DETAIL

资讯详情

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

Apache Flink 学习指南:从流处理、事件时间到有状态容错的四大核心概念

Apache Flink 学习指南:从流处理、事件时间到有状态容错的四大核心概念 大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载本指南是 Learn Flink 动手实践培训 的开篇围绕 Apache Flink 流式编程最核心的四个概念展开无界流处理、事件时间、有状态流处理与状态快照。读完本文你将掌握 Flink 应用由 Source、算子、Sink 构成的流式数据流的基本形态理解并行与数据重分发的原理并能在源码层面讲清Flink 凭什么做到精确一次exactly-once容错。培训的目标与范围这套培训旨在用最少的篇幅让你上手编写可扩展的流式 ETL、实时分析与事件驱动应用刻意省略了大量虽然最终重要的细节将重点放在Flink 管理状态与时间的两套 API 上如何实现流式数据处理管道pipelineFlink 为什么以及如何管理状态如何用事件时间持续、一致地计算准确的分析结果如何在连续流上构建事件驱动应用Flink 如何以精确一次语义提供容错的有状态流处理。本页只负责引入四个核心概念流式数据的连续处理、事件时间、有状态流处理、状态快照。每一个小节结尾都会给出深入学习与配套动手练习的指引。流处理数据天然生活在流里无论是 Web 服务器的日志事件、证券交易所的成交记录还是工厂车间里机器的传感器读数数据在诞生之时就是以流的形式存在的。分析数据时你可以围绕**有界bounded或无界unbounded**的流来组织处理逻辑选择哪种范式会产生深远影响。批处理Batch processing处理的是有界数据流。你可以先完整摄入整个数据集再产生任何结果因此可以对数据排序、计算全局统计量、或者生成一份汇总全部输入的最终报表。流处理Stream processing面对的是无界数据流。从概念上讲输入永远不会结束你必须随到随处理。在 Flink 中应用由可以被用户自定义**算子operators转换的流式数据流streaming dataflows**组成。这些数据流构成有向图从一个或多个Source数据源开始以一个或多个Sink数据出口结束。通常情况下程序中每个转换transformation与数据流中的每个算子是一一对应的但有时一个转换也会由多个算子组成。一个 Flink 应用既可以消费来自消息队列或分布式日志如 Apache Kafka、Kinesis的实时数据也可以从多种数据源消费有界的历史数据产出的结果流同样可以发送到各种可作为 Sink 连接的系统。并行数据流Flink 程序天生是并行、分布式执行的。执行期间一条流有一个或多个流分区stream partitions每个算子有一个或多个算子子任务operator subtasks。子任务彼此独立运行在不同的线程中甚至可能分布在不同的机器或容器上。算子子任务的个数就是该算子的并行度parallelism同一个程序里不同算子可以拥有不同的并行度。两条算子之间的数据传输有两种模式一对一one-to-one或 forwarding保留元素的分区与顺序。例如上图中 Source 与map()算子之间map()的 subtask[1] 看到的元素与 Source 的 subtask[1] 产出的元素顺序完全一致。重分发redistributing改变流的分区方式。算子子任务会根据所选的转换向不同的目标子任务发送数据典型例子有keyBy()按 key 哈希重新分区、broadcast()广播和rebalance()随机轮询重分区。在重分发交换中元素之间的顺序只在每一对发送子任务—接收子任务之间被保留例如map()的 subtask[1] 与keyBy/window的 subtask[2] 之间。因此上图keyBy/window与 Sink 之间的重分发会引入不同 key 的聚合结果到达 Sink 时的顺序不确定性。从源码看keyBy与时间戳分配等核心算子都定义在 DataStream.java 中例如keyBy(KeySelector)负责按键重分区assignTimestampsAndWatermarks(...)负责注入水位线watermark——这正是下一节及时流处理的编程入口。及时流处理让历史数据与实时数据共用一套代码对大多数流式应用而言能用处理实时数据的同一套代码重放历史数据并且无论怎样都产出确定、一致的结果是极具价值的。同时关注事件发生的顺序而非事件被投递处理的顺序并能推断一组事件何时应该完整也至关重要——比如一笔电商交易或一笔金融成交所涉及的一串事件。这些及时流处理的需求可以通过使用**记录在数据流中的事件时间时间戳event time timestamps**来满足而不是使用处理数据的机器时钟。事件时间让迟到、乱序的数据也能被正确归并到所属的时间窗口从而保证分析结果的一致性围绕assignTimestampsAndWatermarks注入时间戳与水位线、再用window算子开窗正是 streaming_analytics.md 与 event_driven.md 两篇后续章节要深入的主题。有状态流处理分布式环境下的分片键值存储Flink 的算子可以是有状态的一个事件如何处理取决于它之前所有事件累积产生的影响。状态可以很简单——比如每分钟统计事件数展示在仪表盘上也可以很复杂——比如为欺诈检测模型计算特征。一个 Flink 应用在分布式集群上并行运行同一算子的多个并行实例彼此独立执行、运行在不同的线程甚至不同的机器上。因此一组有状态算子的并行实例本质上是一个分片的键值存储sharded key-value store。每个并行实例负责处理特定一组 key 的事件这些 key 的状态保存在本地。下图是一个并行度示例作业图前三个算子并行度为 2最后的 Sink 并行度为 1。第三个算子是有状态的第二与第三个算子之间发生了全连接fully-connected的 shuffle 网络重分区目的是按某个 key 切分流让所有需要一起处理的事件被送到同一个实例。状态始终在本机访问这让 Flink 应用得以获得高吞吐与低延迟。你可以选择把状态放在 JVM 堆上也可以当状态过大时放入组织高效、位于磁盘上的数据结构中本地访问、按 key 分片这一设计落实在源码上就是 flink-core-api/src/main/java/org/apache/flink/api/common/state/ 目录下的一系列状态接口ValueStateT单值状态、ListStateT列表状态、MapStateK, V映射状态、ReducingStateT归约状态、AggregatingStateIN, OUT聚合状态等。ValueState的状态值本身是存储在 StateBackend 中的字节数据读取与写入都要经过状态后端的序列化/反序列化与本地访问路径。状态后端堆内存还是磁盘Flink 管理的状态存放在**状态后端state backend**中当前提供两种实现详见 fault_tolerance.md名称工作状态存放位置快照方式EmbeddedRocksDBStateBackend本地磁盘临时目录全量 / 增量HashMapStateBackendJVM 堆全量EmbeddedRocksDBStateBackend内嵌的 RocksDB 键值存储工作状态在磁盘上。支持超过可用内存大小的状态经验法则比堆内存后端慢约 10 倍。HashMapStateBackend工作状态在 JVM 堆内存中。速度快但需要足够大的堆受 GC 影响。在堆内存后端中状态的访问与更新是直接读写堆上的对象而在EmbeddedRocksDBStateBackend中访问与更新涉及序列化与反序列化代价更高但可容纳的状态大小只受本地磁盘限制。另外只有 RocksDB 后端支持增量快照这对状态巨大且变化缓慢的应用是显著优势。两种后端都支持异步快照——即在不阻塞流处理的前提下拍照对应源码中的 HashMapStateBackend 与 EmbeddedRocksDBStateBackend。通过状态快照实现容错精确一次的秘密Flink 通过状态快照state snapshots与流重放stream replay的组合提供容错且精确一次exactly-once的语义。快照捕获分布式管线的完整状态既记录输入队列中的 offset也记录作业图中截至此刻已摄入数据所产生的全部状态。一旦发生故障Source 被回卷、状态被恢复、处理得以继续。如前文所述这些状态快照是异步捕获的不会阻塞正在进行的处理。快照的持久化检查点存储Flink 会周期性地对每个算子的全部状态做持久化快照并复制到更持久的存储如分布式文件系统。存放位置由作业的**检查点存储checkpoint storage**决定同样有两种实现名称状态备份位置FileSystemCheckpointStorage分布式文件系统JobManagerCheckpointStorageJobManager JVM 堆FileSystemCheckpointStorage支持非常大的状态、持久性高生产部署推荐JobManagerCheckpointStorage适合本地测试与小状态实验。快照与检查点相关术语快照Snapshot泛指 Flink 作业状态的一个全局一致的镜像包含每个数据源的指针如文件或 Kafka 分区的 offset以及各状态算子处理到该位置所产生的全部状态副本。检查点CheckpointFlink 为故障恢复目的自动触发的快照。检查点可以是增量的并被优化为可快速恢复。外部化检查点Externalized Checkpoint默认情况下检查点并非给用户操作使用——作业运行期间 Flink 只保留最近 n 个n 可配置作业取消时删除但可以配置为保留从而支持手动从检查点恢复。保存点Savepoint由用户或 API 调用手动触发、服务于运维目的的快照例如有状态的重部署/升级/扩缩容。保存点总是完整的并被优化为便于运维灵活操作。状态快照如何工作异步屏障快照Flink 使用Chandy-Lamport 算法的一个变体称为异步屏障快照asynchronous barrier snapshotting检查点协调器checkpoint coordinator属于 JobManager 的一部分指示 TaskManager 开始一次检查点所有 Source 记录自身 offset并向流中插入带编号的检查点屏障checkpoint barrier这些屏障流经整个作业图标出每个检查点前后的流区间。检查点 n 将包含每个算子在消费了屏障 n 之前的所有事件、且不含其之后任何事件时的状态。每个算子收到屏障后记录自身状态具有两个输入流的算子如CoProcessFunction会执行屏障对齐barrier alignment确保快照反映的是两个输入流都消费到各自屏障且未越过时的状态。Flink 的状态后端使用**写时复制copy-on-write**机制在旧版本状态被异步快照期间流处理可以畅通无阻地继续只有快照被持久化后旧版本状态才会被垃圾回收。源码层面的佐证很清晰屏障以CheckpointBarrier这一运行时事件的形式在网络栈中传播CheckpointBarrier.java而周期性触发快照的调度与协调逻辑集中在 CheckpointCoordinator.java其triggerCheckpoint(boolean isPeriodic)等方法负责发起检查点并跟踪完成情况。精确一次的三种语义与端到端前提流处理应用出问题时可能丢失结果、也可能产生重复结果。取决于应用与集群的配置Flink 可能表现为At most once最多一次Flink 不做恢复努力At least once至少一次什么都不丢但可能出现重复结果Exactly once精确一次既不丢也不重复。由于 Flink 靠回卷并重放源数据流来恢复所谓exactly once 并不是说每个事件都被恰好处理一次而是每个事件对 Flink 所管理状态的影响恰好发生一次。屏障对齐只为提供精确一次保证所需如果不需要可以通过配置CheckpointingMode.AT_LEAST_ONCE禁用屏障对齐来换取性能。而要做到端到端end-to-end精确一次源头的每个事件恰好影响 Sink 一次必须同时满足你的 Source可重放replayable你的 Sink是事务性的或幂等的transactional or idempotent。下一步动手练习与学习路线本页只是起点。完整培训路径按以下顺序展开均位于 docs/content/docs/learn-flink/DataStream API 入门——理解可流式化的数据类型基本类型、Tuple、POJO、Scala case class、Kryo 兜底与 Avro 支持、完整示例、StreamExecutionEnvironment、基本 Source/Sink 与 IDE 内调试数据管道与 ETL——map()、flatmap()等无状态转换再到有状态转换与窗口流式分析——事件时间、水位线与窗口聚合事件驱动应用——连续流上的事件驱动编程容错——状态后端、检查点存储、快照机制与精确一次保证的完整讲解。每节末尾都附有配套的动手练习链接。实践环节可参考 Flink Operations Playground 中的Observing Failure Recovery一节亲身体验故障注入后 Flink 如何基于状态快照自动恢复。掌握了本文的四个核心概念再去阅读更详细的参考文档如 Checkpointing 配置、Checkpoints、Savepoints、Large State Tuning会事半功倍。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Apache Flink 流处理核心概念与实践训练从流式 Dataflow 到有状态容错Apache Flink 流处理核心概念与实践训练从流式 Dataflow 到有状态容错 导读 本文基于 Apache Flink 官方实践训练Hands大数据流处理批处理数据工程Apache Flink 编程模型概念透析从有状态流处理到 SQL 的四层 API 抽象Apache Flink 编程模型概念透析从有状态流处理到 SQL 的四层 API 抽象 Flink 为流式/批式处理应用程序的开发提供了从底层有状态流处理到大数据流处理批处理数据工程Apache Flink 核心概念术语全解从集群架构、数据流模型到状态与容错机制Apache Flink 核心概念术语全解从集群架构、数据流模型到状态与容错机制 导读 Flink 的官方文档、源码注释与社区讨论中充斥着大量术语——Flin大数据流处理批处理数据工程上一篇终极指南如何使用Neural Amp Modeler插件获得专业吉他音色下一篇qui终极指南一键部署qBittorrent现代化管理界面创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表