news 2026/9/13 19:16:33

PostHog Funnel UDF 实战指南:基于 Rust 与 ClickHouse executable_pool 的漏斗聚合加速方案

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
PostHog Funnel UDF 实战指南:基于 Rust 与 ClickHouse executable_pool 的漏斗聚合加速方案

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 明确记录了两个关键设计决策:

  1. Python 会导致 ClickHouse 节点崩溃:团队最初尝试用 Python 编写 UDF,但 Python 垃圾回收器在长期运行、高负载的场景下跟不上内存分配节奏,最终导致 ClickHouse 节点崩溃。因此全部函数改用 Rust 编写。
  2. 使用 cross-rs 交叉编译:ClickHouse 部署环境同时存在 x86_64 与 aarch64(ARM)两种架构,PostHog 使用 cross-rs 这一跨平台 Rust 编译器工具链来产出两套二进制。

这一决策在 funnel-udf/Cargo.toml 中也有呼应:发布配置采用lto = "fat"codegen-units = 1panic = "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_funnelsteps --variant=plain标准漏斗,按字符串 breakdown
aggregate_funnel_cohortsteps --variant=cohort漏斗按人群(cohort ID)拆分
aggregate_funnel_arraysteps --variant=array漏斗按字符串数组(多值属性)拆分
aggregate_funnel_trendstrends --variant=plain漏斗趋势(按区间)
aggregate_funnel_array_trendstrends --variant=array多值属性漏斗趋势
aggregate_funnel_cohort_trendstrends --variant=cohort人群漏斗趋势

此外还保留了一组*_json镜像函数与 Python 调试用 UDF(aggregate_funnel_testaggregate_funnel_array_trends_test),供人工 benchmark 与排障使用。

三、开发工作流:本地构建二进制

README 给出的开发工作流非常简洁,配合 funnel-udf/build.sh 可以还原完整步骤:

  1. 若在 flox 环境中,先输入exit退出(NixOS 不受 cross-rs 支持,见下文 Troubleshooting);
  2. 切换到funnel-udf子目录;
  3. 执行构建脚本./build.sh
  4. 重启 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-gnuaarch64-unknown-linux-gnu两个 release 二进制,并复制到 posthog/user_scripts 目录下,命名为aggregate_funnel_x86_64aggregate_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 查询构建器切换到新版本

  1. 开发:使用user_scripts顶层(未版本化)的二进制文件进行开发,schema 定义在 docker/clickhouse/user_defined_function.xml;
  2. 生成版本:发布前,在 posthog/udf_versioner.py 中递增版本号并运行该脚本。脚本会生成带版本号的二进制目录(如v12/)以及latest_user_defined_function.xml
  3. 提交 PR(第一部分):把更新后的user_scripts目录合入 master——包括新的版本目录(如v12/)与重新生成的latest_user_defined_function.xml注意:此 PR 不要包含UDF_VERSION的 bump,它属于第 5 步的独立 PR
  4. 验证部署:二进制与 XML 会在下一次配置运行时自动部署到 ClickHouse。在 Metabase 中执行SELECT aggregate_funnel_vXX()验证:
    • 返回"invalid arguments"(参数无效)说明函数已注册成功(只是缺少合法参数);
    • 返回"unknown function"(未知函数)说明部署尚未生效;
    • EU 与 US 两个区域都要验证
  5. 提交 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", }

脚本的工作流程为:

  1. 切换到posthog/user_scripts目录,创建版本目录v{UDF_VERSION}(已存在时可用-f/--force覆盖);
  2. 把顶层所有非.xml文件复制进版本目录(即aggregate_funnel_x86_64aggregate_funnel_aarch64等二进制);
  3. 解析 docker/clickhouse/user_defined_function.xml(ACTIVE_XML_CONFIG)作为当前 schema;
  4. 读取上一版本目录(如v11/)内的user_defined_function.xml作为基底,删除其中所有版本化函数(匹配_v(\d+)$命名、且不在UNVERSIONED_FUNCTIONS中、且版本号 ≥EARLIEST_UDF_VERSION的条目);
  5. 把当前 schema 中所有非免版本化函数重命名为函数名_v12形式,并把command改写为v12/aggregate_funnel ...(指向版本目录内的二进制);
  6. 写入v12/user_defined_function.xmllatest_user_defined_function.xml(后者供 ansible 自动拉取)。

这里可以看到一个细节:aggregate_funnel这类核心函数是带版本号的(部署后名为aggregate_funnel_v12),而 JSON 清洗、解压类工具函数(JSONClean*decompress等)是免版本化的,因此它们在 XML 中始终以原名存在,不受版本切换影响。从 posthog/user_scripts 目录结构可以看到历史上保留了v5v12共 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_typeunordered时走独立的 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_typestep_开头时(如step_1),归属到指定步骤,只有当事件 breakdown 与 prop_val 匹配时才计入该步骤。

输出Result(最终步骤, breakdown 值, 各步骤耗时数组, 各步骤事件 UUID 数组, 命中的步骤位掩码),与 XML 中return_typeArray(Tuple(Int8, ..., Array(Float64), Array(Array(UUID)), UInt32))一一对应。

7.2 trends:按时间区间输出转换结果

funnel-udf/src/trends.rs 在 steps 算法基础上增加了from_stepto_stepinterval_start(区间起始时间戳)维度。每个事件携带interval_start字段,IntervalDataHashMap<u64, IntervalData>按区间维护各自的max_stepentered_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)与BreakdownShapeNullableString/ArrayString/U64),支持字符串、多值数组、数字三种 breakdown 形态的序列化,其测试用例覆盖了 2^52(NOT_IN_COHORT_ID哨兵值)等边界输入。二进制 I/O(列头、行二进制编解码)位于 funnel-udf/src/io 下的steps_io.rstrends_io.rspropval.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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/13 19:01:22

深入解析 lo.DropWhile:用 Go 1.18+ 泛型按谓词丢弃切片前缀

深入解析 lo.DropWhile&#xff1a;用 Go 1.18 泛型按谓词丢弃切片前缀 【免费下载链接】lo &#x1f4a5; A Lodash-style Go library based on Go 1.18 Generics (map, filter, contains, find...) 项目地址: https://gitcode.com/GitHub_Trending/lo/lo lo 是一个基于…

作者头像 李华
网站建设 2026/9/13 19:00:10

基于gVisor的LLM代码安全执行架构设计与实践

1. 项目背景与核心挑战在AI技术快速发展的今天&#xff0c;大语言模型(LLM)的代码解释能力正逐步从实验室走向生产环境。作为西南总部AI调度官团队的技术负责人&#xff0c;我们面临着一个关键挑战&#xff1a;如何在保证系统安全的前提下&#xff0c;充分发挥LLM的代码生成与执…

作者头像 李华
网站建设 2026/9/13 18:58:54

SSD主控固件DDR初始化实战:数据结构布局、耗时优化与避坑指南

1. 这不是教科书里的“初始化”——而是主控固件在上电瞬间的生死抉择SSD 主控固件启动时需要在 DDR 中初始化哪些数据结构&#xff1f;各自的规模和耗时如何&#xff1f;——这个问题看似只是嵌入式系统里一个技术细节&#xff0c;但实际是 SSD 可靠性、性能与寿命的底层分水岭…

作者头像 李华
网站建设 2026/9/13 18:58:48

MMC实时仿真三大避坑指南:模型、求解器与硬件协同优化

1. 项目概述&#xff1a;为什么MMC实时仿真不是“把模型拖进去跑一下”那么简单做MMC&#xff08;模块化多电平换流器&#xff09;的实时仿真&#xff0c;我最初也以为就是照着教科书搭个拓扑、选个求解器、设个步长&#xff0c;点下运行——结果前三个小时全在报错里打转。第一…

作者头像 李华