Agent-Reach:智能体触达引擎的决策、频控与降级实践

Agent-Reach:智能体触达引擎的决策、频控与降级实践 做用户触达这行的朋友大概都经历过这种场面运营在后台配规则表配到第四十几条的时候已经没人说得清哪条跟哪条会撞车。用户早上收到一条您有优惠券即将过期中午收到一条好久不见回来看看晚上又来一条完成问卷领积分三连击下去退订率当天就给你颜色看。Agent-Reach 就是我在这个背景下折腾出来的一套智能体触达引擎它要解决的核心问题非常朴素——把什么时候、用什么渠道、跟谁、说什么话这套决策从人肉维护的规则表里挪到一个能感知、能规划、能复盘、还能被硬约束管住的 Agent 手里。名字里的 Reach 就是触达的意思但它不是又一个消息推送 SDK也不是对话机器人框架。它更像一个触达中台上面接着业务场景下面接着短信、邮件、站内信、企业 IM 这些具体通道中间那层由 Agent 负责排序、择时、组稿和节流。写这篇东西的目的很直接如果你正在做用户增长中后台、消息推送系统或者在做面向 C 端的 Agent 应用这套设计思路和踩坑记录应该能让你少走几个月弯路。整套方案我按整体设计 → 核心机制 → 落地实操 → 问题排查 → 调优经验的顺序展开代码给的是可直接跑的骨架参数给的是我实际用过的取值你可以按自己的量级裁剪。1. 项目整体设计与思路拆解1.1 触达这件事到底难在哪很多人以为触达系统的技术难点在怎么把消息发出去其实发送本身是最简单的一环接个 SDK 就完事了。真正难的是四件事决策、节流、组稿、归因。决策难在组合爆炸。一个稍微复杂的业务就有几十个触达场景每个场景有自己的触发条件、目标人群、优先级。你用规则表配规则之间会互相干扰你用纯模型排序又会出现模型今天心情不好一条都不发这种没法向业务交代的情况。节流难在它是多维度的。用户维度每人每天几条、通道维度供应商 QPS 上限、场景维度同一场景 7 天内不重复、全局维度大盘突刺要能刹住车这四个维度必须同时满足任何一个漏掉都会出事。我见过最典型的翻车就是只做了用户维度频控结果一次运营活动把短信通道打爆通道商直接限流连验证码都发不出去。组稿难在个性化与可控性的平衡。用 LLM 生成文案确实能做到千人千面但你没有约束的话它会给用户编出不存在的优惠金额、写错订单号、甚至生成一段语气奇怪的话。归因难在数据链路长。消息发出去到用户点开、到最终转化中间隔着好几个系统口径对不上是家常便饭。Agent-Reach 的整体设计就是围绕这四个难点来的。1.2 四层架构与各自的职责边界我把整个系统拆成四层层与层之间只通过明确定义的接口通信这样任何一层要换实现都不会牵连其他层。层次职责关键产出换了会牵连谁渠道适配层抹平各家通道差异处理重试、幂等、回执统一的 SendResult不牵连上层逻辑决策层判断这条任务该不该发、什么时候发、走哪个通道决策日志 放行/拒绝原因上层业务无感内容层组装最终文案做变量替换和安全校验RenderedContent只影响文案质量治理层频控、熔断、灰度、审计、指标采集大盘监控与配额全局生效先说渠道适配层。这一层的价值在于把短信有长度限制且要分片计费、邮件有 HTML 模板和退信机制、站内信有已读回执、企业 IM 有机器人频率限制这些差异全部吃进去对上层只暴露三个方法send、health、estimate_cost。为什么要暴露 estimate_cost因为决策层需要知道发这一条要花多少钱在配额紧张的时候成本是重要的排序因子。决策层是整个系统的脑子。我采用的是规则做硬约束、模型做软排序的混合架构。所有会出事的判断静默期、黑名单、频控超限、任务过期全部用确定性规则实现这些判断不允许模型插手而那些发哪条更好什么时间打开率更高的判断交给一个轻量排序模型或者干脆是一个打分的函数。这样做的原因是模型可以错但底线不能破。内容层需要理解一个关键点——它不是简单的模板渲染引擎。它要处理变量缺失、长度超限、敏感词、链接白名单、多语言、降级链。我给它设计了三段式降级LLM 生成 → 模板渲染 → 干脆不发。中间任何一步失败就往下掉一级绝不能出现因为组稿失败而阻塞了整个发送队列的情况。治理层是最容易被忽略但最重要的一层。它管着四件事配额、熔断、灰度、审计。其中审计我要特别强调每一条发出的消息都必须能在日志里追溯到哪个场景、哪个决策、用了哪版模板、通过了哪些校验这是出问题时唯一能救你的东西。1.3 关键选型背后的取舍逻辑这里我把自己做过的几个选型决策和当时的思考过程摊开说因为选型这东西没有标准答案只有合不合适。消息队列选 Kafka 还是 Redis Streams我的选择是 Kafka 做主链路Redis 只做配额计数和缓存。原因很实际触达任务需要可重放。某次内容校验逻辑写错了把一批正常的消息拦了下来这时候你需要把过去 6 小时的任务重新投递一遍Kafka 的分区偏移量可以让你精确地做到这一点Redis Streams 虽然也支持但持久化和堆积能力在千万级任务下会吃力。代价是运维复杂度上去了如果你的日发送量在百万级以下Redis Streams 完全够用别为了架构好看给自己找麻烦。为什么用状态机而不是纯事件流一条触达任务从创建到最终送达会经历待决策 → 已放行 → 已发送 → 已送达 → 已点击 → 已转化这些状态也可能走到已拒绝已过期发送失败。如果只记录事件不维护状态你每次想知道这个用户到底收到没有就得重放全部事件查询成本极高。我在任务表上保留了一个 state 字段作为物化的当前状态事件流另存一张表用于审计两者通过同一个 task_id 关联。这是典型的空间换时间。为什么必须单独建决策日志表刚开始我把拒绝原因直接写在任务表的一个字段里后来发现业务方天天来问为什么这个用户没收到一个用户可能同时命中五条规则只记一个原因根本说不清。改成独立日志表后一次决策可以写多行记录每条规则的判定结果和当时的关键变量比如当天已发 2 条剩余配额 1排查效率提升非常明显。2. 核心机制拆解与关键参数怎么定2.1 渠道适配层的统一抽象怎么写先说接口定义。用一个抽象基类约束所有通道强制实现四个方法这是保证上层逻辑不被污染的关键。from abc import ABC, abstractmethod from dataclasses import dataclass dataclass class SendResult: ok: bool channel_msg_id: str retryable: bool False # 是否值得重试网络抖动True号码无效False error_code: str cost_cent: int 0 # 本条成本单位分 class ChannelAdapter(ABC): name: str daily_qps: int 100 # 该通道的对外承诺 QPS用于二级令牌桶 abstractmethod def send(self, user_id: str, content, idem_key: str) - SendResult: idem_key 由上层生成通道侧必须用它做去重 abstractmethod def health(self) - bool: 返回 False 时决策层会触发熔断不再向该通道派发 abstractmethod def estimate_cost(self, user_id: str, content) - int: 预估成本用于配额紧张时的排序 def recall(self, channel_msg_id: str) - bool: 可选实现站内信支持撤回短信基本不支持 return False这里有两个设计细节值得展开。第一是retryable 字段。接口返回的失败不能一概而论号码格式错误这种失败重试一万次也没用反而会污染通道的质量评分而超时、限流这类失败是值得退避重试的。我最初的版本只有一个 ok 布尔值结果重试队列里堆了一大堆永远不可能成功的任务白白占资源。第二是idem_key 的生成规则。这个键必须满足同一业务意图重放多次产生同一个键。我的做法是md5(f{user_id}:{scene}:{biz_id})其中 biz_id 是业务侧的唯一标识比如订单号、活动 ID如果没有明确的业务 ID就退化成md5(f{user_id}:{scene}:{日期})保证同一天同一场景只发一次。千万不要用时间戳或者随机数做幂等键那等于没做。注意幂等键的判重存储不能只放在通道适配层。我强烈建议在决策层入口就做一次判重因为有些通道尤其是自建的站内信根本没有幂等能力等消息到了适配层再去拦可能已经产生副作用了。2.2 频控、静默期与优先级抢占的参数怎么算这一节是整套系统里最容易算错的地方我把我的计算方法完整写出来。先确定约束条件。假设我们的目标人群是 200 万活跃用户业务侧约定每人每天最多收到 3 条营销类消息同时考虑实际的拒绝率频控拦截、静默期、内容校验失败大约会拦掉 40%那么计划投递量要按 1000 万来准备而不是 600 万。这个 1.67 倍的冗余系数是我实测出来的经验值真实环境里拦截率通常在 30% 到 50% 之间波动。然后是投递窗口。营销类消息集中在 9:00 到 21:00扣掉午休和晚饭时段大约 2 小时有效窗口约 10 小时也就是 36000 秒。平均速率 1000 万 / 36000 ≈ 278 条/秒。但流量从来不是均匀的早高峰和晚高峰的瞬时峰值能到平均值的 3 倍左右所以峰值速率要按 850 条/秒设计。令牌桶的参数就这样定速率设为需求峰值的 1.1 倍桶容量设为 10 秒的流量。也就是速率 300 令牌/秒容量 3000。容量的选择逻辑是这样——容量太小比如 300会导致突发流量一来就被打散任务出现明显延迟容量太大比如 30000则失去削峰意义通道商那边照样被打爆。10 秒是我试过比较舒服的折中值。二级令牌桶是必须的。除了系统整体的桶每个通道还要有自己的桶速率取通道商承诺值的 80%。比如短信供应商承诺 500 QPS通道桶就设 400 令牌/秒站内信是自己家的服务可以设到 2000。留 20% 余量的原因是通道商的实际能力会波动而且它可能同时在服务其他业务方你按满额打过去最先被限流的一定是你。静默期的处理要区分优先级。默认 22:00 到次日 08:00 不发营销消息但验证码、支付结果、物流异常这类事务型消息必须能穿透静默期。所以我在决策层做了一个判断只有 priority 为 critical 的任务才能无视静默期high 及以下一律推迟到静默期结束后再投递而不是直接丢弃。优先级抢占是个容易被设计错的机制。我的做法是当一条 critical 任务到达、但通道桶的令牌已经耗尽时允许它借令牌同时把队列尾部等量的低优任务顺延到下一个时间片。实现上就是给 critical 任务一个独立的备用桶容量是主桶的 20%平时不用只在紧急时消耗。用户维度的配额用滚动窗口而不是自然日。为什么自然日在零点重置会诱导用户集中在 23:50 和 00:10 收到两条消息体验很差。改成 24 小时滚动窗口用 Redis ZSET 记录每次发送的时间戳score 为时间戳查询时移除超过 24 小时的成员并统计剩余数量就能规避这个问题。2.3 内容生成的三段式降级与校验器内容层的核心不是生成得多漂亮而是生成得多安全。我给它加了三道闸。第一道是槽位约束。所有需要 LLM 填写的部分我都不让它自由发挥而是给定固定槽位和取值范围。比如优惠金额这个槽位我会把业务系统里的真实数值作为上下文传进去并在 prompt 里明确你只能使用给定的数值不得自行计算或推断。这比事后再校验要可靠得多因为大部分幻觉在源头就被掐住了。第二道是确定性校验器。不管你用什么方式生成最终文案必须过一遍校验函数import re SENSITIVE {最低价, 绝对, 百分百} # 示例词表实际按业务规范维护 LINK_WHITELIST (https://example.com/, https://m.example.com/) def validate(text: str, max_len: int 70) - tuple[bool, str]: if not text or not text.strip(): return False, empty_content if len(text) max_len: return False, too_long if re.search(r\{\{.*?\}\}, text): return False, unresolved_placeholder # 变量没替换掉 for w in SENSITIVE: if w in text: return False, fsensitive:{w} for url in re.findall(rhttps?://[^\s], text): if not url.startswith(LINK_WHITELIST): return False, link_not_allowed return True, okunresolved_placeholder这条检查帮我拦掉过最尴尬的事故——模板变量因为用户资料缺失没替换成功用户收到一条亲爱的{{nickname}}您的订单已发货。第三道是降级链。LLM 组稿失败或校验不通过时直接回退到该场景预设的静态模板模板里只有最基础的变量昵称、时间、金额保证一定能出内容如果静态模板也渲染失败比如用户没有任何可用昵称那就放弃这条消息写入content_failed状态不阻塞队列。实操心得给 LLM 组稿设一个硬超时我这里设的是 800 毫秒。超过就直接走降级链。原因是大促期间模型服务的排队延迟会飙升如果不设超时整个发送链路会被拖垮而一条稍微不那么个性化的文案远好过一条迟到的文案。2.4 归因口径与指标定义指标体系我分成三层每层的定义必须写死在文档里不然各部门对数字的争吵会永无止境。指标定义计算方式注意事项触达率成功送达设备的任务占比送达数 / 放行数站内信无送达回执用拉取到近似打开率消息被点开的任务占比打开数 / 送达数要把 App 冷启动拉起算进去口径要统一转化率触达后产生目标行为的占比转化数 / 送达数必须和自然转化率做差退订率触发退订或关闭推送的占比退订数 / 送达数这是最重要的负向指标超过阈值要自动降频投诉率被用户主动举报的占比投诉数 / 送达数通道商会据此调整你的通道质量分级关于转化率我要强调一个很多人会忽略的做法必须留对照组。我通常会随机留出 1% 的用户不做任何触达用这 1% 的自然转化率作为基线。否则你永远不知道那 5% 的转化率里有多少是消息带来的有多少是用户本来就要下单的。这个 1% 的留白几乎不影响整体收益但给你的决策依据带来的价值巨大。指标采集的埋点要跟着 task_id 一路走。我见过太多系统在归因阶段靠 user_id 时间窗口去匹配结果同一天发了两条消息点击到底算谁的完全说不清。正确的做法是每条消息的落地页链接里都带上 task_id点击回传时直接精确匹配。3. 实操过程从零搭一条可跑的触达链路3.1 环境与依赖准备这套系统的依赖很轻本地开发用 docker compose 就能全部拉起来PostgreSQL 存任务和决策日志Redis 做配额计数和令牌桶Kafka 做任务队列本地开发可以单节点再加一个 Python 服务跑决策引擎。version: 3.8 services: pg: image: postgres:16-alpine environment: POSTGRES_DB: reach POSTGRES_PASSWORD: reach_dev ports: [5432:5432] volumes: [pgdata:/var/lib/postgresql/data] redis: image: redis:7-alpine command: redis-server --appendonly yes --maxmemory 1gb --maxmemory-policy noeviction ports: [6379:6379] kafka: image: bitnami/kafka:3.7 environment: KAFKA_CFG_NODE_ID: 0 KAFKA_CFG_PROCESS_ROLES: controller,broker KAFKA_CFG_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093 KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: 0kafka:9093 KAFKA_CFG_CONTROLLER_LISTENER_NAMES: CONTROLLER ports: [9092:9092] volumes: pgdata: {}Redis 的maxmemory-policy一定要设成noeviction。这一点我踩过坑默认策略在内存吃满时会把 key 随机淘汰令牌桶的 key 一旦被淘汰系统会认为桶是满的然后放行一整桶流量出去直接把通道打爆。宁可让 Redis 写入报错也不能让它悄悄淘汰计数器。Python 侧我用的依赖是redis、psycopg[binary]、confluent-kafka或者aiokafka看你是同步还是异步技术栈。版本尽量锁死消息中间件客户端的小版本升级经常带来行为变化。3.2 数据模型设计表结构我精简到四张核心表每一张都有明确的用途。-- 触达任务表一条任务代表想给某个用户发某条消息 CREATE TABLE reach_task ( task_id BIGSERIAL PRIMARY KEY, user_id BIGINT NOT NULL, scene TEXT NOT NULL, biz_id TEXT, idem_key TEXT NOT NULL, channel TEXT NOT NULL, priority SMALLINT NOT NULL DEFAULT 30, template_id TEXT, payload JSONB NOT NULL DEFAULT {}, state TEXT NOT NULL DEFAULT pending, expire_at TIMESTAMPTZ, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), updated_at TIMESTAMPTZ NOT NULL DEFAULT now() ); -- 幂等键唯一索引这是防重复的第一道防线 CREATE UNIQUE INDEX uk_task_idem ON reach_task (idem_key); CREATE INDEX idx_task_state_created ON reach_task (state, created_at); -- 决策日志表一次决策可以写多行逐条记录规则判定结果 CREATE TABLE reach_decision_log ( id BIGSERIAL PRIMARY KEY, task_id BIGINT NOT NULL, rule_name TEXT NOT NULL, passed BOOLEAN NOT NULL, detail JSONB NOT NULL DEFAULT {}, decided_at TIMESTAMPTZ NOT NULL DEFAULT now() ); CREATE INDEX idx_decision_task ON reach_decision_log (task_id); -- 通道配置表运行期可以热更新不用重启服务 CREATE TABLE channel_config ( channel TEXT PRIMARY KEY, enabled BOOLEAN NOT NULL DEFAULT true, qps_limit INT NOT NULL, burst INT NOT NULL, daily_cap INT NOT NULL, unit_cost_cent INT NOT NULL DEFAULT 0 ); -- 用户触达偏好表 CREATE TABLE user_reach_pref ( user_id BIGINT PRIMARY KEY, quiet_start TIME, quiet_end TIME, daily_limit SMALLINT NOT NULL DEFAULT 3, unsubscribed BOOLEAN NOT NULL DEFAULT false, updated_at TIMESTAMPTZ NOT NULL DEFAULT now() );有几个细节值得说。idem_key上的唯一索引是硬防线插入冲突时捕获UniqueViolation直接判定为重复任务这个判断放在整个流程的最前面成本最低。channel_config做成表而不是配置文件是为了支持运行期热更新——大促当天通道商临时通知你 QPS 降到 200你得能在不重启服务的情况下改掉。reach_decision_log我加了detail的 JSONB 字段用来存当天已发 2 条剩余配额 1这种上下文。JSONB 的好处是查询时可以用-remain直接提取排查问题特别顺手。3.3 决策引擎核心实现先看令牌桶的 Lua 脚本必须用 Lua 才能保证读取-计算-写回是一个原子操作。-- KEYS[1]: 桶的 key -- ARGV[1]: 速率令牌/秒 -- ARGV[2]: 桶容量 -- ARGV[3]: 当前时间戳秒可带小数 -- ARGV[4]: 本次需要的令牌数 local rate tonumber(ARGV[1]) local cap tonumber(ARGV[2]) local now tonumber(ARGV[3]) local need tonumber(ARGV[4]) local data redis.call(HMGET, KEYS[1], tokens, ts) local tokens tonumber(data[1]) or cap local ts tonumber(data[2]) or now local delta math.max(0, now - ts) tokens math.min(cap, tokens delta * rate) local ttl math.ceil(cap / rate) 60 if tokens need then redis.call(HMSET, KEYS[1], tokens, tokens, ts, now) redis.call(EXPIRE, KEYS[1], ttl) return 0 end tokens tokens - need redis.call(HMSET, KEYS[1], tokens, tokens, ts, now) redis.call(EXPIRE, KEYS[1], ttl) return 1注意tokens初始值用的是cap桶默认是满的这个设计是为了让服务刚启动时能立刻处理一批积压任务而不是匀速慢慢放。如果你希望启动时就严格限速把初始值改成 0 即可。接着是决策引擎的主体逻辑。我把它写成纯函数风格输入候选任务和用户上下文输出决策结果不依赖任何全局状态这样单元测试很好写。import time import redis from dataclasses import dataclass from typing import Optional TOKEN_BUCKET_LUA open(lua/token_bucket.lua).read() PRIORITY {critical: 100, high: 60, normal: 30, low: 10} dataclass class Candidate: user_id: int scene: str channel: str priority: str normal payload: dict None expire_at: Optional[int] None dataclass class Decision: accepted: bool reason: str channel: str delay_seconds: int 0 class DecisionCore: def __init__(self, r: redis.Redis, cfg: dict, pref: dict): self.r r self.cfg cfg # {channel: {qps, burst, ...}} self.pref pref # user_reach_pref 的一行 self._bucket r.register_script(TOKEN_BUCKET_LUA) def decide(self, c: Candidate, now: Optional[float] None) - Decision: now now or time.time() # 1. 任务过期直接丢 if c.expire_at and now c.expire_at: return Decision(False, task_expired) # 2. 用户退订除了 critical 一律不发 if self.pref.get(unsubscribed) and c.priority ! critical: return Decision(False, user_unsubscribed) # 3. 静默期判断 if not self._quiet_passed(now, c.priority): delay self._seconds_to_quiet_end(now) return Decision(False, quiet_hours, delay_secondsdelay) # 4. 通道熔断 if self.r.get(fcircuit:{c.channel}) bopen: return Decision(False, circuit_open) # 5. 通道令牌桶二级限流 ch_cfg self.cfg[c.channel] got self._bucket( keys[fbucket:ch:{c.channel}], args[ch_cfg[qps], ch_cfg[burst], now, 1], ) if got ! 1: return Decision(False, channel_throttled, delay_seconds2) # 6. 用户 24 小时滚动窗口配额 ok, remain self._consume_user_quota(c.user_id, now) if not ok: return Decision(False, fuser_quota_exceeded:{remain}) # 7. 场景级去重同场景 7 天内不重复 scene_key fscene:{c.scene}:{c.user_id} if not self.r.set(scene_key, 1, nxTrue, ex7 * 86400): return Decision(False, scene_dedup) return Decision(True, accepted, channelc.channel) def _consume_user_quota(self, user_id: int, now: float): key fquota:user:{user_id} pipe self.r.pipeline() pipe.zremrangebyscore(key, 0, now - 86400) # 清掉 24 小时前的记录 pipe.zcard(key) _, cnt pipe.execute() limit int(self.pref.get(daily_limit, 3)) if cnt limit: return False, limit - cnt self.r.zadd(key, {f{now}: now}) self.r.expire(key, 90000) return True, limit - cnt - 1 def _quiet_passed(self, now: float, priority: str) - bool: if priority critical: return True t time.localtime(now) hm t.tm_hour * 60 t.tm_min start self._to_minutes(self.pref.get(quiet_start) or 22:00) end self._to_minutes(self.pref.get(quiet_end) or 08:00) if start end: return not (start hm end) return not (hm start or hm end) # 跨零点 staticmethod def _to_minutes(s: str) - int: h, m s.split(:) return int(h) * 60 int(m) def _seconds_to_quiet_end(self, now: float) - int: t time.localtime(now) hm t.tm_hour * 60 t.tm_min end self._to_minutes(self.pref.get(quiet_end) or 08:00) delta (end - hm) % (24 * 60) return max(delta * 60, 60)这段代码里有几个判断顺序的讲究。过期判断放在最前面因为它最便宜且能直接省掉后面所有的 Redis 调用静默期放在熔断和限流之前因为静默期会给出 delay_seconds这些任务是要延后投递而不是丢弃的如果先做了限流判断会白白消耗令牌桶里的令牌等于把配额浪费在了本来就不会发的消息上。第 7 步的场景去重我用的是 Redis 的 SETNX而不是「先查再写」。这是为了避免并发场景下的竞态同一用户的两条相同场景任务同时进入决策先查再写会导致两条都通过。SETNX 是原子的天然免疫这个问题。TTL 设成 7 天键的数量量级大约是日活 × 人均场景数200 万日活按人均 10 个场景算就是 2000 万键Redis 内存完全扛得住。3.4 灰度与上线检查单这套系统上线绝对不能一把梭。我的灰度策略是四个维度逐步放开先按通道灰度站内信优先因为它是自建的出问题影响可控再按场景灰度先跑低优先级的营销场景再按用户比例灰度1% → 5% → 20% → 100%最后才放开高优先级场景。每一档至少观察 24 小时重点看三个指标退订率有没有涨、通道错误率有没有涨、决策拒绝率是否在预期范围内。上线检查单我列了一份每次发版都过一遍检查项通过标准怎么验证幂等键唯一索引存在数据库里能查到 uk_task_idem\d reach_taskRedis 淘汰策略noevictionCONFIG GET maxmemory-policy令牌桶预热冷启动第一批任务不被全部拒绝压测时看首批通过率静默期跨零点22:00-08:00 正确判定为静默单测覆盖 start end 分支降级链可用手动让 LLM 超时消息仍能发出注入延迟故障决策日志落库每条拒绝都有对应日志行抽样查询 task_id熔断开关可手动触发能通过配置立即停止某通道灰度环境验证成本预估误差实际成本与预估偏差 15%对比日账单4. 常见问题与排查技巧实录4.1 用户重复收到同一条消息这是最经典也最招骂的问题。排查顺序我总结成三步先看reach_task里同一idem_key有几行如果有多行说明唯一索引没生效或者被绕过了如果只有一行但用户说收到了两条那就往下游查大概率是通道侧重试导致的。我遇到过的真实原因有三种。第一种是某些通道适配器在超时后自动重试但重试时没有带上 idem_key通道商那边就当新消息处理了。解决办法是强制要求在适配器里透传 idem_key并在通道商后台配置去重窗口一般支持 5 分钟到 24 小时。第二种是业务侧异步补偿任务重复触发了任务创建这种要靠业务侧的幂等来兜。第三种最隐蔽——主从切换期间决策层的 Redis 连接短暂失败代码走了异常分支但没打日志任务被重新投递了一次。这个坑让我明白一件事所有的异常分支都必须有日志和指标静默失败是运维最大的敌人。避坑技巧在通道适配层加一个影子计数。每次实际调用通道前先对idem_key做一次INCR如果返回值大于 1 就告警。这个计数不参与业务逻辑纯粹用于发现重复能帮你提前发现很多隐蔽问题。4.2 队列积压与突发流量积压的典型信号是kafka_consumergroup_lag持续上涨同时任务表里pending状态的行数在堆积。这时候不要急着加消费者先判断是「生产太快」还是「消费太慢」。判断方法很简单看决策日志里各条规则的拒绝率。如果拒绝率明显下降说明是生产端出问题了通常是业务方误配置导致任务量暴增如果拒绝率正常但消费跟不上那就是消费端的问题可能是下游通道变慢了。有一次我们的短信通道商出故障响应时间从 80ms 涨到 3s消费者线程全部堵在等待响应上队列在 20 分钟内积压了 40 万条。当时的处理是先手动打开该通道的熔断开关把任务快速拒绝并标记为可重试让消费者线程释放出来然后把熔断期间被拒的任务批量重投到延迟队列。这个流程我们后来做成了自动化——通道健康检查连续 3 次失败自动熔断30 分钟后半开试探。延迟队列的实现我用的是 Kafka 的重试 topic 配合消费端的时间判断简单粗暴但足够可靠。别用 Redis 的 ZSET 做延迟队列在消息量大的时候定时扫描的开销很可观。4.3 内容校验不通过率异常升高内容校验的失败率正常应该在 2% 以内超过 5% 就要查。最常见的三个原因变量未替换。用户资料缺失导致{{nickname}}没被替换。这类问题要看具体是哪个变量缺失率高如果是昵称就设置默认值比如您好如果是金额这种强业务变量那必须降级到静态模板。长度超限。短信按 70 字符计费超过就分片成本翻倍。我给校验器加了一个长度分布监控如果 P95 长度接近上限说明 LLM 的输出风格变了比如换了个模型版本需要调整 prompt 里的字数约束。敏感词误伤。词表太激进会把正常文案拦下来比如把绝对新鲜这种商品描述也给拦了。解决办法是给词表分级硬禁词直接拒绝软敏感词只告警不拦截人工确认后再调整。4.4 归因数据对不上账业务方说我们后台看到 5000 个转化你这里说 8000差在哪。这种情况九成是口径问题我先讲怎么定位再讲怎么从设计上避免。定位方法拉出这 5000 个转化的明细逐条比对 task_id。如果发现有些转化的 task_id 为空说明这些是自然转化被算进了触达转化需要剔除对照组如果 task_id 存在但不在我们的送达列表里说明用户是通过其他路径进来的比如直接打开 App这种情况要看归因窗口设置。从设计上避免我的做法是三条。第一每条消息带独立的 task_id 落地页点击回传直接精确匹配不依赖时间窗口。第二明确归因窗口我一般设 24 小时超过窗口的行为不算这条消息的功劳。第三固定对照组1% 的 holdout 用户在所有报表里单独列出。这三条做完归因争议基本就消失了。4.5 问题排查速查表把上面这些整理成一张表出问题的时候直接对着看。现象最可能的原因第一步查什么处理动作用户重复收到消息幂等键未透传到通道通道调用的 idem_key 日志补透传 通道侧配去重队列持续积压下游通道变慢各通道 P99 响应时间熔断慢通道 重投延迟队列通道被限流令牌桶参数配错或被绕过bucket:ch:*的剩余令牌核对 qps 配置 检查有无旁路某场景完全没发出去场景去重键误判scene:*键的 TTL 分布清理误写键 修去重逻辑校验失败率飙升变量缺失或模型输出变长失败原因分布补默认值 调 prompt静默期消息漏发跨零点逻辑写错判定分支的单测覆盖修_quiet_passed归因数对不上归因窗口或口径不一致抽样比对 task_id统一口径 加对照组成本超预算短信分片导致成本翻倍长度分布 P95压缩文案 加长度告警提示把这张表做成值班手册的一部分并且给每条现象都绑定一个监控告警。排查速度的提升靠的不是经验丰富的人而是让新手也能按图索骥。5. 效果调优与经验沉淀5.1 频控额度到底该调到多少人均每天 3 条这个数字不是拍脑袋定的但也不是普适的。我的调优方法是做阶梯实验把用户随机分成四组日限额分别是 1、2、3、5 条跑两周同时盯住退订率和转化率两条曲线。实测下来的规律是从 1 条到 2 条转化率提升明显大约 20%从 2 条到 3 条提升收窄到 8% 左右从 3 条到 5 条转化率基本不再增长但退订率涨了 60%。所以 3 条是性价比的拐点。不同业务的最优拐点不一样工具类和电商类的差距可能很大一定要自己跑实验。还有一个容易被忽略的维度是场景优先级。日限额 3 条不等于所有场景平分。我的做法是给场景设权重critical 场景事务通知不占营销配额high 场景比如购物车召回优先占用剩下的额度才分给 normal 和 low。这样能保证重要的消息一定发得出去。5.2 时段与渠道组合的实测结果时段方面我跑过一个月的分时段投放实验把同一批用户按小时切分。结果显示上午 10 点到 11 点、晚上 19 点到 20 点这两个窗口的打开率最高但晚上的退订率也偏高反而是下午 14 点到 16 点这个低谷打开率虽然低 15%但转化质量更好下单客单价高 8%。所以最终我们把高价值场景放在下午高频低值场景放在晚高峰。渠道方面组合策略比单渠道效果好很多但组合的逻辑不是多通道同时发而是分层递进先发站内信成本几乎为零2 小时未读再发推送24 小时仍未转化才考虑短信。这套分层策略让我们的短信用量下降了 40%而整体转化率只降了 3%。省下的成本相当可观。我给渠道选择加了一个简单的打分函数帮你理解决策层的排序逻辑def channel_score(channel: str, user: dict, scene: str, cost_cent: int) - float: base {inbox: 0.6, push: 0.75, im: 0.7, sms: 0.85, email: 0.4}[channel] # 用户历史打开偏好取值 0~1 affinity user.get(channel_affinity, {}).get(channel, 0.5) # 成本惩罚项单位是分除以 100 归一化 cost_penalty cost_cent / 100.0 * 0.3 return base * 0.5 affinity * 0.5 - cost_penalty这个打分函数很粗糙但它把通道基础能力 用户偏好 成本三个因子都考虑进去了。实际生产中我用的是一个离线训练的 LightGBM 排序模型特征包括用户活跃度、历史点击率、时段、场景、通道、成本离线 AUC 大概 0.78。不过说实话早期阶段用上面这个手工打分函数效果能有模型版本的八成先上线验证价值比调模型重要得多。5.3 我踩过的几个坑最后分享几个印象深刻的教训都是文档里不会写的。不要相信低优先级可以慢慢发。我一开始把 low 优先级的任务扔进一个低优先队列想着空闲时再处理结果它永远处理不到——高峰期通道全被高优任务占满低谷期又没有调度触发。后来改成统一队列 优先级加权轮询critical 权重 8high 4normal 2low 1问题才解决。优先级应该是排序权重而不是物理隔离的队列。决策日志要记当时的变量值不是变量的名字。我第一版日志记的是命中规则 user_quota看着挺清楚但排查时没人知道当时用户到底发了几条。改成记录当天已发 2 条剩余配额 1之后所有扯皮都消失了。压测必须用真实的通道或者至少模拟真实延迟。我们用 mock 通道压过 5000 QPS 毫无压力但真实短信通道响应时间 80ms 到 3000ms 不等一压就崩。后来我要求在压测环境接入真实的通道 SDK但把消息发送目标改成一个内部测试号段这样既能压出真实延迟又不会真的发出去。给所有阈值都加上运行期可改的能力。大促当天凌晨两点通道商打电话说要把 QPS 降到 150如果你需要发版才能改那就完蛋了。我把所有关键阈值通道 QPS、桶容量、日限额、静默期、熔断开关全部做成了配置中心的热更新项改动 5 秒生效。留一个全局的一键停发开关。这个开关平时不用但当内容出问题、通道出问题、或者舆情出现的时候能在 3 秒内停掉全部营销类消息只保留事务型消息。我用它救过一次场——某个模板的变量写错了导致一批消息里出现了错误的金额从发现问题到停发只用了 40 秒把影响范围控制在了几千条以内。这套系统从第一版跑到现在大概一年半中间大改过两次架构小调整不计其数。最大的体会是触达系统的核心竞争力不在于发得多而在于知道什么时候不该发。把拒绝的理由记清楚、把配额算明白、把降级链搭扎实比多接几个通道、多上几个模型有价值得多。