PostHog Weekly Digest 架构解析:Temporal 工作流 + Redis 中间层驱动的个性化周报引擎
【免费下载链接】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 为所有客户提供的默认产品能力之一。本文以仓库中 posthog/temporal/weekly_digest/README.md 为主干,结合该目录下的完整源码,深入讲解这套周报系统的整体架构、Redis 存储结构、三类数据(团队级 / 组织级 / 用户级)的组织方式,以及如何向周报中新增"团队级字段"与"用户级字段"。读完本文,你将掌握一个大规模多租户邮件系统的经典实现范式:用 Temporal 子工作流做两阶段流水线(先生成、后发送),用 Redis 作为批次间的数据交换层,用 Pydantic 模型做类型化的数据契约。
一、系统定位与总体架构
Weekly Digest 是 PostHog 面向所有客户发送的一封周期性邮件,汇总其 PostHog 项目中过去一周发生的关键活动,例如新建的仪表盘、新定义的事件、启动/完成的实验、新增的特性开关(Feature Flag)、有趣的 Session Replay 筛选器、即将过期的录制、启动的调研问卷、事件量与活跃用户的周环比(WoW)趋势、以及新增的 Error Tracking 问题等。
整套流程由两个 Temporal 工作流协作完成(定义见 workflows.py):
- GenerateDigestDataWorkflow(Temporal 名称
generate-digest-data):负责生成所有摘要数据,并写入 Redis; - SendWeeklyDigestWorkflow(Temporal 名称
send-weekly-digest):从 Redis 读取数据,为每个用户渲染并发送个性化邮件。
两者之上还有一个编排入口WeeklyDigestWorkflow(Temporal 名称weekly-digest),它负责计算摘要周期(上周时间窗与digest_key)、解析输入参数,然后依次以子工作流方式调用上述两个工作流。
三个工作流均继承自PostHogWorkflow基类(posthog/temporal/common/base.py),并被统一注册在init.py 的WORKFLOWS/ACTIVITIES列表中,供start_temporal_worker与start_temporal_workflow等管理命令加载(参见 start_temporal_workflow.py 中from posthog.temporal.weekly_digest import WORKFLOWS as WEEKLY_DIGEST_WORKFLOWS)。
从源码结构看,这套
weekly_digest模块同时是其他产品线周报的复用基座:例如 products/web_analytics/backend/temporal/weekly_digest/ 实现了 Web Analytics 的周报工作流wa_weekly_digest,其管理命令测试甚至专门断言了"WA digests 已随 weekly digest 一起注册"(test_start_temporal_worker.py)。
为什么需要两阶段 + Redis 中间层?
周报要处理的是"全量组织 → 全量团队 → 全量用户"的扇出(fan-out)计算。若在发送阶段实时聚合每个用户的数据,会产生海量重复查询;而把"团队级数据的计算"与"用户级邮件的发送"解耦成两个可独立重试、独立扩容的工作流,中间用 Redis 落盘中间结果,则可以:
- 团队级数据只计算一次,所有用户共享同一份
OrganizationDigest; - 生成阶段可以按团队 id 区间并行切分,发送阶段可以按组织区间并行切分;
- 生成失败时只需重跑生成工作流,发送阶段完全幂等。
二、Redis 存储结构:三类 Key 与类型安全约定
所有 Redis Key 都通过 keys.py 中的辅助函数生成,严禁在代码中硬编码 Key 字符串,必须使用类型化枚举与生成函数。Key 统一以{digest_key}为前缀,例如weekly-digest-2024-01(由 ISO 年-周号生成,见workflows.py中f"weekly-digest-{year}-{week}")。
2.1 团队级数据(Team-level)
通过team_data_key(digest_key, TeamDataKey.*, team_id)生成,映射关系与枚举定义完全一致(对比 keys.py):
TeamDataKey枚举 | Key 模式 | 内容 |
|---|---|---|
DASHBOARDS | {digest_key}-dashboards-{team_id} | 新建的仪表盘 |
EVENT_DEFINITIONS | {digest_key}-event-definitions-{team_id} | 新定义的事件 |
EXPERIMENTS_LAUNCHED | {digest_key}-experiments-launched-{team_id} | 启动的实验 |
EXPERIMENTS_COMPLETED | {digest_key}-experiments-completed-{team_id} | 完成的实验 |
EXTERNAL_DATA_SOURCES | {digest_key}-external-data-sources-{team_id} | 新增的外部数据源 |
FEATURE_FLAGS | {digest_key}-feature-flags-{team_id} | 新增的特性开关 |
SAVED_FILTERS | {digest_key}-saved-filters-{team_id} | 有趣的 Replay 筛选器 |
EXPIRING_RECORDINGS | {digest_key}-expiring-recordings-{team_id} | 即将过期录制的数量 |
SURVEYS_LAUNCHED | {digest_key}-surveys-launched-{team_id} | 启动的调研问卷 |
USAGE_TRENDS | {digest_key}-usage-trends-{team_id} | 事件量 + 活跃用户,周环比 |
ERROR_ISSUES | {digest_key}-error-issues-{team_id} | 新增的 Error Tracking 问题 |
其中USAGE_TRENDS与ERROR_ISSUES是注释中标注的"新信号"(new signals):TeamDigest为它们提供了默认值(UsageTrends()与ErrorIssueList(root=[])),保证旧调用方与旧测试仍能正常构造对象(见 types.py)。
2.2 组织级数据(Organization-level)
通过org_digest_key(digest_key, org_id)生成:
| Key 模式 | 内容 |
|---|---|
{digest_key}-{org_id} | 一个OrganizationDigest,内含该组织下所有团队的摘要 |
OrganizationDigest在生成阶段由generate_organization_digest_batch活动一次性聚合写盘(活动名generate-organization-digest-batch),发送阶段每个用户都从这份共享数据出发做个性化裁剪。
2.3 用户级数据(User-level)
通过user_data_key(digest_key, UserDataKey.*, user_id)生成(见 keys.py):
UserDataKey枚举 | Key 模式 | 内容 |
|---|---|---|
NOTIFY_TEAMS | {digest_key}-user-notify-{user_id} | Redis SET,记录该用户需要接收通知的团队 id 集合 |
PRODUCT_SUGGESTION | {digest_key}-product-suggestion-{user_id} | 单个DigestProductSuggestion,用于产品推荐 |
注意NOTIFY_TEAMS使用r.sadd写入(一个用户可属于多个团队),而其余数据使用r.setex写入 JSON 字符串,两者都设置了CommonInput.redis_ttl的 TTL——默认 3 天(3600 * 24 * 3),见 types.py。
2.4 公共参数CommonInput
所有活动输入都携带一个CommonInput(types.py):
| 字段 | 默认值 | 说明 |
|---|---|---|
redis_ttl | 3600 * 24 * 3(3 天) | Redis Key 的过期时间 |
redis_host | None(运行时回退到环境变量WEEKLY_DIGEST_REDIS_HOST,默认localhost) | 周报专用 Redis 主机 |
redis_port | None(运行时回退到WEEKLY_DIGEST_REDIS_PORT,默认6379) | 周报专用 Redis 端口 |
batch_size | 2500 | 团队/组织分批切分的批大小 |
django_redis_url | None(运行时回退到settings.REDIS_URL) | Django 缓存 Redis 地址,仅 Replay 筛选器计数需要 |
环境变量回退逻辑位于WeeklyDigestWorkflow.run(workflows.py):WEEKLY_DIGEST_REDIS_HOST、WEEKLY_DIGEST_REDIS_PORT用于指定周报专用 Redis(避免与业务缓存互相污染),django_redis_url则用于读取 Django 缓存中已算好的播放列表计数。
三、数据流:从团队切分到邮件发送
README 给出的完整数据流如下(此图与 workflows.py 中的实际编排一一对应):
┌─────────────────────────────────────────────────────────────────────────────┐ │ GenerateDigestDataWorkflow │ ├─────────────────────────────────────────────────────────────────────────────┤ │ │ │ 1. Cut teams into id-range batches, count orgs for batching │ │ │ │ 2. Generate team-level data (parallel per batch): │ │ ├── generate_dashboard_lookup │ │ ├── generate_event_definition_lookup │ │ ├── generate_experiment_launched_lookup │ │ ├── generate_experiment_completed_lookup │ │ ├── generate_external_data_source_lookup │ │ ├── generate_feature_flag_lookup │ │ ├── generate_survey_lookup │ │ ├── generate_filter_lookup │ │ ├── generate_recording_lookup │ │ ├── generate_user_notification_lookup │ │ └── generate_product_suggestion_lookup │ │ │ │ 3. Aggregate into org digests: │ │ └── generate_organization_digest_batch │ │ │ └─────────────────────────────────────────────────────────────────────────────┘ │ ▼ Redis Storage │ ▼ ┌─────────────────────────────────────────────────────────────────────────────┐ │ SendWeeklyDigestWorkflow │ ├─────────────────────────────────────────────────────────────────────────────┤ │ │ │ For each organization (batched): │ │ 1. Load OrganizationDigest from Redis │ │ 2. For each org member: │ │ a. Load user's notification team set │ │ b. Load user's product suggestion │ │ c. Create UserSpecificDigest via org_digest.for_user() │ │ d. Render payload and send via PostHog capture event │ │ │ └─────────────────────────────────────────────────────────────────────────────┘3.1 生成阶段的关键设计
团队 id 区间切分:list_team_id_ranges活动(activities.py)一次性对query_team_ids_for_digest()做主键索引扫描,按batch_size切成若干[start, end)半开区间(_cut_team_id_ranges)。这样每个 generator 活动都通过主键谓词filter(id__gte=start, id__lt=end)直接定位到批首行,避免 LIMIT/OFFSET 每批都要构建并丢弃批前所有行的性能问题(activities.py)。
并行扇出:GenerateDigestDataWorkflow.run用itertools.product(team_id_ranges, generators)生成"区间 × 生成器"的全笛卡尔组合,通过asyncio.gather并发调度所有活动;每个活动配置了start_to_close_timeout=1 小时、heartbeat_timeout=2 分钟、RetryPolicy(maximum_attempts=2, initial_interval=1 分钟)(workflows.py)。活动内部全部使用Heartbeater()保持心跳,配合日志上下文bind_contextvars记录digest_key、period_start、period_end与团队区间。
通用的 lookup 模板:除 Replay 相关与 usage trends 外,绝大多数团队级数据都走同一个模板函数generate_digest_data_lookup(activities.py),它接收三个参数:key_kind(TeamDataKey枚举)、query_func(从 queries.py 传入的查询函数)、resource_type(对应的 Pydantic RootModel 列表类型)。流程为:查询该团队在周期内新增的数据 → 封装成列表模型 →r.setex写入team_data_key(...)并设置 TTL;单团队失败只记录 warning 并跳过,不影响批次。
3.2 特殊数据源的生成逻辑
三个团队级字段没有走通用模板,值得单独说明:
- Replay 筛选器(
generate-filter-lookup):除周报 Redis 外,还连接django_redis_url指定的 Django 缓存,从PLAYLIST_COUNT_REDIS_PREFIX前缀下批量mget播放列表的录制计数(PlaylistCount),把recording_count与more_available回填到筛选器上,最后按录制数降序排列(FilterList.order_by_recording_count)。解析失败视同无计数(activities.py)。 - 即将过期的录制(
generate-recording-lookup):使用 ClickHouse 客户端执行SessionReplayEvents.count_soon_to_expire_sessions_query,参数含ttl_threshold=10(天),响应经ClickHouseResponse模型解析后取出RecordingCount(activities.py)。 - 用量趋势(
generate-usage-trends-lookup):在离线集群(Workload.OFFLINE)上执行一段内嵌的 HogQL 查询(USAGE_TRENDS_QUERY),用countIf/uniqExactIf(person_id, ...)在一次扫描中同时统计"本周 vs 上周"的事件量与去重活跃用户数;filterTestAccounts=True继承团队自己的测试账号过滤规则,使数字与团队在 Product Analytics 中看到的一致。当前周无事件(events_current == 0)的团队直接跳过不写。周环比变化率通过compute_week_over_week_change(来自 posthog/tasks/email_utils.py)计算并归一为direction(up/down/flat)与change_pct;has_baseline=False表示上周无数据可比,邮件模板会渲染"new"而非把 0→N 增长混同于"无变化"(activities.py)。
此外,generate_usage_trends_lookup有一个值得借鉴的健壮性处理:若某批次attempted > 0且全部失败(例如离线集群故障或查询语句被写坏),它会主动raise RuntimeError,让活动重试并最终让运行显式失败——否则"全部失败"会与"没有活跃团队"表象相同,从而静默发出空的 usage 区块(activities.py)。
3.3 组织级聚合
所有团队级数据生成完毕后,count_organizations统计组织总数,按batch_size切成(start, end)批次,再由asyncio.gather并发执行generate_organization_digest_batch。该活动对每个组织下的每个团队用一次r.mget批量读取全部 11 个团队级 Key,与defaults列表(空列表默认值)逐位 zip;某 Key 缺失(result is None)时用对应类型的空默认值兜底。随后组装出TeamDigest列表并写入org_digest_key(...)(activities.py)。
3.4 发送阶段
SendWeeklyDigestWorkflow先count_organizations再按组织区间切批,并发执行send_weekly_digest_batch(workflows.py)。send_weekly_digest_batch(activities.py)内部的关键逻辑:
- 从 Redis 读取
OrganizationDigest,缺失则跳过; - 内容门槛:
org_digest.is_empty() or org_digest.count_items() < DIGEST_ITEM_COUNT_THRESHOLD(阈值为 4)时跳过。注意TeamDigest._fields()有意不统计usage_trends——它是几乎每个活跃团队都有的环境性上下文,计入会把几乎所有组织都推过阈值从而人人收到邮件;而error_issues被计入,因为"出现新的生产错误"本身就值得一封邮件(types.py); - 幂等防重发:通过
MessagingRecord.objects.aget_or_create(raw_email=f"org_{organization.id}", campaign_key=input.digest.key)记录发送状态,若sent_at已存在且未设置allow_already_sent则跳过;成功后以 100 条一批abulk_update回写sent_at(RECORD_BATCH_SIZE = 100); - 遍历组织成员:读取用户的
NOTIFY_TEAMS集合(r.smembers)与PRODUCT_SUGGESTION,构建UserDigestContext,调用org_digest.for_user(user_notify_teams, user_context)得到个性化UserSpecificDigest,同样应用"空或少于 4 项则跳过"的用户级门槛; - 渲染 payload 后,实际发送并非直接调用 SMTP,而是通过 PostHog 自身的 capture 事件:
ph_client.capture(distinct_id=user.distinct_id, event="transactional email", properties=payload, groups={...})——注释说明只有 US 部署会把邮件事件转发给 customer.io(activities.py)。dry_run模式下只打印日志不发送。
四、关键类型(Key Types)
4.1Digest与WeeklyDigestInput
Digest(types.py)描述一次摘要周期,含key(如weekly-digest-2026-37)、period_start、period_end,并提供render_payload()把周期信息序列化进邮件 payload。WeeklyDigestInput(types.py)是入口工作流的 CLI 输入,支持四个开关:
| 字段 | 默认值 | 说明 |
|---|---|---|
dry_run | False | 试运行,只打日志不真正发送 |
skip_generate | False | 跳过生成阶段,直接复用 Redis 中已生成的数据 |
digest_key_override | None | 覆盖自动生成的weekly-digest-{年}-{周}Key 前缀 |
allow_already_sent | False | 允许对已发送过的组织再次发送(绕过幂等检查) |
WeeklyDigestWorkflow.parse_inputs从管理命令 CLI 接收 JSON 字符串数组并model_validate_json解析(workflows.py),这说明可通过start_temporal_workflow管理命令以 JSON 形式传入参数启动(参见 start_temporal_workflow.py 的workflow+inputs参数设计)。
4.2OrganizationDigest/UserDigestContext/UserSpecificDigest
OrganizationDigest(types.py):组织级基础摘要,含id、name、created_at与team_digests列表,整体序列化后存 Redis。其for_user(user_teams, context)方法把team_digests过滤为team_digest.id in user_teams的子集,并挂上用户上下文,返回UserSpecificDigest。UserDigestContext(types.py):所有用户级数据的容器,当前只有product_suggestion一个字段,注释明确指示"在这里添加新的用户级字段":class UserDigestContext(BaseModel): product_suggestion: DigestProductSuggestion | None = None # Add new user-specific fields hereUserSpecificDigest:发送时刻动态计算、不落 Redis的个性化视图。其render_payload(digest)遍历非空团队摘要,仅当产品建议与团队匹配(product_suggestion.team_id == td.id)时才把建议挂到该团队的 payload 上,最终输出含organization_name、organization_id、teams、scope: "user"、template_name: "weekly_digest_report"、period与digest_region(get_instance_region())的完整邮件载荷(types.py)。
4.3TeamDigest与资源模型
TeamDigest(types.py)聚合单个团队的全部摘要字段,并提供三个重要方法:
_fields():返回参与"内容计数"的资源列表(排除 usage_trends,理由见 3.4);is_empty()/count_items():判断摘要是否为空及内容条数;render_payload(product_suggestion):把各类列表model_dump()后组织成report字典(new_dashboards、new_event_definitions、new_experiments_launched等),可选用产品建议追加new_product_suggestion。
各类资源都是轻量 Pydantic 模型:DigestDashboard(name/id)、DigestEventDefinition(name/UUID id)、DigestExperiment(name/id/start_date/end_date)、DigestExternalDataSource(source_type/UUID id)、DigestFeatureFlag(name/id/key)、DigestFilter(name/short_id/view_count/recording_count/more_available)、DigestSurvey(name/UUID id/description/start_date)、DigestErrorIssue(name/UUID id,名称为空时回退为 "Untitled issue")、以及UsageTrendMetric/UsageTrends。列表形态统一用RootModel包裹(如DashboardList、FilterList),并汇总为DigestResourceTypeTypeAlias(types.py)。
4.4DigestProductSuggestion
产品推荐模型(types.py),包含team_id(建议针对的项目)、product_path(产品名,如 "Session replay")、reason_text(人类可读的推荐理由)。当 campaign 没有手写推广文案时,回退到与导航卡片一致的默认文案DEFAULT_PRODUCT_SUGGESTION_TEXT。其生成逻辑在generate_product_suggestion_lookup(activities.py):按组织缓存ProductPushCampaign(仅status=ACTIVE且started_at <= period_end的在跑 campaign,见 queries.py),通过resolve_product_path/project_uses_product(products/growth/backend/product_push/selection.py)判断团队是否已使用该产品(已使用则不推送),且尊重user.allow_sidebar_suggestions is False的退出开关,每个用户只存第一个建议。
五、如何新增一个"团队级字段"(扩展指南)
团队级字段属于"项目/团队维度的数据"(仪表盘、特性开关等)。每种字段类型应有独立的 activity 方法,以便关注点分离并支持并行执行。以新增ALERTS(告警)为例,按 README 与源码的 7 个步骤:
1. 在keys.py中为TeamDataKey增加枚举值:
class TeamDataKey(StrEnum): # ... existing keys ... ALERTS = "alerts" # new2. 在types.py中创建新数据类型的 Pydantic 模型(例如DigestAlert、AlertList),其中列表类需继承RootModel以匹配DigestResourceType。
3. 在queries.py中新增查询函数,从数据库按周期窗口取数。参考既有实现(queries.py)的约定:时间窗口统一为created_at__gt=period_start, created_at__lte=period_end(或实验类的start_date),返回带team_id的.values(...):
def query_new_alerts(period_start: datetime, period_end: datetime) -> QuerySet: return Alert.objects.filter( created_at__gt=period_start, created_at__lte=period_end ).values("team_id", "name", "id")4. 在activities.py中创建新 activity,复用通用模板generate_digest_data_lookup:
@activity.defn(name="generate-alert-lookup") async def generate_alert_lookup(input: GenerateDigestDataBatchInput) -> None: return await generate_digest_data_lookup( input, key_kind=TeamDataKey.ALERTS, query_func=query_new_alerts, resource_type=AlertList, )若需要"每个团队最多 N 条",可传per_team_limit=N(参考generate_error_issue_lookup中NEW_ERROR_ISSUES_PER_TEAM_LIMIT = 5的用法,见 activities.py)。
5. 在workflows.py中注册:把generate_alert_lookup加入GenerateDigestDataWorkflow.run的generators列表(workflows.py),并同步在init.py 的ACTIVITIES中导出。generator 列表顺序无关紧要,因为asyncio.gather会并发执行所有"区间 × 生成器"组合。
6. 把字段加到TeamDigest(types.py),并更新generate_organization_digest_batch的聚合逻辑:在r.mget的 Key 列表、defaults列表与TeamDigest(...)构造参数三处按TeamDataKey枚举顺序对齐追加——README 特别强调"all_team_data_keys中的 Key 顺序与TeamDataKey枚举顺序一致",这是最容易出错的一步(activities.py)。
7. 更新TeamDigest.render_payload():在report字典中输出新字段,邮件模板即可消费(types.py)。若新字段属于"值得单独触发邮件"的内容,记得加入_fields()计数列表;若属于"环境性上下文",则保持排除。
六、如何新增一个"用户级字段"(扩展指南)
用户级字段是"按用户个性化"的数据(产品建议、通知偏好等)。同样每种字段应有独立 activity 在生成阶段加载数据。以新增RECOMMENDATIONS为例:
1. 在keys.py中为UserDataKey增加枚举值:
class UserDataKey(StrEnum): # ... existing keys ... RECOMMENDATIONS = "recommendations" # new2. 在types.py中创建 Pydantic 模型(如UserRecommendations)。
3. 在queries.py中新增查询函数。
4. 在activities.py中创建 activity,遍历团队/用户,查询数据并写入 Redis:
@activity.defn(name="generate-user-recommendations-lookup") async def generate_user_recommendations_lookup(input: GenerateDigestDataBatchInput) -> None: # Iterate through teams/users, query data, store in Redis using: key = user_data_key(input.digest.key, UserDataKey.RECOMMENDATIONS, user.id)参考generate_user_notification_lookup(activities.py)的写法:它遍历区间内每个团队的所有有访问权限用户(team.all_users_with_access()),通过should_send_notification(user, NotificationSetting.WEEKLY_PROJECT_DIGEST.value, team.id)判断通知开关(实现在 posthog/tasks/email.py:先检查全局all_weekly_digest_disabled,再按团队检查project_weekly_digest_disabled),命中的用户用r.sadd把团队 id 加入其 NOTIFY_TEAMS 集合并设置 TTL。
5. 在workflows.py中注册:加入generators列表。
6. 把字段加到UserDigestContext(types.py):
class UserDigestContext(BaseModel): product_suggestion: DigestProductSuggestion | None = None recommendations: UserRecommendations | None = None # new field7. 在send_weekly_digest_batch中加载并挂到 context 上(参考 activities.py 中 product_suggestion 的读取范式):
raw_recommendations = await r.get( user_data_key(input.digest.key, UserDataKey.RECOMMENDATIONS, user.id) ) recommendations = UserRecommendations.model_validate_json(raw_recommendations) if raw_recommendations else None user_context = UserDigestContext( product_suggestion=product_suggestion, recommendations=recommendations, )8. 在UserSpecificDigest.render_payload()中使用,把新字段渲染进邮件 payload(types.py)。
核心不变量:for_user()的签名永远不变——所有用户级数据都经由UserDigestContext流入,这保证了OrganizationDigest的存储格式稳定、发送阶段无需感知新增字段的细节。
七、查询层与邮件触达的底层约定
7.1 查询层的过滤规则
queries.py 集中了所有数据访问,几个值得注意的约定:
query_teams_for_digest()排除内部指标组织(for_internal_metrics=True)与演示项目(is_demo=True),只加载必要列并按id排序;query_new_dashboards排除名称含 "Generated Dashboard" 的自动生成仪表盘;query_new_feature_flags排除实验自动创建的特性开关(名称含 "Feature Flag for Experiment")与调研目标 flag("Targeting flag for survey");query_saved_filters排除未命名/派生名为 "(Untitled)"/"Unnamed" 的播放列表、deleted=True的、默认播放列表(DEFAULT_PLAYLIST_NAMES)与type="collection"的,并annotate出周期内的view_count;- 实验的"启动"与"完成"是互斥的两个查询:launched 排除同时期内也完成的实验,completed 只看
end_date落在窗口内的。
7.2 发送渠道:transactional email 事件
值得强调的是,这套周报的"发送"最终落到 PostHog 自身的 product analytics 事件上(event="transactional email",按organization与instance分组),而不是直接调用邮件服务。这意味着:邮件内容、收件人、个性化数据都作为事件属性被 PostHog 记录,后续由 US 部署转发到 customer.io 真正投递;这也让整个流程天然可观测——每次发送都是一条可查询的事件。dry_run模式则让你在不上线真实邮件的情况下,通过日志审查每个用户将要收到的完整 payload。
八、小结:这套架构的可复用要点
从 posthog/temporal/weekly_digest 这个模块可以提炼出几个大规模周期性通知系统的通用设计:
- 两阶段流水线:数据生成与消息发送解耦为两个 Temporal 工作流,以 Redis 为中间层,各自可独立重试、独立扩容、按区间并行切分;
- 类型化 Key 管理:所有 Redis Key 由枚举 + 生成函数派生,杜绝硬编码散落各处;
- Pydantic 数据契约:Redis 中流动的是可校验的 JSON 模型,读取时
model_validate_json,缺省时有明确的空默认值兜底; - 内容门槛与幂等:
DIGEST_ITEM_COUNT_THRESHOLD防止"空周报"打扰用户,MessagingRecord.sent_at保证同一周期不会重复发送; - 用户个性化与数据生成解耦:
OrganizationDigest存公共数据,for_user()+UserDigestContext在发送时按用户裁剪,新增用户级字段不需要改动存储结构与for_user签名; - 健壮性细节:单团队失败跳过并告警、整批失败显式抛错(usage trends)、解析失败的缓存数据视同缺失(playlist counts)、活动全程心跳 + 重试策略。
对于任何"周期性地向海量多租户用户生成并投递个性化内容"的系统(周报、对账单、告警摘要等),这套"Temporal 编排 + Redis 中转 + Pydantic 契约 + 事件化投递"的组合都是可以直接借鉴的成熟范本。
【免费下载链接】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),仅供参考