ARTICLE DETAIL

资讯详情

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

RisingWave 流式引擎架构全解:Actor 模型、共享状态存储与 Barrier 一致性机制

RisingWave 流式引擎架构全解:Actor 模型、共享状态存储与 Barrier 一致性机制 数据库流处理后端数据工程【免费下载链接】risingwaveEvent streaming platform for agentic AI. Continuously ingest, transform, and serve event streams in real time, at scale.项目地址https://gitcode.com/gh_mirrors/ri/risingwave点击查看免费下载本文基于 RisingWave 仓库内设计文档 docs/dev/src/design/streaming-overview.md 展开。RisingWave 是一个面向实时分析的事件流处理平台其核心能力在于用户只需定义物化视图Materialized ViewMV系统便会自动根据最新数据更新视图使查询始终反映实时分析结果——这一自动刷新工作正是由 RisingWave 流式引擎streaming engine完成的。本篇技术指南将系统讲解 RisingWave 流式引擎的整体架构、三大核心设计原则Actor 模型执行引擎、共享存储状态、一切皆表、一切皆状态、从 CREATE MATERIALIZED VIEW 到流式管线就绪的完整链路以及保证一致性与故障恢复的 Barrier 检查点机制。读完本文你将理解 RisingWave 如何在分布式环境下高效维护任意 SQL 物化视图并能结合仓库源码进一步深入各模块实现。设计总览流式引擎的三大核心原则RisingWave 通过定义物化视图MV来提供实时分析能力所有物化视图都会随最新数据更新自动刷新查询物化视图即可获得实时的分析结果。负责这一刷新任务的就是本文的主角——RisingWave 流式引擎。其核心设计原则可概括为三条基于 Actor 模型的执行引擎Actor model based execution engine。系统创建一组 Actor每个 Actor 响应自己的输入消息包括数据更新与控制信号。以此构建一个高并发、高效率的流式引擎。从源码看Actor 的执行模型实现于 src/stream/src/executor/actor.rs其运行循环依托 tokio 异步运行时。状态共享存储Shared storage for states。状态存储的骨干基于共享云对象存储当前为 AWS S3这带来计算弹性、廉价且近乎无限的存储容量以及在配置变更时的简洁性。对应的存储引擎 Hummock 详见 docs/dev/src/design/state-store-overview.md。一切皆表、一切皆状态Everything is a table, everything is a state。内部存储中的每个对象既是一张逻辑表也是内部状态。因此它们可以被 catalog 有效管理并在统一流式引擎中以一致性保证进行更新。这三条原则分别回答了如何算Actor 并发模型、状态放哪共享云存储和状态怎么管统一为表 catalog 治理三个核心问题是理解后文架构与一致性的钥匙。系统架构前端、计算节点与元数据服务RisingWave 流式引擎由三组节点构成前端节点Frontend由服务层serving layer构成负责并发处理用户的 SQL 请求。它是 SQL 解析、计划生成和查询服务的入口。计算节点Compute nodes构成处理层processing layer。每个计算节点承载一组长期运行的 Actor 用于流处理所有 Actor 共享一个持久化存储层当前为 AWS S3作为各自的状态存储。元数据服务Meta service维护全部元信息并协调整个集群包括流图切分、调度、检查点协调等职责。其设计细节可参考 docs/dev/src/design/meta-service.md。从 CREATE MATERIALIZED VIEW 到流式管线四步构建流程当前端收到一条CREATE MATERIALIZED VIEW语句时物化视图及对应的流式管线按下述四步构建构建流计划Building a stream plan流计划是由逻辑算子编码数据流dataflow的逻辑计划由前端节点的流式规划器streaming planner完成。切分Fragmentation元数据服务上的流切分器stream fragmenter将生成的逻辑流计划切成若干流片段stream fragment并在必要时对片段进行复制。一个流片段持有流计划中的部分节点每个片段可以通过构建多个 Actor 实现数据并行。源码位于 src/meta/src/stream/stream_graph。调度流片段Scheduling plan fragments元数据服务将不同片段分发到不同计算节点并让所有计算节点构建各自的本地 Actor。调度相关实现可见 src/meta/src/stream/stream_graph/schedule.rs。后端初始化作业Initializing the job at the backend元数据服务通知所有计算节点开始运行流式管线。可以推断该流程中逻辑流计划 → 片段 → Actor的逐级物化正是分布式并行度每个片段可并行出多个 Actor的根源关于并行度如何通过streaming_parallelism系列会话参数控制可参见 docs/dev/src/design/streaming-parallelism.md。Actor、Executor 与状态Actor可被调度的最小单元Actor 是流式引擎中可被调度的最小单元Actor 内部不再并行。一个典型 Actor 由三部分组成Merger可选将来自不同上游 Actor 的消息合并到单一通道使下游 Executor 能顺序处理消息同时负责对齐 Barrier 以支持检查点详见后文。Executor 链Chain of executors每个 Executor 是增量计算delta computation的基本单元。Dispatcher可选按特定策略如哈希打散 hash shuffling 或轮询 round-robin将收到的消息发送给不同下游 Actor。Actor 的执行由 tokio 异步运行时承载。Actor 启动后运行一个无限循环持续执行异步函数生成输出直到收到停止消息。相关实现见 src/stream/src/executor/actor.rs。消息传递方面两个本地 Actor 之间通过 channel 传输消息对于位于不同计算节点上的 Actor消息会被重定向到交换服务exchange service交换服务之间通过 RPC 请求持续交换消息。Executor增量计算的基本单元Executor 是流式引擎中的基本计算单元。每个 Executor 响应收到的消息并原子地计算输出消息即每个 Executor 内部的计算不会被进一步拆解。RisingWave 流式系统底层的算法框架是传统的变更传播框架change propagation framework。给定一个待维护的物化视图系统构建一组 Executor每个 Executor 对应一个关系算子含基表。当任一基表收到更新时流式引擎从叶子节点到根节点递归计算每个物化视图的变更每个节点从某个子节点收到更新计算局部更新并传播给父节点。由于保证每个 Executor 的正确性便得到一个可组合composable的框架用以维护任意 SQL 查询。从源码结构看各算子对应 Executor 位于 src/stream/src/executor 目录例如filter.rs过滤、hash_join.rs哈希连接、hash_agg.rs哈希聚合、top_n.rsTopN、merge.rs对应 Merger与dispatch.rs对应 Dispatcher等。一致性、检查点与故障恢复一致性模型这里用一致性表示查询物化视图的完整性与正确性模型。系统保证查询结果始终是查询发起时间戳之前某个时间点 t 的一致快照consistent snapshot后续查询总是从更晚的时间点取得一致快照。所谓 t 时刻的一致快照要求所有不晚于 t 的消息都被恰好一次地反映在快照中而所有晚于 t 的消息都未被反映。更形式化的描述一致性 单调性可参见 docs/dev/src/design/checkpoint.md——该文档同时指出 RisingWave不保证读己之写read-after-write一致性用户可通过FLUSH命令确保变更在读取前已生效。基于 Barrier 的检查点Barrier based checkpoint为保证一致性RisingWave 引入 Chandy–Lamport 风格的分布式一致性快照算法作为检查点方案。该过程保证被刷入存储的每个状态是一致的与源端某个 Barrier 对应。因此当批处理引擎读取存储上视图和表的一致快照时一致性天然得到保证。每个 Barrier 也称为一个 epoch两者经常互换使用因为数据流被切分为若干 epoch。换言之对数据库的写入只有在通过检查点提交到存储后才可见。检查点详细流程见 docs/dev/src/design/checkpoint.md为元数据服务定期间隔由barrier_interval_ms配置初始化 Barrier 并广播给所有源 ActorBarrier 沿流图流经每个算子——对Dispatch等扇出算子复制给所有下游对Merge、Join等扇入算子收集齐所有上游的 Barrier 后再发出其余算子则触发一次检查点操作将变更刷入存储当计算节点所有脏状态都刷完后节点向元数据服务发送完成信号元数据服务收到所有节点的完成信号后通知存储管理器提交检查点。为提高效率同一计算节点上的所有脏状态先汇入一个共享缓冲区shared buffer计算节点再将整个共享缓冲区异步刷成存储中的单个 SST 文件使检查点过程不会阻塞流处理。共享缓冲区的引入还带来额外收益一个计算节点内的写批次可先压缩成单个 SSTable 再上传显著减少 L0 层 SSTable 文件数量。结合存储层视角看Executor 收到 Barrier 后即进入新 epochepoch 1 的数据以 epoch 1 写入收到 Barrier 也会使读写 epoch 切换到新值只有等元数据服务收集到下一 Barrier如 epoch 2后才能确认前一个 epochepoch 1的数据已全部写入共享缓冲区并启动其检查点。详见 docs/dev/src/design/state-store-overview.md。故障恢复当流式引擎崩溃时系统必须整体回滚到之前某个一致快照。为此只要元数据服务检测到某些计算节点故障failover或某个检查点过程正在进行便会触发恢复流程recovery process在重建流式管线后每个 Executor 从存储中的一致快照重置本地状态恢复计算。总结与延伸阅读RisingWave 流式引擎以Actor 模型实现高并发流执行以共享云存储 Hummock 存储引擎承载全部状态以一切皆表、一切皆状态统一逻辑表与内部状态的管理物化视图通过流计划 → 流片段 → Actor的流水线构建并借助Barrier/epoch 检查点提供一致快照与故障恢复能力。若想进一步深入建议继续阅读仓库内以下设计文档检查点与一致性细节docs/dev/src/design/checkpoint.md状态存储与 Hummock 存储引擎docs/dev/src/design/state-store-overview.md元数据服务与集群协调docs/dev/src/design/meta-service.md流式并行度配置docs/dev/src/design/streaming-parallelism.md执行算子源码src/stream/src/executor流图切分与调度源码src/meta/src/stream/stream_graph赞分享数据库流处理后端数据工程【免费下载链接】risingwaveEvent streaming platform for agentic AI. Continuously ingest, transform, and serve event streams in real time, at scale.项目地址https://gitcode.com/gh_mirrors/ri/risingwave点击查看免费下载相关推荐Dapr Workflow 引擎wfengine完全指南内部 Actor 架构、状态存储与容错机制Dapr Workflow 引擎wfengine完全指南内部 Actor 架构、状态存储与容错机制 导读 本文聚焦 Dapr 工作流引擎Dapr Wor后端微服务云原生消息队列AI AgentMarmite短代码实战YouTube嵌入、Spotify播放器和社交卡片制作Marmite短代码实战YouTube嵌入、Spotify播放器和社交卡片制作 Marmite是一款基于Markdown的静态网站生成器专为博客设计。本文将Leptos状态管理响应式存储与状态共享方案Leptos状态管理响应式存储与状态共享方案 痛点传统状态管理的困境 在现代Web应用开发中状态管理一直是开发者面临的核心挑战。你是否曾经遇到过这些问题前端后端Web框架SSR上一篇完整指南用10分钟音频训练一个RVC变声器模型下一篇革命性位置编码技术Open Location Code完全指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表