PostHog Funnel UDF 实战指南:基于 Rust 与 ClickHouse executable_pool 的漏斗聚合加速方案
【免费下载链接】posthog:hedgehog: PostHog is the leading platform for building self-driving products. Our developer tools – AI observability, analytics, session replay, flags, experiments, error tracking, logs, and more – capture all the context agents need to diagnose problems, uncover opportunities, and ship fixes. Steer it all from Slack, web, desktop, or the MCP.项目地址: https://gitcode.com/GitHub_Trending/po/posthog
导读
本文以 PostHog 仓库中的 funnel-udf/README.md 为核心,系统讲解 PostHog 如何将产品分析中最核心、最昂贵的漏斗(Funnel)与趋势(Trends)聚合计算,从 Python 脚本迁移到 Rust 编写的 ClickHouse 可执行 UDF(User Defined Function),并完整梳理其本地构建、版本化部署、线上验证与问题排查的完整流程。读完本文,你将掌握:为什么 PostHog 选择用 Rust 而非 Python 编写 UDF、build.sh如何跨架构交叉编译、udf_versioner.py如何完成带版本号的灰度发布、以及如何在 ClickHouse 与 Metabase 中验证新 UDF 是否生效。
一、背景与设计决策:为什么弃用 Python,改用 Rust
漏斗分析是 PostHog 产品分析的核心能力,需要把每个用户的原始事件流按漏斗步骤做顺序匹配、时间窗判定、排除事件与可选步骤处理,计算量大且要求极高的吞吐。funnel-udf 的 README 明确记录了两个关键设计决策:
- Python 会导致 ClickHouse 节点崩溃:团队最初尝试用 Python 编写 UDF,但 Python 垃圾回收器在长期运行、高负载的场景下跟不上内存分配节奏,最终导致 ClickHouse 节点崩溃。因此全部函数改用 Rust 编写。
- 使用 cross-rs 交叉编译:ClickHouse 部署环境同时存在 x86_64 与 aarch64(ARM)两种架构,PostHog 使用 cross-rs 这一跨平台 Rust 编译器工具链来产出两套二进制。
这一决策在 funnel-udf/Cargo.toml 中也有呼应:发布配置采用lto = "fat"、codegen-units = 1、panic = "abort"、strip = "symbols",注释明确说明"UDF 运行在长生命周期的进程池中,但构建体积与启动速度仍然影响进程重启路径",可见其对运行时稳定性与内存布局的极致追求。
二、架构总览:一个 Rust 二进制,六种线上函数,两种计算模式
从 docker/clickhouse/user_defined_function.xml 可以看到,所有漏斗 UDF 共用同一个 Rust 可执行文件funnels(构建产物被复制为aggregate_funnel_*),通过命令行参数区分运行模式:
aggregate_funnel <steps|trends> --variant=<plain|cohort|array> [--json]核心入口在 funnel-udf/src/main.rs 的parse_cli中:
- 模式(mode):
steps(漏斗步骤结果)或trends(按时间区间划分的趋势结果); - 变体(--variant):固定 breakdown 的线上线格式,
plain对应Nullable(String)、cohort对应UInt64(人群 ID)、array对应Array(String); - 格式(format):默认
RowBinaryWithNamesAndTypes(高性能二进制协议),--json切换到JSONEachRow,仅用于调试与基准对比。
README 与 XML 中均强调一个关键设计:--variant在进程生命周期内固定不变。因为 XML 中每个<function>块都拥有独立的executable_pool,进程启动时变体已确定,调用期间不会变化,因此可以针对每种线格式做专门优化,无需在运行时动态判断。
六种线上函数及其命令映射为:
| 函数名 | 命令 | 用途 |
|---|---|---|
aggregate_funnel | steps --variant=plain | 标准漏斗,按字符串 breakdown |
aggregate_funnel_cohort | steps --variant=cohort | 漏斗按人群(cohort ID)拆分 |
aggregate_funnel_array | steps --variant=array | 漏斗按字符串数组(多值属性)拆分 |
aggregate_funnel_trends | trends --variant=plain | 漏斗趋势(按区间) |
aggregate_funnel_array_trends | trends --variant=array | 多值属性漏斗趋势 |
aggregate_funnel_cohort_trends | trends --variant=cohort | 人群漏斗趋势 |
此外还保留了一组*_json镜像函数与 Python 调试用 UDF(aggregate_funnel_test、aggregate_funnel_array_trends_test),供人工 benchmark 与排障使用。
三、开发工作流:本地构建二进制
README 给出的开发工作流非常简洁,配合 funnel-udf/build.sh 可以还原完整步骤:
- 若在 flox 环境中,先输入
exit退出(NixOS 不受 cross-rs 支持,见下文 Troubleshooting); - 切换到
funnel-udf子目录; - 执行构建脚本
./build.sh; - 重启 ClickHouse docker 容器使新函数生效。
build.sh的实际行为如下:
#!/bin/sh # 若在 flox 环境则直接退出 if [ -n "$FLOX_ENV" ]; then echo "⚠️ Please exit the flox environment by typing 'exit' before connecting to a toolbox." exit 0 fi echo "Installing cross" cargo install cross --git https://github.com/cross-rs/cross echo "Building for x86_64" cross build --target x86_64-unknown-linux-gnu --release echo "Building for aarch64" cross build --target aarch64-unknown-linux-gnu --release echo "Copying executables to posthog/user_scripts" cp target/x86_64-unknown-linux-gnu/release/funnels ../posthog/user_scripts/aggregate_funnel_x86_64 cp target/aarch64-unknown-linux-gnu/release/funnels ../posthog/user_scripts/aggregate_funnel_aarch64可以看到脚本会自动安装 cross-rs,分别产出x86_64-unknown-linux-gnu与aarch64-unknown-linux-gnu两个 release 二进制,并复制到 posthog/user_scripts 目录下,命名为aggregate_funnel_x86_64与aggregate_funnel_aarch64(这也与 XML 中command里出现的aggregate_funnel可执行名对应——同目录下还保留了一个无架构后缀的aggregate_funnel副本,供本地测试直接调用)。
四、部署工作流:自动拉取 + 版本化灰度发布
4.1 线上 XML 自动同步
README 指出,部署到 ClickHouse 的user_defined_function.xml不再需要手动上传,而是由 ansible 从本仓库master分支自动拉取:
- 拉取逻辑位于
posthog-cloud-infra/ansible/roles/clickhouse/tasks/setup_udfs.yml(外部 infra 仓库,本仓库不包含); - 它下载
posthog/user_scripts/latest_user_defined_function.xml,并将其部署为 ClickHouse 的user_defined_function.xml。
也就是说,只要把新版本文件合入 master 分支,下一次配置运行就会自动部署。
4.2 可回滚的云上发布流程(五步法)
README 给出了一套兼顾"线上注册"与"查询侧切换"分离的可回滚发布流程,核心思路是先让新版本函数在 ClickHouse 注册并验证,再让 PostHog 查询构建器切换到新版本:
- 开发:使用
user_scripts顶层(未版本化)的二进制文件进行开发,schema 定义在 docker/clickhouse/user_defined_function.xml; - 生成版本:发布前,在 posthog/udf_versioner.py 中递增版本号并运行该脚本。脚本会生成带版本号的二进制目录(如
v12/)以及latest_user_defined_function.xml; - 提交 PR(第一部分):把更新后的
user_scripts目录合入 master——包括新的版本目录(如v12/)与重新生成的latest_user_defined_function.xml。注意:此 PR 不要包含UDF_VERSION的 bump,它属于第 5 步的独立 PR; - 验证部署:二进制与 XML 会在下一次配置运行时自动部署到 ClickHouse。在 Metabase 中执行
SELECT aggregate_funnel_vXX()验证:- 返回"invalid arguments"(参数无效)说明函数已注册成功(只是缺少合法参数);
- 返回"unknown function"(未知函数)说明部署尚未生效;
- EU 与 US 两个区域都要验证;
- 提交 PR(第二部分):部署确认后,再单独提交一个 PR,把 posthog/udf_versioner.py 中的
UDF_VERSION递增到新版本。这一步才会让 PostHog 查询构建器切换到新 UDF。
这种"先注册、后切换"的两阶段发布策略,使得新函数与旧函数可以并存(旧版本函数仍保留在 XML 中),一旦新版本出现问题,查询侧回滚只需改回UDF_VERSION,无需重新部署 ClickHouse。
五、udf_versioner.py 版本化机制详解
posthog/udf_versioner.py 是发布流程的核心工具,其关键逻辑值得展开:
UDF_VERSION = 12 # 当前版本 EARLIEST_UDF_VERSION = 11 # 早于此版本的函数会被清理 UNVERSIONED_FUNCTIONS = { "decompress", "JSONCleanPostHogEventProperties", "JSONCleanPostHogPersonProperties", "JSONCleanPostHogTemporaryProperties", "JSONStripEmptyStringsAndNulls", }脚本的工作流程为:
- 切换到
posthog/user_scripts目录,创建版本目录v{UDF_VERSION}(已存在时可用-f/--force覆盖); - 把顶层所有非
.xml文件复制进版本目录(即aggregate_funnel_x86_64、aggregate_funnel_aarch64等二进制); - 解析 docker/clickhouse/user_defined_function.xml(
ACTIVE_XML_CONFIG)作为当前 schema; - 读取上一版本目录(如
v11/)内的user_defined_function.xml作为基底,删除其中所有版本化函数(匹配_v(\d+)$命名、且不在UNVERSIONED_FUNCTIONS中、且版本号 ≥EARLIEST_UDF_VERSION的条目); - 把当前 schema 中所有非免版本化函数重命名为
函数名_v12形式,并把command改写为v12/aggregate_funnel ...(指向版本目录内的二进制); - 写入
v12/user_defined_function.xml与latest_user_defined_function.xml(后者供 ansible 自动拉取)。
这里可以看到一个细节:aggregate_funnel这类核心函数是带版本号的(部署后名为aggregate_funnel_v12),而 JSON 清洗、解压类工具函数(JSONClean*、decompress等)是免版本化的,因此它们在 XML 中始终以原名存在,不受版本切换影响。从 posthog/user_scripts 目录结构可以看到历史上保留了v5~v12共 8 个版本目录,印证了这一机制在持续运行。
六、XML 注册配置:executable_pool 的关键参数
以aggregate_funnel为例,docker/clickhouse/user_defined_function.xml 中完整的注册配置如下:
<function> <type>executable_pool</type> <name>aggregate_funnel</name> <return_type>Array(Tuple(Int8, Nullable(String), Array(Float64), Array(Array(UUID)), UInt32))</return_type> <return_name>result</return_name> <argument> <type>UInt8</type> <name>num_steps</name> </argument> <argument> <type>UInt64</type> <name>conversion_window_limit</name> </argument> <argument> <type>String</type> <name>breakdown_attribution_type</name> </argument> <argument> <type>String</type> <name>funnel_order_type</name> </argument> <argument> <type>Array(Nullable(String))</type> <name>prop_vals</name> </argument> <argument> <type>Array(Int8)</type> <name>optional_steps</name> </argument> <argument> <type>Array(Tuple(Nullable(Float64), UUID, Nullable(String), Array(Int8)))</type> <name>value</name> </argument> <format>RowBinaryWithNamesAndTypes</format> <send_chunk_header>true</send_chunk_header> <stderr_reaction>throw</stderr_reaction> <command>aggregate_funnel steps --variant=plain</command> <lifetime>600</lifetime> </function>关键参数含义:
<type>executable_pool</type>:使用进程池模式,进程常驻、跨多次调用复用,避免每次调用都 fork/exec,这是高性能的关键(在 funnel-udf/src/main.rs 的注释中有明确说明);<format>RowBinaryWithNamesAndTypes</format>:输入输出采用 ClickHouse 二进制行格式,比 JSON 快得多;<send_chunk_header>true</send_chunk_header>:ClickHouse 在每批数据前发送一个十进制行数与换行符,Rust 侧在 funnel-udf/src/codec/chunk.rs 的read_chunk_header中解析,遇到空行会报错而非静默当作 EOF,防止丢数据;<stderr_reaction>throw</stderr_reaction>:进程写 stderr 时直接抛出异常,便于快速暴露问题;<lifetime>600</lifetime>:进程池中进程的最大存活秒数(600 秒),到期后重启。
七、源码实现原理:steps 与 trends 的聚合算法
7.1 steps:顺序/无序/严格漏斗 + 排除事件 + 可选步骤
funnel-udf/src/steps.rs 实现了标准漏斗聚合。输入结构Args包含num_steps(漏斗步骤数)、conversion_window_limit(转换窗口,单位秒)、breakdown_attribution_type(breakdown 归属方式)、funnel_order_type(漏斗顺序类型)、prop_vals(拆分属性值列表)、optional_steps(可选步骤)以及value(每个用户的事件序列)。
核心算法要点(从源码可以确认的实现细节):
- 三种顺序类型:
funnel_order_type为unordered时走独立的 funnel-udf/src/unordered_steps.rs 无序实现;strict(严格)时,在处理事件后会把所有未命中的步骤清空(process_event末尾的逻辑),强制步骤必须严格按序发生;默认则为有序(ordered)匹配; - 转换窗口:判断事件时间戳与前一步命中时间戳之差是否
<= conversion_window_limit,超出窗口的匹配无效; - 排除事件:事件
steps中的负值表示排除事件(代码中exclusion = true; -*step),命中排除事件会把对应步骤标记为excluded,最终结果中该用户输出Result(-1, ...)表示被排除; - 可选步骤:
optional_steps列表中的步骤可以跳过,匹配时允许向前回溯查找更早的已命中步骤,最终结果数会减去已命中的可选步骤数; - 同时间戳多事件:对同一时间戳的多个不同事件做排列处理(
chunk_by(|e| e.timestamp)+ 排序),处理最复杂的同时刻命中场景,并支持提前终止(命中最后一步即 break); - breakdown 归属:
breakdown_attribution_type以step_开头时(如step_1),归属到指定步骤,只有当事件 breakdown 与 prop_val 匹配时才计入该步骤。
输出Result为(最终步骤, breakdown 值, 各步骤耗时数组, 各步骤事件 UUID 数组, 命中的步骤位掩码),与 XML 中return_type的Array(Tuple(Int8, ..., Array(Float64), Array(Array(UUID)), UInt32))一一对应。
7.2 trends:按时间区间输出转换结果
funnel-udf/src/trends.rs 在 steps 算法基础上增加了from_step、to_step与interval_start(区间起始时间戳)维度。每个事件携带interval_start字段,IntervalData以HashMap<u64, IntervalData>按区间维护各自的max_step与entered_timestamp,最终输出(interval_start, 1 或 -1, breakdown 值, 事件 UUID)四元组,对应 XML 中 trends 系函数的Array(Tuple(UInt64, Int8, ..., UUID))返回类型。其中Exclusion枚举(Not/Partial/Full)区分了"仅命中排除事件"(部分排除)与"排除后再命中事件"(完全排除)两种语义。
7.3 类型与线格式
funnel-udf/src/types.rs 定义了PropVal枚举(String(Bytes)/Vec/Int/VecInt)与BreakdownShape(NullableString/ArrayString/U64),支持字符串、多值数组、数字三种 breakdown 形态的序列化,其测试用例覆盖了 2^52(NOT_IN_COHORT_ID哨兵值)等边界输入。二进制 I/O(列头、行二进制编解码)位于 funnel-udf/src/io 下的steps_io.rs、trends_io.rs、propval.rs等文件中,与 XML 中的参数类型定义严格对应。
八、常见问题排查(Troubleshooting)
在 flox 环境中运行./build.sh失败
若你使用 flox 进行开发,必须先退出 flox 环境(NixOS 不受 cross-rs 支持),否则会得到如下报错:
Error: 0: could not determine os in target triplet 1: unsupported os in target, abi: "1.82.0", system: "rustc"build.sh开头也内置了防御性检查:检测到FLOX_ENV环境变量时直接提示先执行exit再连接 toolbox,避免在错误环境中浪费时间。这同时解释了 README 开发工作流第一步"Exit flox:exit"的原因。
部署后验证函数是否生效
按 README 的建议,在 Metabase 中对 EU 与 US 分别执行SELECT aggregate_funnel_vXX():
- "invalid arguments"→ 函数已注册,只差合法参数(部署成功);
- "unknown function"→ 部署还没生效(XML 尚未同步到 ClickHouse)。
九、总结
PostHog 的 funnel-udf 是一套"以 Rust 替代 Python、以版本化 XML 管理线上注册、以两阶段 PR 完成灰度切换"的 ClickHouse 可执行 UDF 完整实践:
- 性能与稳定:Rust +
executable_pool+RowBinaryWithNamesAndTypes二进制协议,配合 aggressive release profile,解决 Python 垃圾回收导致的节点崩溃问题; - 可回滚发布:
udf_versioner.py生成vXX/版本目录与latest_user_defined_function.xml,ansible 自动拉取部署,先验证注册、再切换查询侧; - 架构双维度:steps/trends 两种计算模式 × plain/cohort/array 三种 breakdown 线格式,覆盖 PostHog 漏斗分析的全部场景。
对于任何需要在 ClickHouse 中承载高复杂度、高吞吐聚合计算的团队,这套"独立 Rust 二进制 + 版本化 XML + 自动化部署"的 UDF 工程化方案都有直接的可复制价值。
相关文件索引
- 文档:funnel-udf/README.md
- 构建脚本:funnel-udf/build.sh
- 入口与 CLI:funnel-udf/src/main.rs
- 漏斗算法:funnel-udf/src/steps.rs、funnel-udf/src/trends.rs、funnel-udf/src/unordered_steps.rs
- 协议编解码:funnel-udf/src/codec/chunk.rs、funnel-udf/src/codec
- 函数注册 schema:docker/clickhouse/user_defined_function.xml
- 版本化发布工具:posthog/udf_versioner.py
- 已发布的版本化产物:posthog/user_scripts
【免费下载链接】posthog:hedgehog: PostHog is the leading platform for building self-driving products. Our developer tools – AI observability, analytics, session replay, flags, experiments, error tracking, logs, and more – capture all the context agents need to diagnose problems, uncover opportunities, and ship fixes. Steer it all from Slack, web, desktop, or the MCP.项目地址: https://gitcode.com/GitHub_Trending/po/posthog
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考