PostHog PersonHog gRPC 开发指南:Proto 消息约定、一致性路由与新增 RPC 完整实操
本文围绕 PersonHog 服务的 Proto 开发规范展开:从消息类型约定(int64/bytes/optional)、ReadOptions 一致性级别设计,到按 service → replica → leader 三层 Proto 结构新增一个 RPC 的完整六步流程。读完后,你能够独立完成一次符合仓库规范的 PersonHog gRPC 接口设计与落地,包括消息定义、代码生成、Rust 实现与 Python/Node.js 客户端封装。
PersonHog 的三层 Proto 架构
PersonHog 是 PostHog 中专门负责 person、distinct ID、group、group type mapping、cohort membership 以及 feature flag hash key override 数据的 gRPC 服务栈。它的 API 表面被拆分为三层,分别对应不同的调用方:
| Proto 文件 | 作用 |
|---|---|
| proto/personhog/types/v1/person.proto | Person、DistinctId 消息及所有 person 相关 request/response 类型 |
| proto/personhog/types/v1/group.proto | Group、GroupTypeMapping 消息及所有 group 相关 request/response 类型 |
| proto/personhog/types/v1/cohort.proto | CohortMembership 消息及 cohort 相关 request/response 类型 |
| proto/personhog/types/v1/feature_flag.proto | HashKeyOverride 消息及 feature flag request/response 类型 |
| proto/personhog/types/v1/common.proto | ReadOptions、ConsistencyLevel、TeamDistinctId 等跨领域共享类型 |
| proto/personhog/service/v1/service.proto | PersonHogService——公共 API,客户端直接调用 |
| proto/personhog/replica/v1/replica.proto | PersonHogReplica——内部 API,router 判定归属后调用的副本池 |
| proto/personhog/leader/v1/leader.proto | PersonHogLeader——person 数据的内部写 API |
调用链的关键设计是:客户端只面向 PersonHogService,router 负责透明地把请求委派给 replica(或 leader)。从 service.proto 的注释可以看到:"PersonHogService is the public API exposed by the router. Clients call this service; the router handles backend selection and routing." 而 replica.proto 则说明 "The router calls this service after determining which replica owns the vnode"。
这带来一条硬约束:service 与 replica 两个 Proto 中必须同时声明同一签名的 RPC。对比两个文件可以发现,GetPerson、GetGroups、CheckCohortMembership 等 RPC 在两边逐行一致;leader 侧(leader.proto)只包含 UpdatePersonProperties、GetPerson(强一致读)以及生命周期围栏相关的 FencePerson/ReleaseFence/FoldPersonDocument 等写路径 RPC。
消息类型约定(Message conventions)
仓库对消息定义有一套明确的命名与类型规范,新增消息时必须遵守:
- ID、时间戳(Unix 毫秒)、版本号一律用
int64; - JSON 数据一律用
bytes,不用string——适用于 properties、properties_last_updated_at、properties_last_operation、default_columns 这类字段; - 可空字段用
optional; - 附加式变更保持在 v1 内演进,只有破坏性变更才升级到 v2;
- 字段号删除后严禁复用;
- 请求消息命名为
<RpcName>Request,响应消息命名为<RpcName>Response(如共享响应类型PersonsResponse也可以复用); - 需要一致性控制的读请求必须携带
ReadOptions read_options字段,且约定放在消息的最后一个字段。
这些约定在现有 Proto 中随处可见。以 person.proto 的 Person 消息为例:
message Person {
int64 id = 1;
string uuid = 2;
int64 team_id = 3;
bytes properties = 4; // JSON 用 bytes,不是 string
bytes properties_last_updated_at = 5;
bytes properties_last_operation = 6;
// Epoch milliseconds, as are all timestamps in this message.
int64 created_at = 7; // 时间戳为 Unix 毫秒
int64 version = 8; // 版本号 int64
bool is_identified = 9;
optional bool is_user_id = 10; // 可空字段用 optional
optional int64 last_seen_at = 11;
bool is_deleted = 12;
}
再看读请求消息的 ReadOptions 位置约定,ListGroupsRequest(group.proto)与 GetPersonsRequest(person.proto)都把 ReadOptions read_options 作为最后一个字段,正是上面约定的直接体现:
message GetPersonsRequest {
int64 team_id = 1;
// Max 250 per request.
repeated int64 person_ids = 2;
ReadOptions read_options = 3; // last field
}
ReadOptions 与一致性级别
ReadOptions 已在 common.proto 中定义,规范明确要求:import 它,不要重新定义。其实际内容(common.proto)如下:
import "personhog/types/v1/common.proto";
enum ConsistencyLevel {
CONSISTENCY_LEVEL_UNSPECIFIED = 0;
CONSISTENCY_LEVEL_EVENTUAL = 1;
CONSISTENCY_LEVEL_STRONG = 2;
}
// ReadOptions configures read operation behavior.
message ReadOptions {
ConsistencyLevel consistency = 1;
repeated string field_mask = 2; // 响应中要包含的顶层字段名
}
ConsistencyLevel 两个取值的路由语义是:
EVENTUAL(默认)——查询 replica 副本池,可能存在复制延迟;STRONG——查询 primary 主库池(只有对 PersonData 的读取才有意义)。
对于 NonPersonData(group、group type mapping、cohort membership、hash key override),一致性由副本内部自行保证——router 无论一致性级别如何都直接把请求发给 replica。这一分流规则在 SKILL.md 的资格检查表中也有对应描述:posthog_person/posthog_persondistinctid 等 PersonData 表走 "读: replica(eventual) 或 leader(strong),写: leader",而四个 NonPersonData 表则 "All ops: replica"。
在 router 的 Rust 实现中,一致性级别通过请求元数据传递:proxy.rs 中会读取 x-read-consistency 请求头。而 leader.proto 的注释进一步说明了 leader 路由的机制:router 对 (team_id, person_id) 做哈希得到分区,再通过协调路由表解析出归属 pod,分区号放在 x-partition 请求元数据头中而不是请求体里——这样 router 可以逐字节转发客户端请求,无需解码再编码。
值得注意的是 ReadOptions 里还有一个超出基础约定文档范围的字段 field_mask:它允许批量/列表类 RPC 按顶层字段名裁剪返回内容,当 properties 相关字段全部被排除时服务端会在 SELECT 子句中跳过取数;未知字段名会被静默忽略。这是"附加式变更保持在 v1"这一演进策略的典型案例——不破坏旧客户端,同时给新客户端提供性能优化通道。
实战示例:为 PersonHog 新增 GetPersonCount RPC
规范文档给出了一个端到端的完整示例:新增一个统计 person 数量的 GetPersonCount 接口。下面按六步展开,并结合仓库实际情况补充代码生成与验证环节。
第 1 步:在 types/v1/person.proto 中定义消息
由于统计对象是 person,消息放入 person 领域文件(proto/personhog/types/v1/person.proto):
message GetPersonCountRequest {
int64 team_id = 1;
ReadOptions read_options = 2; // 读请求,read_options 放最后
}
message GetPersonCountResponse {
int64 count = 1;
}
第 2 步:在 service.proto 中声明 RPC
service PersonHogService {
// ... existing RPCs ...
rpc GetPersonCount(personhog.types.v1.GetPersonCountRequest)
returns (personhog.types.v1.GetPersonCountResponse);
}
第 3 步:在 replica.proto 中声明同名 RPC
service PersonHogReplica {
// ... existing RPCs ...
rpc GetPersonCount(personhog.types.v1.GetPersonCountRequest)
returns (personhog.types.v1.GetPersonCountResponse);
}
两边的 RPC 签名必须完全一致——router 依赖这一契约把公共请求透明转发给副本池。参考 service.proto 与 replica.proto 中现有的 CountCohortMembers、CountGroupTypeMappings 等计数类 RPC,两侧声明逐行对应。
第 4 步:确定路由策略
GetPersonCount 是一次 PersonData 读操作,路由规则为:EVENTUAL 一致性走 replica,STRONG 一致性走 leader。按规范文档,router 侧使用形如以下的路由调用(route_request 来自 router 的 routing.rs 模块):
let route = route_request(
DataCategory::PersonData,
OperationType::Read,
get_consistency(&request.read_options),
)?;
这里 DataCategory::PersonData 决定了它属于"按一致性分流"的类别;如果换成 GetGroupTypeMappingsByTeamId 这类 NonPersonData 读,则无论一致性级别都固定发往 replica。若新增的是 person 数据写操作,则还需在 leader.proto 中追加 RPC,并且必须满足 at-least-once 投递安全约束:leader 路径的请求可能被重复投递(客户端在失败歧义后会重试 UNAVAILABLE,router 内部还会重放被围栏/传输层弹回的请求),因此每个 leader RPC 必须在重复投递下收敛——只读查询、可重放的 merge、墓碑式删除、max-merge 版本地板,或携带显式操作标识以便去重。既不收敛也不带幂等键的操作(无保护的自增、追加)禁止添加,需要时先重新设计成携带幂等键的形态。这一点在 SKILL.md 的 "Leader RPCs must be safe under at-least-once delivery" 一节中有完整论证,leader.proto 中 FoldPersonDocument 的注释也是该原则的现实例证。
第 5 步:生成客户端桩代码并封装 Python 客户端
Proto 定义完成后,各语言客户端需要生成桩代码。Python 侧运行:
bin/generate_personhog_proto.sh
然后更新三处:
posthog/personhog_client/proto/__init__.py——为新 request/response 消息补充 re-export;- posthog/personhog_client/client.py——按既有方法模式新增包装方法:
def get_person_count(self, request: GetPersonCountRequest) -> GetPersonCountResponse:
return self._stub.GetPersonCount(request, timeout=self._timeout)
posthog/personhog_client/fake_client.py——为测试实现该方法,记录调用并返回构造的响应:
def get_person_count(self, request: Any) -> Any:
team_id = request.team_id
count = sum(1 for (tid, _) in self._persons_by_id if tid == team_id)
call = _Call("get_person_count", request)
resp = person_pb2.GetPersonCountResponse(count=count)
call.response = resp
self.calls.append(call)
return resp
Node.js 侧则运行:
cd nodejs && pnpm run generate:personhog-proto
随后在 nodejs/src/common/personhog/persons.ts 的对应 operations 类中按既有模式补充包装方法,并在 nodejs/src/common/personhog/client.test.ts 的 SERVICE_DEFAULTS 中为新 RPC 加默认桩。Rust 侧无需生成步骤——tonic 会在 cargo build 时重新生成,但必须实现新 RPC,否则编译直接失败。
第 6 步:Rust 侧实现与验证
Rust 实现分三层,编译器会引导你完成——Proto 定义后 cargo build 的报错会精确指出缺失的 trait 方法:
- 存储层(personhog-replica):在
rust/personhog-replica/src/storage/traits/<domain>.rs加 trait 方法,在storage/postgres/<domain>.rs用sqlx::query_as!/sqlx::query!实现查询,读用replica_pool、写用primary_pool,并补上DB_QUERY_DURATION/DB_ROWS_RETURNED指标埋点; - 服务层(personhog-replica):在
rust/personhog-replica/src/service/mod.rs添加 RPC handler——提取 Proto 字段、调用 storage、把存储结果转成 Proto 响应、把存储错误映射为 tonicStatus码; - 路由接线(personhog-router):在
rust/personhog-router/src/router/mod.rs加路由方法(使用route_request与正确的DataCategory/OperationType,用call_backend!宏做埋点),在rust/personhog-router/src/service/mod.rs中用route_request!宏委派,并在backend/mod.rs的 backend trait 与replica.rs中补齐实现。
测试分别放在 rust/personhog-replica/tests/storage_tests.rs、rust/personhog-replica/tests/service_tests.rs 与 rust/personhog-router/tests/;同一行为的多种变体推荐用 rstest 参数化测试。
前置约束:查询必须命中索引
在写任何 Proto 之前,规范还要求先确定所需 SQL 查询并对照可用索引验证——每条查询都必须是索引扫描,绝不允许顺序扫描。索引清单以 rust/persons_migrations/ 中的 SQL 迁移为唯一事实来源,references/database-indexes.md 给出了汇总。例如 posthog_person 按 team_id 哈希分区为 64 个分区,复合主键为 (team_id, id),典型查询模式 WHERE team_id = $1 AND id = $2、WHERE team_id = $1 AND uuid = $2 均可走分区裁剪的索引扫描。这也解释了为何 PersonHog 所有请求消息都强制携带 team_id——它是分区裁剪与索引命中的前提。
交付前检查清单
完成一次 PersonHog RPC 新增后,按 SKILL.md 的清单逐项核对:
- 查询使用既有索引(无 seq scan);
- Proto 消息已加入
types/v1/<domain>.proto,RPC 已同时加入service.proto与replica.proto(必要时加leader.proto); - leader RPC 满足 at-least-once 投递安全;
- Python 桩已生成、
proto/__init__.py/client.py/fake_client.py已更新; - Node.js 桩已生成、
client.test.ts的SERVICE_DEFAULTS已更新; - Rust 存储 trait、postgres 实现、service handler、router 接线全部落地;
- 存储、服务、路由三层测试已补齐;
- 三个 crate 的构建与测试全部通过:
cargo build -p personhog-proto
cargo build -p personhog-replica
cargo build -p personhog-router
cargo test -p personhog-replica
cargo test -p personhog-router
综上,PersonHog 的 Proto 体系可以概括为三句话:类型层按领域分文件、遵循 int64/bytes/optional 与 ReadOptions 尾部约定;服务层与副本层 RPC 签名强制镜像;写路径必须满足重复投递收敛。掌握这三点后,新增任何 PersonHog gRPC 接口都有明确且可验证的落地路径。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0624
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00