Lambda 写入矩阵
每个 Lambda 对每张 Neon 表的字段级写入职责清单,给开发者改 Lambda 写入逻辑前对照。字段所有权是硬约束 — 违反会导致 last-write-wins 竞态或数据不一致。
::: details 实施溯源 + 落地状态(2026-05-04) 设计文档:
- Pipeline 设计:
.claude/specs/2026-03-25-task-pipeline-write-flow.md - Flow 3(AI-driven task decisions):
.claude/specs/2026-03-28-flow3-contact-analysis-design.md
Schema 状态(callytics-common 0.34.0 + migration 0014-0023):
- 主业务表
contacts/calls/messages/leads/contact_timeline/tasks/staff,辅助表rc-stores/store-config/user-stores均具备store_id列 task_events仍只有account_id,无store_id— audit 表暂保留 account-level isolation,store-level 终态后续补contacts复合主键 =(phone, store_id),store_id NOT NULL(migration 0014)staff.store_idNOT NULL(migration 0015)site_id→account_idrename(migration 0016)- 其余表
store_id仍 nullable(calls/leads/messages/contact_timeline/tasks),writer 全写但保留 NULL 兼容历史行,各表均有idx_<table>_store_id WHERE store_id IS NOT NULLpartial index
Lambda writer 当前 store_id 写入状态(2026-05-04):
transcribe-processor→calls.store_id+contact_timeline.store_id(call.status_changed,Phase 1 起 dual writer)ai-analysis-processor→contacts/calls/contact_timeline.store_idcontacts-analyzer→contacts/tasks/contact_timeline.store_idmessage-processor→messages/contacts/contact_timeline.store_idlead-tracking(独立 repo)→leads/contacts/tasks/contact_timeline.store_id
franchise_id 来源(读源代码核对):
message-processor:parseFranchiseFromClientId(_client_id, _account_id)解析,不是硬编码ai-analysis-processor:从 SQS 上游 metadata 传入contacts-analyzer:从 contact 行的franchiseId读取lead-tracking:StoresV2 leadEmail 映射 →parseBrandFromFranchise()- 当前所有 writer 都从外部数据派生
franchise_id,不硬编码(历史 P0 #632 fix:阻止"_client_id整串含 site 后缀直接写入"的 bug)。
:::
全局数据流
图例:实线 = 写入 · 虚线 = 读取 / DEFERRED · 🆕 = Task 系统新增 · * = DEFERRED(V1 不实现)
一、写入模式总览
1.1 所有 Lambda 统一用 Neon HTTP driver + Drizzle ORM
连接方式:@neondatabase/serverless HTTP driver → drizzle-orm/neon-http。每条 SQL 是一个 HTTP POST。
1.2 多表写入:db.batch() vs 分开 await
db.batch()底层调用 Neon 的sql.transaction()— 所有 SQL 打包成 1 个 HTTP 请求,DB 内执行 BEGIN/COMMIT/ROLLBACK。
1.3 五个 Flow 写入模式
所有 Flow 的 Neon 写入统一用 db.batch() 打包(2026-03-28 Neon-primary 决策)。DDB 写入是过渡期冗余,独立标记不参与 batch;SQS 发送在 batch 外,batch 成功后才发,非阻塞。
1.4 踩坑注意事项
::: warning db.batch() 生产注意(社区踩坑总结)
db.batch()内不能Promise.all— 并行语句会丢 transaction scope(Drizzle #2200)$onUpdate(() => new Date())在 batch 内不生效 — Drizzle 的$onUpdate是 query builder hook,不是 DB trigger。batch 内所有 UPDATE 必须显式写updatedAt: sql\NOW()``ON CONFLICT DO NOTHING不导致 batch 回滚 — PostgreSQL 的 conflict 处理是静默跳过,不抛错。dedup index 在 batch 内安全工作- Pin driver 版本 —
@neondatabase/serverless和drizzle-orm之间有版本兼容性问题,社区报告升级后运行时报错(Drizzle #5208) - SQS 提供 Lambda 层重试 —
db.batch()解决 SQL 原子性,SQSbatchItemFailures解决 Lambda crash 重试。两者配合 = 幂等 + 原子
:::
二、总览矩阵
✅ = 现有 🆕 = V1 新增 * = DEFERRED — = 不读 Neon(数据全在触发输入里)
Per-Call Analysis 读取:当前代码从 DDB 读(
customer-history.ts),计划迁移到 Neon。如果customerSummary存在就读,不存在就跳过(graceful degradation)。Contact Analysis 读取:customerSummary 是 compressed memory。
lastContactAnalysisAt为 NULL(首次)→ 读 30d 全量建立基线。有值 → 读lastContactAnalysisAt之后的增量。详见 contact-analysis-writes.md。
三、Contact Analysis 触发源
6 个触发源,统一发 SQS 消息到同一 schema(Zod-validated by SqsMessageBodySchema):
franchiseId/accountId不再传 — Consumer(contacts-analyzer)从 contact 行通过 PK lookup 自行读取。上游 producer 只需要phone + storeId。
技术约束:SQS ESM filtering 只支持 body key(不支持 MessageAttributes)。
FIFO Queue 已上线 — dailyBatchQueue 从 Standard 切换为 FIFO 并迁移到
SqsStack。触发源发送消息必须提供MessageGroupId和MessageDeduplicationId,例如 ai-analysis-processor 用${phone}#${storeId}作 group key、per_call_analysis-${phone}-${callId}作 dedup id。
* expressQueue 尚未创建(Phase D5)。V1 期间 ❷❸❹ 用 dailyBatchQueue。 ❷ On-demand 后端 API 已存在,但前端 Refresh AI 按钮随员工 UI 一起 DEFERRED。
字段来源
V1 不做 Cooldown
不做应用层 cooldown,重复分析由 DB index 保护(uq_tasks_pending_contact_category 防重复 task,uq_tasks_analysis_run_category 防 SQS 重试)。成本 ~200 calls/day × $0.005 ≈ $1/day,可接受。lastContactAnalysisAt 每次分析仍写 NOW(),未来加 cooldown 时可直接读取。
T5/T6 推荐实现(DEFERRED)
Contact Analysis 自检测 leadStatus 变化 → EventBridge Scheduler 延迟触发(showed=4h, trialed=2-4h) — SQS DelaySeconds 最大 15min 不够用。Loop guard:source 已是 showed/trialed 时不再检测。
四、contacts 表写入详情
contacts 是中心化客户档案,4 个 Lambda 各写不同字段。
字段所有权
Per-Call Analysis 只写"活动标记"(什么时候发生了什么),Contact Analysis 只写"AI 判断"(客户状态是什么) — 两者字段互不重叠。
三个 UPSERT contacts 对比
lead-tracking / Per-Call Analysis / message-processor 都 UPSERT contacts,但写的字段不同 — 数据源决定能写什么。
INSERT 时写入:
ON CONFLICT 时更新:
详见各 per-Lambda 文件的 Step ② section。
五、contact_timeline 表写入详情
contact_timeline 记录所有客户互动和状态变更事件,5 个 Lambda 共写 9 种事件类型。每条都包含公共必填字段(contactPhone / franchiseId / siteId / eventType / eventCategory / occurredAt / actorType)。
* = DEFERRED
occurred_at 规则:非 transaction 写入(lead / call / SMS)用源记录时间防排序漂移;transaction 内(Contact Analysis Step ⑥)用
NOW(),所有语句共享同一时间。错误处理:5 个 Lambda 的 timeline INSERT 都在
db.batch()内,batch 阻塞(失败 → 整个 batch 回滚 → SQS 重试)。lead-tracking 的 prod Neon 写入阻塞,test Neon fire-and-forget;其他 Lambda 的 SQS send / DDB 旁路写入是非阻塞。studio-api Close 当前是 raw SQL × 2(UPDATE tasks + INSERT task_events)未原子化;不写 contact_timeline。
六、字段清理计划
有了 contact_timeline 后,contacts 表上部分字段可从 timeline 派生 → 计划移除。temperature 相关字段和 timeline 事件也将一并移除。
已移除的 contacts 字段
待移除的 contacts 字段(需 engineer review,callytics-common#32)
Phase 1:timeline 有数据后可移除
Phase 2:Task Pipeline 完成 + 前端迁移后可移除
待移除的 timeline 事件
— 随 temperature 字段一起移除。contact.temperature_changed
为什么 timeline 比 contacts 列更好
重复进入场景(客户 3 月来 → 冷却 → 9 月再来)时,contacts 列的 acquiredAt 永远冻结在 3 月。Timeline 可查当前周期的 lead.created,准确算 Speed to Lead。
相关决策变更
- typeCategory 门控:
temperature≠cold→ 改为 AI prompt 指令("lastActivityAt > 5d 不建 lead_follow_up") - Stale override:>30d → neglected → 改为"30 天无活动 + 无 pending task"查询
当前状态:新代码(lead-tracking / message-processor)不写这些字段。Contact Analysis 现有代码仍写,等 timeline 有数据 + 验证后停止。
七、已知问题与代码隐患
Expert Panel 代码验证发现 10 个待修复问题,按严重度排序。