☰
RisingWave 中的 prost-helpers:用过程宏为 prost 生成消息自动生成类型安全的 getter
2026/9/25 8:00:05 网站建设 项目流程
  • 数据库
  • 流处理
  • 后端
  • 数据工程

【免费下载链接】risingwave

Event 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 是面向 Agentic AI 的实时事件流平台,其内部大量使用 gRPC/Protobuf 消息在 Frontend、Meta、Compute 与 Compactor 等节点之间传递计划、目录与状态信息。为了让成千上万个由 prost 中实现了一个名为prost-helpers的过程宏 crate。本文围绕该 crate 的官方说明文档,深入讲解AnyPBderive macro 的四种 getter 生成规则、底层实现原理、错误处理机制,以及它与risingwave_pb构建流程的集成方式。读完本文,你将掌握如何在自己的 prost 生成代码中应用这套模式,彻底告别Option样板代码与裸i32枚举带来的心智负担。

背景:prost 生成消息的 Rust 使用痛点

prost 根据.proto文件生成的 Rust 结构体遵循 proto3 的语义,这给调用方带来了两类常见的样板代码:

  • 可选字段(optional / singular message)在 Rust 侧被映射为Option<T>。读取时必须先处理Option:match、as_ref()、unwrap_or…… 每个字段都要写一遍,代码冗长且容易出错。
  • 枚举字段在 proto3 中被映射为裸i32。枚举的真实类型(如data_type::TypeName)只在生成的模块里存在,调用方需要手动执行TypeName::from_i32(x)转换,且 proto3 中枚举值为0时语义上是“未设置/未指定”,这一层校验逻辑也散落在各处。

prost-helpers正是为解决这些问题而生。它提供一个 derive macroAnyPB,只要把它挂在 prost 生成的消息上,就会自动为每个字段生成对应的get_xxx()方法。

核心机制:AnyPB的三种字段处理规则

根据 README 的说明,AnyPB为字段生成 getter 的规则如下:

字段类型生成的 getter 行为
可选字段(Option<T>)getter 返回Result<&T>,字段缺失时返回错误,简化Option处理样板代码
枚举字段(i32+#[prost(enumeration=...)])getter 自动把i32转换为枚举类型,并额外校验枚举的零值——零值在 proto3 中意味着“该枚举字段未设置”
其他字段(标量、消息引用、repeated 等)getter 与直接字段访问等价,只是包了一层命名统一的访问方法

这一“统一入口 + 按需类型转换 + 显式错误”的设计,让消息消费方可以放心地通过get_xxx()访问字段,不必每次都关心底层存储形态。

完整示例:五类字段的 getter 生成效果

原文档给出了一个完整的示例。假设有一个FooMessage,同时挂上prost_helpers::AnyPB与prost::Message两个 derive:

#[derive(prost_helpers::AnyPB)] #[derive(Clone, PartialEq, ::prost::Message)] pub struct FooMessage { #[prost(message, optional, tag="1")] pub field: ::core::option::Option<Field>, #[prost(enumeration="foo_message::EnumFieldType", tag="2")] pub enum_field: i32, #[prost(uint32, tag="3")] pub uint32_field: u32, #[prost(uint32, optional, tag="4")] pub optional_uint32_field: ::core::option::Option<u32>, #[prost(message, repeated, tag="5")] pub repeated_field: ::prost::alloc::vec::Vec<Field>, }

AnyPB会自动为它展开出如下五个方法:

impl FooMessage { // 可选 message 字段:Result 包装,避免手动处理 Option pub fn get_field(&self) -> Result<&Field> { self.field .as_ref() .ok_or_else(|| crate::ProstFieldNotFound(stringify!(field))) } // 枚举字段:先校验零值,再执行 i32 -> 枚举转换 pub fn get_enum_field(&self) -> Result<foo_message::EnumFieldType> { if self.enum_field.eq(&0) { return Err(crate::ProstFieldNotFound(stringify!(enum_field))); } foo_message::EnumFieldType::from_i32(self.enum_field) .ok_or_else(|| crate::ProstFieldNotFound(stringify!(enum_field))) } // 非可选标量:直接返回值(值类型拷贝) pub fn get_uint32_field(&self) -> u32 { self.uint32_field } // 可选标量:Result 返回引用 pub fn get_optional_uint32_field(&self) -> Result<&u32> { self.optional_uint32_field .as_ref() .ok_or_else(|| crate::ProstFieldNotFound(stringify!(optional_uint32_field))) } // repeated 字段:直接返回切片引用 pub fn get_repeated_field(&self) -> &Vec<Field> { &self.repeated_field } }

需要注意:文档示例中的错误类型写作crate::ProstFieldNotFound,这是早期命名;当前仓库实际定义为PbFieldNotFound(见 src/prost/src/lib.rs),命名上二者一致,具体以当前仓库为准。

源码级实现:getter 是如何被逐字段生成的

AnyPB的过程宏入口位于 src/prost/helpers/src/lib.rs:它把 derive 输入解析为syn::DeriveInput,调用produce()后输出展开代码。produce()的核心逻辑是先判断目标是不是结构体,若是则遍历每个字段并交给 src/prost/helpers/src/generate.rs 的implement()生成对应方法。

implement()对字段类型的判定顺序,正好对应前文的三条规则,其判定逻辑值得展开:

  1. 枚举字段判定:extract_enum_type_from_field()会检查字段类型是否为i32,并解析#[prost(...)]属性中的enumeration = "path::EnumName"参数,还原出真实枚举类型(generate.rs)。命中后生成的 getter 分两步:先判断self.field == 0(proto3 枚举零值 = 未设置),再调用EnumName::from_i32(...),任何一步失败都返回PbFieldNotFound。

  2. Option<T>判定:通过extract_type_from_option()取出泛型参数T(generate.rs),生成返回Result<&T>的 getter。

  3. 值类型判定:代码中维护了一份白名单——u32, u64, f32, f64, i32, i64, bool这些基础类型直接按值返回;此外crate::id::...下的TypedId类型也被视为按值返回(generate.rs)。TypedId是 RisingWave 在 src/prost/src/id.rs 中定义的强类型 ID 封装,通过#[repr(transparent)]与底层类型保持内存布局一致,并实现prost::TransparentOver以支持透明序列化。

  4. 兜底规则:其余所有字段(如 repeated、嵌套 message 引用、字符串等)一律返回&T引用。

另外两个细节值得注意:

  • 所有生成的 getter 都带有#[inline(always)]属性(generate.rs),这些访问方法在热路径上不会产生额外调用开销。
  • 若字段上带有#[deprecated]标记,生成的 getter 也会同步继承#[deprecated]属性(generate.rs),保持 API 弃用状态一致。

produce()还做了一个额外动作:为每个消息类型生成一个Pb前缀的类型别名,例如pub type PbFooMessage = FooMessage;(lib.rs)。这使得代码库可以统一使用PbXxx命名风格引用 protobuf 类型,且 rust-analyzer 会把文档自动转发到原始类型。

错误处理:PbFieldNotFound与 tonic 的衔接

每个失败路径返回的错误类型都是PbFieldNotFound(pub &'static str),定义在 src/prost/src/lib.rs,实现了thiserror::Error,错误消息为`field `{0}` not found`。它携带的静态字符串来自stringify!(field_name),也就是字段名本身。

更关键的是,该类型实现了From<PbFieldNotFound> for tonic::Status:

impl From<PbFieldNotFound> for tonic::Status { fn from(e: PbFieldNotFound) -> Self { e.to_status_unnamed(tonic::Code::Internal) } }

见 src/prost/src/lib.rs。这意味着在 gRPC 服务实现中,任何调用get_xxx()得到的PbFieldNotFound都可以通过?运算符直接转换为tonic::Status(内部错误码Internal),错误会自动带上缺失字段名,便于排查。这条链路让“读取 Protobuf 字段 → 校验缺失 → 返回 gRPC 错误”成为一行代码。

配套宏:StreamNodeBodyVariants与Version

AnyPB并非这个 crate 唯一的宏。src/prost/helpers/src/lib.rs 还导出了另外两个:

StreamNodeBodyVariants专门服务于stream_plan::stream_node::NodeBody枚举。它要求目标必须是名为NodeBody的枚举,否则编译报错(lib.rs)。它为该枚举的每个变体生成一个同名的零大小标记类型(如SourceVariant、ProjectVariant),并导出一个#[doc(hidden)]的__dispatch_stream_node_body!宏,用于在match中根据NodeBody变体分发到不同的执行逻辑,同时利用标记类型辅助类型推断。该宏在 lib.rs 中有对应的单元测试,验证生成代码的展开结果与预期一致。

Version面向版本枚举,为枚举生成一个LATEST关联常量,指向枚举的最后一个变体(lib.rs)。RisingWave 用它来表达“当前最新协议版本”,用于消息格式演进场景。

在 RisingWave 中的集成:一次 build.rs 全量生效

prost-helpers不是手工逐个标注的,而是通过risingwave_pbcrate 的构建脚本统一注入。src/prost/build.rs 在tonic_build::configure()中写入:

.type_attribute(".", "#[derive(prost_helpers::AnyPB)]")

这条规则对整个 protobuf 文件描述符集合中的所有类型生效,也就是说所有由proto/*.proto生成的消息结构体都会自动获得AnyPB派生与配套的get_xxx()方法。同样的机制还被用于注入其他派生:

  • stream_plan.StreamNode.node_body字段获得StreamNodeBodyVariants(build.rs);
  • stream_plan.AggNodeVersion、stream_plan.PausableAggNodeVersion、expr.UdfExprVersion等版本枚举获得Version(build.rs)。

risingwave_pb的依赖关系在 src/prost/Cargo.toml 中通过prost-helpers = { path = "helpers" }声明,而prost-helpers本身是proc-macro = true的过程宏 crate,仅依赖proc-macro2、quote、syn三个解析与代码生成库(见 src/prost/helpers/Cargo.toml),编译期开销被控制在最小范围。

真实调用与测试验证

生成的 getter 在仓库中有大量真实使用,可以直接观察其效果。以 src/prost/src/lib.rs 中的stream_plan::MaterializeNode为例:

impl stream_plan::MaterializeNode { pub fn dist_key_indices(&self) -> Vec<u32> { self.get_table() .unwrap() .distribution_key .iter() .map(|i| *i as u32) .collect() } pub fn column_descs(&self) -> Vec<plan_common::PbColumnDesc> { self.get_table() .unwrap() .columns .iter() .map(|c| c.get_column_desc().unwrap().clone()) .collect() } }

这里get_table()是AnyPB为可选 message 字段生成的 getter,直接返回Result;get_column_desc()同理。调用方仅用.unwrap()即可完成校验,无需手写as_ref()+ok_or_else。

src/prost/src/lib.rs 的测试模块还对各类 getter 做了直接验证:

  • test_getter:对可选字段get_data_type()的Result返回值做断言;
  • test_enum_getter:设置type_name = TypeName::Double as i32后,get_type_name().unwrap()成功还原出枚举;
  • test_enum_unspecified:当枚举值为TypeUnspecified(零值)时,get_type_name()正确返回Err,印证了“零值即未设置”的校验语义;
  • test_primitive_getter:验证get_is_nullable()这类标量 getter 按值返回。

这些测试直接映射到 README 描述的三条规则,是理解行为契约的可靠参考。

使用注意事项与边界

  1. 错误即契约:对于 optional 字段,Result的Err表示字段未设置,这是正常业务分支,而非异常。调用方应显式决定是unwrap(确信存在)、?向上传播,还是unwrap_or提供默认值。
  2. 枚举零值语义:proto3 中枚举值0约定为“未指定”,AnyPB对枚举 getter 的零值校验是硬性行为,不区分枚举名(即使显式赋了零值枚举名,getter 仍返回Err)。
  3. 标量按值、引用按借用:基础标量(含TypedId)返回副本,其余字段返回引用,避免无谓 clone,调用方也无需承担所有权负担。
  4. 不可对已存在的自定义get_xxx方法重名:derive 展开会与手写 impl 方法冲突,自定义扩展应放在额外 impl 块中,且不要与生成的get_前缀方法重名。

总结

prost-helpers是 RisingWave 在“protobuf 生成代码可用性”上的一次系统性工程实践:通过AnyPB统一 getter 命名与返回形态、用Result折叠Option样板、在 getter 内完成枚举零值校验与类型转换,并通过PbFieldNotFound → tonic::Status的自动转换融入 gRPC 错误链。借助build.rs的全局type_attribute注入,整个risingwave_pb的所有消息无需人工改动即可获得这套能力,配合StreamNodeBodyVariants(执行节点分发)与Version(协议版本演进)两个配套宏,构成了一个完整、可复用的 prost 消息增强工具箱。相关实现与文档可继续查阅:helpers README、宏实现 lib.rs、getter 生成逻辑 generate.rs、构建注入 build.rs。

  • 数据库
  • 流处理
  • 后端
  • 数据工程

【免费下载链接】risingwave

Event streaming platform for agentic AI. Continuously ingest, transform, and serve event streams in real time, at scale.

项目地址:https://gitcode.com/gh_mirrors/ri/risingwave
点击查看免费下载
上一篇:5倍速处理PB级数据:Windmill如何用Parquet和DataFusion重构大数据工作流
下一篇:libcimbar:用一块屏幕和一部手机,跑出 850 Kbps 的气隙传输

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询