ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

flink的架构

flink的架构 Flink 是一个面向有状态计算的分布式流处理框架也支持批处理。它的架构可以从“运行时组件、数据流模型、容错机制”三个层面理解。一、整体架构一个典型的 Flink 集群包含以下组件客户端 Client | | 提交 Job v 作业管理器 JobManager | | 调度任务、协调检查点、管理元数据 v 任务管理器 TaskManager 集群 | | 执行算子、处理数据、保存本地状态 v 外部系统Kafka / MySQL / Elasticsearch / HDFS / OSS1. ClientClient 负责解析用户编写的 Flink 程序构建 JobGraph将作业提交给 JobManager提交后可以退出作业仍由集群继续运行Client 通常不是长期运行时组件。2. JobManagerJobManager 是作业的控制中心主要负责接收和管理作业将作业转换为可执行任务调度任务到 TaskManager管理 Checkpoint 和 Savepoint处理故障恢复管理作业状态和资源需求在现代 Flink 架构中JobManager 内部通常可以理解为几个逻辑角色Dispatcher接收 REST、命令行等提交请求并启动 JobMasterJobMaster负责单个作业的调度和协调ResourceManager管理集群资源并向外部资源管理系统申请资源3. TaskManagerTaskManager 是真正执行计算任务的工作节点主要负责执行算子逻辑处理输入和输出数据管理网络数据交换保存算子状态向 JobManager 汇报任务状态一个 TaskManager 中可以有多个Task Slot。Slot 是 Flink 对 TaskManager 资源进行逻辑隔离和分配的单位。二、Flink 的作业执行模型用户程序通常经过以下转换用户代码 - StreamGraph - JobGraph - ExecutionGraph - TaskManager 上的执行任务1. StreamGraphStreamGraph 是对用户逻辑的直接描述。例如Source - Map - Filter - Sink它包含 Source、Transformation、Sink 等节点。2. JobGraphJobGraph 是提交给 JobManager 的作业描述。Flink 会在这一阶段进行算子链合并等优化把可以连续执行的算子合并为一个 JobVertex。例如Source - Map - Filter如果它们之间没有发生数据重分区可能会被合并到同一个 Operator Chain 中从而减少线程切换和网络通信。3. ExecutionGraphExecutionGraph 是 JobManager 根据并行度和资源情况生成的实际执行计划。它会把 JobGraph 中的逻辑节点展开为多个并行实例。例如并行度为 3Map 算子 Map-0 Map-1 Map-24. Operator Chain多个连续算子可以在同一个线程中执行形成 Operator ChainSource - Map - Filter这样可以避免不必要的序列化、网络传输和线程切换提升性能。如果两个算子之间存在 Shuffle、KeyBy、广播或并行度变化通常不能直接链在一起。三、Task、SubTask 和 Slot假设一个作业如下Source - Map - KeyBy - Reduce - Sink并行度为 2 时每个算子会产生两个 SubTaskSource-0 - Map-0 - Reduce-0 - Sink-0 Source-1 - Map-1 - Reduce-1 - Sink-1其中Operator算子的逻辑定义SubTask算子按照并行度展开后的一个实例Task一个或多个 Operator Chain 的执行单元Task SlotTaskManager 提供的逻辑资源槽位需要注意Slot 主要用于资源调度和隔离并不等于一个独立线程也不一定对应一个 SubTask。四、数据流模型Flink 的核心是连续不断的数据流。即使处理有限数据集也可以看作一个有界流。Source - Transformation - Sink常见组件如下Source从 Kafka、文件、数据库等读取数据Transformation执行 Map、Filter、Join、Window 等计算Sink将结果写入数据库、消息队列或文件系统KeyBy会按照 Key 对数据进行逻辑分区使同一个 Key 的数据进入同一个下游并行实例相同 userId 的数据 | v 同一个 Reduce SubTask这为按 Key 维护状态提供了基础。五、状态管理Flink 与普通无状态计算框架的一个重要区别是它原生支持大规模、有一致性保障的状态管理。状态通常分为Keyed State绑定到某个 Key例如每个用户的累计金额Operator State绑定到算子实例例如 Kafka Source 的分区消费进度Flink 的状态可以存储在堆内存RocksDB 等嵌入式状态后端其他状态后端实现状态后端负责状态的实际存储而 Checkpoint 负责定期持久化和恢复状态。六、Checkpoint 容错机制Flink 使用基于 Chandy-Lamport 思想的分布式快照机制实现一致性 Checkpoint。简化流程如下JobManager 发出 Checkpoint Barrier | v Source 注入 Barrier | v Barrier 随数据流向下游传播 | v 各算子保存自己的状态 | v Checkpoint 完成并写入持久化存储Checkpoint 保存的内容通常包括算子状态Source 消费位点其他恢复所需的运行时信息任务失败后Flink 可以从最近一次成功的 Checkpoint 恢复状态恢复 Source 的消费位置重新执行之后的数据继续处理作业常见语义包括At-most-once最多一次可能丢数据At-least-once至少一次可能重复数据Exactly-once端到端恰好一次但需要 Source、Flink 和 Sink 协同支持七、网络与数据交换TaskManager 之间通过网络传输数据。典型的数据交换方式包括Forward上游一个 SubTask 对应下游一个 SubTaskRebalance轮询分发数据Broadcast发送给所有下游实例KeyBy / Hash Partition按 Key 哈希分区Rescale在本地范围内重新分配数据keyBy()往往会产生网络 Shuffle是数据流发生重新分区的关键位置。TaskManager 内部通常通过网络缓冲区、Result Partition 和 Input Gate 组织数据传输。八、部署模式Flink 常见的部署方式有1. Session Cluster预先启动一个 Flink 集群多个作业共享这个集群。优点启动快资源可以复用。缺点作业之间可能互相影响资源隔离较弱。2. Per-Job Cluster每个作业单独创建一个集群作业结束后集群释放。优点作业隔离较好。缺点启动和资源申请成本较高。3. Application Mode将应用程序和 Flink 集群一起部署由集群直接执行应用入口逻辑。常见运行环境包括StandaloneYARNKubernetes在 Kubernetes 环境中JobManager 通常运行在一个或多个 Pod 中TaskManager 以 Pod 形式按需扩缩容。九、Flink 与传统 Lambda 架构的区别传统 Lambda 架构一般将系统拆成Batch Layer Speed Layer Serving Layer同一份数据可能需要分别走批处理和流处理链路。Flink 更强调统一的流批处理模型有界流 无界流无界流适合实时计算有界流适合批处理两者可以共享大量 API、状态和执行机制十、一次作业运行的完整过程可以把 Flink 作业运行概括为1. 用户编写 DataStream 或 Table/SQL 程序 2. Client 生成作业图 3. Client 将作业提交给 JobManager 4. JobManager 申请资源并生成执行计划 5. TaskManager 启动各个 SubTask 6. Source 持续读取数据 7. 数据经过算子链和网络 Shuffle 流转 8. 算子维护本地状态 9. JobManager 定期触发 Checkpoint 10. 发生故障时从 Checkpoint 恢复 11. 结果通过 Sink 写入外部系统一句话总结JobManager 负责控制和调度TaskManager 负责执行Operator 描述计算逻辑State 保存中间结果Checkpoint 提供一致性容错网络 Shuffle 负责并行实例之间的数据交换。
返回列表