Friend 后端每日摘要收件人选择切换(Daily-Summary Cutover):从代码层过滤迁移到 Firestore 索引查询的完整实战指南
2026/9/15 16:29:34 网站建设 项目流程

Friend 后端每日摘要收件人选择切换(Daily-Summary Cutover):从代码层过滤迁移到 Firestore 索引查询的完整实战指南

【免费下载链接】FriendAI that sees your screen, listens to your conversations and tells you what to do项目地址: https://gitcode.com/GitHub_Trending/fr/Friend

Friend(Omi)是一个"看见你的屏幕、听见你的对话、并告诉你该做什么"的 AI 产品。其后端每天按用户本地时区的设定时间(默认 22:00)为每位用户生成并推送一份"每日摘要"(Daily Summary)。这项能力由一个每小时运行一次的notifications-job定时任务驱动。随着用户规模增长,收件人筛选环节的读取成本成为瓶颈:任务每次执行都要全量扫描所有携带time_zone的用户文档,再在 Python 里逐条过滤daily_summary_enabled/daily_summary_hour_local两个字段。

本文以仓库中的运维手册 backend/docs/runbooks/daily-summary-cutover.md 为主线,完整讲解 Friend 后端如何通过"写时默认值 + 一次性回填 + 服务端复合索引查询 + 分阶段开关"四个步骤,把收件人选择从全量扫描(legacy)切换到索引查询(indexed)的整套工程实践。读完本文,你将掌握:两个字段默认值为何必须在写入时落盘、回填脚本的正确用法与幂等性保证、三种选择模式的语义与日志判读方法、以及分五步安全上线的完整顺序。

背景:为什么需要这次切换

切换前的实现位于 backend/database/notifications.py 的get_users_for_daily_summary函数(本文后续称"legacy 选择器")。它的逻辑是:

  1. 按"当前本地小时"把所有 IANA 时区分成若干组(_get_timezones_grouped_by_hour,见 backend/utils/other/notifications.py);
  2. 对每一组,用time_zone IN chunk查询users集合(FirestoreIN查询最多支持 30 个值,因此时区列表要按 30 切块);
  3. 在 Python 端逐条过滤:daily_summary_enabled显式为False的跳过;daily_summary_hour_local未设置则用代码默认值22判断。

问题在于:大多数用户文档上根本没有这两个字段,它们的默认值(True22)只存在于代码中。这意味着:

  • Firestore 无法用索引直接返回"符合条件的收件人",因为过滤条件是代码运行时算出来的;
  • 每次任务执行都必须做一次users集合的全量读取(凡带time_zone的文档都会被读出),把不满足条件的文档在内存里丢掉;
  • 每小时一次全量扫描,随着用户量增长,read_ops_count会持续累积,成本不可忽视。

切换的目标,就是让 Firestore 直接通过==等值过滤返回收件人集合,把每次执行的读取量从"全部用户"降到"仅收件人"。

四块积木:切换由哪些代码构成

整个方案由四个部分协同组成,对应手册中的 "Pieces" 一节。

1. 写时默认值(Write-time defaults)

在 backend/database/notifications.py 中定义了两个常量与一个辅助函数:

  • DEFAULT_DAILY_SUMMARY_HOUR_LOCAL = 22:默认本地发送小时(晚上 10 点);
  • DEFAULT_DAILY_SUMMARY_ENABLED = True:默认开启;
  • daily_summary_schedule_defaults(user_data):返回"当前文档缺失字段"的默认值补丁。
def daily_summary_schedule_defaults(user_data: Mapping[str, Any]) -> Dict[str, Any]: """Return write-time defaults for whichever schedule fields are absent. Present values, including explicit ``False`` and hour ``0``, are never included. Empty dict when both fields are already on the document. """ patch: Dict[str, Any] = {} if 'daily_summary_enabled' not in user_data: patch['daily_summary_enabled'] = DEFAULT_DAILY_SUMMARY_ENABLED if 'daily_summary_hour_local' not in user_data: patch['daily_summary_hour_local'] = DEFAULT_DAILY_SUMMARY_HOUR_LOCAL return patch

关键不变量:已存在的值绝不覆盖,包括显式的False和小时0——它们都是合法的用户偏好,必须保留。凡是会写入time_zone或偏好字段的写路径,都会顺带补齐缺失的调度字段:

  • save_token(FCM token 注册时写入time_zone)在 backend/database/notifications.py 中通过schedule_patch.update(daily_summary_schedule_defaults(user_data))补齐;
  • set_user_time_zone_if_missing在 backend/database/notifications.py 中同样带上了daily_summary_schedule_defaults
  • set_daily_summary_hour_local/set_daily_summary_enabled两个偏好写入口也复用同一辅助函数(backend/database/notifications.py)。

由此形成不变量:一份携带time_zone的用户文档必然同时携带两个调度字段,索引查询才可能匹配到它。模块顶部 docstring 也明确记录了这条约定(backend/database/notifications.py)。

2. 索引化选择器(Indexed selector)

新增的get_users_for_daily_summary_indexed位于 backend/database/notifications.py。它按小时组、并按不超过 30 个时区切块后,执行服务端等值过滤:

daily_summary_enabled == True AND daily_summary_hour_local == H AND time_zone IN chunk

这段查询不是手写的,而是由索引注册表生成的。在 backend/database/firestore_index_registry.py 中注册了标识符为DAILY_SUMMARY_RECIPIENTS_QUERYFirestoreQuerySpec,其过滤器为daily_summary_enabled ==daily_summary_hour_local ==time_zone in,并声明了对应的复合索引字段。该注册表的定位是"仓库拥有 Firestore 查询与索引需求":一条注册的查询规格既能构建生产查询,又声明了这条查询所需的精确复合索引;复合索引随后被生成到仓库根目录的firestore.indexes.json清单中。

FirestoreQuerySpec.build()(backend/database/firestore_index_registry.py)负责把声明式过滤器翻译成真实的 Firestore 查询。索引化选择器调用它:

query = DAILY_SUMMARY_RECIPIENTS_QUERY.build( db.collection('users'), {'enabled': True, 'hour_local': target_local_hour, 'time_zones': chunk}, field_filter_factory=FieldFilter, )

索引化选择器在其余行为上与 legacy 完全对齐:同样按 30 个时区切块、同样读取fcm_tokens子集合与遗留fcm_token字段、同样保留无 token 用户(仅桌面端使用、没有 FCM token 的用户不能被丢弃)、同样对单个 chunk 异常做 try/except 记录并继续。这一点在它的 docstring 中明确说明(backend/database/notifications.py)。

一个重要边界:时区 → 本地小时的换算仍然在运行时用 tzdata(pytz)求值,由_get_timezones_grouped_by_hour完成(backend/utils/other/notifications.py)。也就是说,索引只负责"该小时组内的收件人"筛选,DST(夏令时)切换无需做任何数据对账——每个小时组都是根据"此刻各时区的真实本地小时"动态算出来的。

3. 回填脚本(Backfill)

存量用户文档上并没有这两个字段,需要一次性回填。脚本位于 backend/scripts/backfill_daily_summary_schedule_fields.py,其特性如下:

  • __name__分页扫描users集合,并用select(list(_PROJECTED_FIELDS))做字段投影,只读需要的字段(backend/scripts/backfill_daily_summary_schedule_fields.py),避免拉取整份文档;
  • 只写缺失字段:对每个文档调用daily_summary_schedule_defaults,patch 为空则跳过;写操作通过set(..., merge=True)完成,绝不会覆盖已有偏好(backend/scripts/backfill_daily_summary_schedule_fields.py);
  • 默认 dry-run:不传--apply时只统计不写入;
  • 幂等:重复执行结果一致,可安全重跑;
  • 可续跑:通过--start-after <uid>从上次的last_uid继续;
  • 可限流--limit N用于小规模冒烟。

命令行用法:

# 冒烟:只扫描前 50 个用户,dry-run 只统计 python scripts/backfill_daily_summary_schedule_fields.py --limit 50 # 全量回填(dry-run 后确认 missing_enabled / missing_hour 数量再执行) python scripts/backfill_daily_summary_schedule_fields.py --apply # 中断后续跑 python scripts/backfill_daily_summary_schedule_fields.py --apply --start-after <uid>

dry-run 的输出包含scannedwith_time_zonemissing_enabledmissing_hourwrittenlast_uidapply等统计项(backend/scripts/backfill_daily_summary_schedule_fields.py),其中missing_enabled/missing_hour是判断回填规模的关键指标。预期结果是每个用户文档大约一次写操作

4. 模式开关(DAILY_SUMMARY_SELECTION_MODE)

运行期通过环境变量DAILY_SUMMARY_SELECTION_MODE控制选择器行为,读取逻辑在 backend/utils/other/notifications.py:

  • legacy代码默认值)——使用全量扫描选择器,即切换前的行为。环境变量未设置时,_selection_mode_from_env()返回'legacy'
  • shadow——以 legacy 结果为权威进行发送,同时额外运行索引化选择器,并按小时组打印一条对比日志;代价是每次执行多一次 legacy 全量扫描;
  • indexed——仅运行索引化选择器。

值会做strip().lower()归一化;遇到未知值则打印daily_summary_selection_mode_unknown告警并回退到legacy(backend/utils/other/notifications.py)。可选值集合_DAILY_SUMMARY_SELECTION_MODES = frozenset({'legacy', 'shadow', 'indexed'})(backend/utils/other/notifications.py)。

_get_users_for_daily_summary(backend/utils/other/notifications.py)中,模式决定实际使用的 selector:

selector = notification_db.get_users_for_daily_summary if DAILY_SUMMARY_SELECTION_MODE == 'indexed': selector = notification_db.get_users_for_daily_summary_indexed

shadow模式下,legacy 结果照常返回,同时并行运行索引化选择器,并调用_log_daily_summary_selection_shadow(backend/utils/other/notifications.py)打印集合差集日志:

daily_summary_selection_shadow hour=... legacy=N indexed=M only_legacy=K only_indexed=J sample_only_legacy=[...] sample_only_indexed=[...]

各字段含义:

字段含义
hour当前处理的小时组
legacy/indexed两种选择器各自返回的用户数
only_legacy只在 legacy 结果中、索引查询漏掉的用户数(回填或漏斗的漏网之鱼)
only_indexed只在索引结果中、legacy 漏掉的用户数
sample_only_legacy/sample_only_indexed差集用户 uid 的前 5 个样本,便于定位

若索引化查询本身失败,会打印daily_summary_selection_shadow_failed告警并继续使用 legacy 结果(backend/utils/other/notifications.py),保证 shadow 模式下的发送始终以 legacy 为权威。

运行锁(Run lock)

切换上线涉及环境变量变更,为防止 Cloud Run Scheduler 的重复触发导致任务重叠执行,start_cron_job(backend/utils/other/notifications.py)在每次执行时通过 RedisSET NX EX抢占notifications_job:run_lock

  • 锁 TTL 为55 分钟
  • 释放采用token 比对 + 删除(compare-and-delete),避免误删他人持有的锁;
  • 抢占失败(重叠触发)时打印notifications_job_run_skipped reason=overlap并跳过通知段,已有的 checkpoint 会在下次执行时续跑尾部;
  • 抢占阶段 Redis 出错则fail open:打印notifications_job_run_lock_acquire_failed、记录daily_summaryfallback 计数器后继续执行,代价是可能重复一次遍历(由每用户按日锁兜底)。

三种选择模式的源码验证

模式选择的行为有单元测试覆盖,见 backend/tests/unit/test_daily_summary_selection_mode.py:

  • test_selection_mode_unset_defaults_to_legacy:未设置环境变量时默认legacy
  • test_selection_mode_normalizes_indexed'indexed'与带空格的' indexed '均被归一化为indexed
  • test_selection_mode_shadowshadow生效;
  • test_selection_mode_garbage_falls_back_to_legacy_with_warning:垃圾值回退legacy并告警;
  • test_legacy_mode_calls_only_legacy_selectorlegacy模式下索引选择器绝不会被调用(否则触发AssertionError);
  • test_indexed_mode_calls_only_indexed_selectorindexed模式下 legacy 选择器绝不会被调用;
  • test_shadow_mode_returns_legacy_and_logs_set_diff:shadow 模式返回 legacy 结果并打印集合差集日志。

索引化选择器本身还有专门的测试文件 backend/tests/unit/test_daily_summary_indexed_selection.py。这些测试把"每种模式下谁会被调用、谁绝不会被调用"固化为契约,是切换安全性的第一道防线。

分五步的上线顺序(Cutover order)

手册给出了严格的切换顺序,每一步都有可验证的判据:

第 1 步:部署写时默认值

先部署后端,让写时默认值逻辑生效。定时任务本身不做任何改动——此时DAILY_SUMMARY_SELECTION_MODE未设置,仍走legacy。这一步只保证"从今以后新写入的文档都带齐两个字段"。

第 2 步:对生产执行一次回填

先 dry-run 并读取missing_enabled/missing_hour两个统计值,确认规模符合预期后执行--apply。预期结果:每个用户文档大约一次写操作。回填完成前,索引化选择器会漏掉所有缺字段的存量文档,所以这一步是后续 shadow / indexed 生效的前提。

第 3 步:shadow 模式对拍

notifications-job的 runtime-env overlay 中把DAILY_SUMMARY_SELECTION_MODE设为shadow,部署后读取一次执行的 shadow 日志。判据:

only_legacy必须对每个小时组都为0。只要出现非零值,就说明存在回填或写路径漏斗漏掉的文档,需要追查并补齐。

在仓库当前的 backend/deploy/runtime_env.yaml 中,notifications-jobDAILY_SUMMARY_SELECTION_MODE已配置为indexed(dev 与 prod 两个环境均为如此,见该文件中notifications-job的 env 段)。这意味着当前线上已经完成切换;如果你需要复现完整流程,可在自己的部署 overlay 中依次设置shadow→ 验证 →indexed

第 4 步:indexed 模式正式生效

DAILY_SUMMARY_SELECTION_MODE设为indexed并部署。判据:观察users集合的read_ops_count,应从每次执行一次全量扫描,降到仅读取收件人文档。这是整个切换的核心收益:每次小时任务只读取该小时组实际需要推送的用户。

第 5 步:后续清理

切换完成后有两项后续工作:

  1. 移除 legacy 选择器get_users_for_daily_summary的全量扫描路径);
  2. 注意保留 08:00 的晨间推送_get_users_in_timezones——它每天仍会读取约 1/24 的用户量(按目标小时筛选时区),这是独立于每日摘要的"穿戴提醒"类推送路径(backend/utils/other/notifications.py),不属于本次清理范围。

设计要点总结

关注点方案实现位置
默认值落盘写时补缺、已存在值永不覆盖backend/database/notifications.py
收件人查询服务端==等值过滤 + 30 时区切块backend/database/notifications.py
索引声明注册表驱动,生成到firestore.indexes.jsonbackend/database/firestore_index_registry.py
存量回填投影分页、只写缺失、幂等、dry-run 默认backend/scripts/backfill_daily_summary_schedule_fields.py
模式切换DAILY_SUMMARY_SELECTION_MODE:legacy / shadow / indexedbackend/utils/other/notifications.py
重叠防护Redis 55 分钟运行锁,token 比对释放,fail-openbackend/utils/other/notifications.py
行为契约模式选择与索引查询的单元测试backend/tests/unit/test_daily_summary_selection_mode.py、backend/tests/unit/test_daily_summary_indexed_selection.py

这套方案的核心工程思想值得借鉴:先把隐式默认值物化成显式数据,再让存储层完成过滤,用影子模式对拍新旧实现,最后以观测指标(read_ops_countonly_legacy)作为切换完成的判据。对于任何"过滤逻辑在应用层、导致全量扫描"的定时任务,这都是一个可复用的降本切换范式。

【免费下载链接】FriendAI that sees your screen, listens to your conversations and tells you what to do项目地址: https://gitcode.com/GitHub_Trending/fr/Friend

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

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

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

立即咨询