十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

智能体触达系统设计:多通道下发、频控与回执

智能体触达系统设计:多通道下发、频控与回执 给智能体装上手这件事我踩过的坑比想象中多得多。模型能写出一封措辞漂亮的召回邮件也能在对话里判断出这个用户该被唤醒了但它写完就停在那儿——没人知道这封邮件该由谁来发、什么时候发、发几次、发失败了算谁的、用户明明点了退订为什么还在收。我后来把这一整层单独抽出来做成一个项目内部叫它 Agent-Reach核心就一句话把智能体的意图翻译成真实世界里可追踪、可限流、可回溯的一条条触达任务。它解决的场景非常具体——多通道消息下发、用户级频次控制、静默期保护、失败重试与回执回流适合做后端、做增长、做智能体工程的同学参考做运营和产品的朋友也能从里面看清一条消息到底走了哪些关卡。1. Agent-Reach 要解决的到底是什么问题1.1 从一个真实的翻车现场说起最早的版本特别朴素Agent 生成文案后直接在业务代码里调用邮件服务的 SDK 发出去。上线第一周就出了三件事。第一件是同一封召回信因为上游 Agent 的重试逻辑和我们自己的重试逻辑叠在一起一个用户收到了七封客服群里直接炸锅。第二件是凌晨两点给一批用户推了短信第二天投诉量翻倍。第三件最要命一个用户明明在三个月前就点了退订但退订状态记在营销系统里而这次调用走的是另一条链路退订名单根本没被查。这三件事指向同一个结论Agent 负责想说什么但它绝对不该负责怎么发。这两件事的失败模式完全不同。前者是模型质量问题改改提示词、换个模型就能缓解后者是工程问题涉及幂等、并发、限流、状态机、回流靠调模型一点用都没有。所以我做的第一件事是把发送能力从 Agent 的调用链里彻底剥离出来中间加一层专门管触达的服务Agent 只提交一个结构化的触达意图剩下的全部由 Agent-Reach 接管。这层抽象带来一个很舒服的副作用Agent 的开发者不需要知道短信服务商的接口长什么样也不需要关心邮件通道的 QPS 上限是多少。他们只需要知道我要发给谁、用哪个模板、传什么变量、优先级多高。反过来通道的接入方也不需要理解 Agent 的推理逻辑。两边通过一份数据契约解耦各自的测试和发布节奏互不干扰。1.2 核心概念拆解意图、任务、通道、回执我把整个模型拆成五个概念这五个概念定清楚了后面所有实现都是顺理成章的。触达意图Intent是 Agent 的输入它描述为什么发。它包含目标用户、场景标识、模板标识、变量字典、期望优先级、以及一个可选的过期时间。意图本身不承诺一定会发出去它只是一个申请。触达任务Task是系统内部真正落地的一条记录是意图经过策略评估、模板渲染、频控校验之后生成的。意图和任务是一对多的关系——一个意图可能因为通道降级被拆成两条任务也可能因为频控被丢弃甚至可能被合并进另一条任务里去。通道Channel是最终的下发手段邮件、短信、站内信、Webhook、企业内部的协作工具机器人每个通道都是一个独立适配器向上暴露统一的发送接口和统一的结果结构。策略Policy是一组可配置的规则决定意图能不能变成任务、什么时候变成任务、以什么通道落地。静默期、用户日上限、通道级配额、退订抑制名单全部属于策略层。回执Receipt是通道返回的、关于这条消息最终命运的事实已送达、被拒、被投诉、被点击。回执是整条链路里唯一能证明事情真的发生了的东西也是策略调优和 Agent 自我评估的唯一依据。把这五个概念摆清楚之后你会发现Agent-Reach 的本质是一个带状态机的调度与治理层它跟 Agent 的智能程度毫无关系。它可以服务一个只会填模板的脚本也可以服务一个每天推理上万次的智能体接口是同一套。1.3 为什么不直接上现成的消息队列加工作流引擎这是被问得最多的一个问题。我们的答案很直接消息队列解决的是任务怎么可靠传递工作流引擎解决的是步骤怎么编排但这两样都没解决这个用户今天还能不能被打扰。我举个具体的例子。假设你用消息队列加消费者来做写一个消费者从队列里拿任务、调用通道、更新状态。看起来没问题但当你需要做用户级频控时问题就出现了消费者是并发的十个消费者同时处理同一个用户的十条任务每个消费者都得去查这个用户今天已经发了几条。查库有延迟查缓存有一致性问题最后你不得不在消费者里加分布式锁。加了锁之后吞吐掉一半而且锁的粒度很难把握——按用户加锁用户量大时锁竞争激烈按分片加锁又会出现跨分片统计不准。更麻烦的是静默期这个需求。它不是简单地晚上不发而是按用户所在时区计算晚上不发并且如果用户跨时区旅行了按他最近一次活跃的时区算。这是一个典型的策略计算问题塞进消费者的业务逻辑里会让代码迅速腐化。我们最后的选择是把任务队列只当传输层真正的决策集中在一个独立的策略评估器里评估器无状态、可水平扩展输入是意图加当前上下文输出是决策结果。这样频控、静默期、降级全部收敛在一处可以单独写测试、单独做灰度。提示如果你现在的项目里频控逻辑散落在三四个服务里各写一遍那你要修的不是某一段代码而是这个分层本身。2. 整体架构与关键选型2.1 三层结构接入层、决策层、执行层最终落地的架构是三层的中间用两段异步边界隔开。接入层负责接收 Agent 的意图提交做鉴权、参数校验、去重键计算然后把意图写入持久化存储并投递一个待评估事件。这一层要足够快因为它可能被 Agent 高频调用所以它不做任何复杂的业务判断只做无脑的落库和投递。决策层是核心它由策略评估器和调度器组成。策略评估器消费待评估事件依次过一遍抑制名单、静默期、频次、通道配额四道闸门输出决策立即执行、延迟到某个时间点执行、降级到备用通道、或者直接丢弃并记录丢弃原因。调度器负责把决策为延迟执行的任务放进时间轮或者延迟队列到点再推入执行队列。执行层就是各个通道适配器加一组工作进程。工作进程从执行队列取任务调用对应的通道适配器拿到结果后更新任务状态并写一条待回执的记录。三段之间的边界为什么用异步因为同步调用会让接入层的可用性绑定到通道的可用性上。短信服务商抖动一下Agent 的整个对话链路就超时了这是不可接受的。异步之后通道挂了只是任务堆积Agent 侧完全无感。2.2 通道适配层怎么设计才不会被新通道拖垮通道适配器是整个系统里最容易失控的部分因为每接一个新通道你都会想这次特殊一点我就在适配器里多写两百行。三次之后适配器之间就没法统一治理了。我给自己定的规矩是适配器只能做三件事——把统一的任务结构翻译成通道需要的请求格式、发起调用并处理传输层的异常、把通道返回的结果翻译成统一的SendResult。所有跟要不要发有关的判断一律不允许出现在适配器里。from abc import ABC, abstractmethod from dataclasses import dataclass dataclass class SendResult: ok: bool channel_msg_id: str | None None error_code: str | None None retryable: bool False raw: dict | None None class Channel(ABC): name: str supports_receipt: bool False abstractmethod async def send(self, task: dict) - SendResult: 只负责把已渲染好的内容投递给通道不做任何策略判断 raise NotImplementedError abstractmethod def capability(self) - dict: 声明通道能力供策略层做降级和配额计算 raise NotImplementedErrorcapability()这个方法看着不起眼但它是降级逻辑的基础。它声明了这个通道支持哪些内容类型、最大长度是多少、是否支持富文本、是否支持回执、单条消息的最大变量数量。当一条富文本站内信因为通道故障需要降级到短信时策略层会先拿短信的 capability 去校验内容发现超长就自动切到短模板而不是硬发一条被截断的乱码。2.3 幂等键的设计决定了你会不会被投诉重复发送是所有触达系统最致命的 bug没有之一。用户能原谅你发得晚但绝对不能原谅你发三遍。幂等的关键在键的设计。我最后用的是这样一组输入意图标识、用户标识、通道名、模板版本号、还有一个时间窗口桶。import hashlib, time def idempotency_key(intent_id: str, user_id: str, channel: str, template_version: str, window_sec: int 3600) - str: bucket int(time.time() // window_sec) raw f{intent_id}|{user_id}|{channel}|{template_version}|{bucket} return hashlib.sha256(raw.encode()).hexdigest()模板版本号必须进键这是血泪教训。有一次我们把模板文案从您的订单已发货改成了您的订单已签收结果因为键没变改版当天的所有通知都被判定为重复而静默丢弃用户当天一条都没收到。窗口桶的粒度也要想清楚。一小时是一个比较稳妥的默认值它允许同一个用户在下一小时收到同类消息但挡住了同一小时内的重复提交。如果你的业务对重复极度敏感可以把窗口拉长到一天但要接受用户当天主动触发的第二次通知会被吞掉这个副作用。我的建议是给不同场景配不同的窗口通知类用一小时营销类用一天验证类干脆不做窗口去重、只做严格幂等同一意图标识只发一次。2.4 技术选型对比表组件我们的选择备选选它的理由什么情况下别选任务存储PostgreSQLMySQL / MongoDB需要事务保证状态机推进的原子性且要按时间范围查回执单表过亿且不便分片时考虑分库队列Redis StreamKafka / RabbitMQ部署轻、延迟低量级没到需要 Kafka 的程度日均任务过千万就要换 Kafka节流计数Redis Lua本地内存 / 数据库脚本原子执行避免查改竞态单机部署且不要求精确时用本地计数更快延迟调度Redis ZSET 时间轮定时任务扫表精度到秒成本低延迟任务量极大时扫表反而更可控状态回写事件总线 幂等消费直接更新数据库回执会乱序到达必须先落事件再推进状态不用回执的场景可以简化配置管理YAML 热加载数据库配置表策略改动频繁走 Git 可评审可回滚需要运营自助改配置时必须上配置表3. 核心模块的实现细节3.1 数据模型三张表撑起整条链路表结构我改过四版最后稳定在意图、任务、回执事件三张主表。意图表记录 Agent 提交的原始申请永久保留它是审计的起点。任务表记录真正下发的东西一条意图可能对应零到多条任务包含最终渲染后的内容快照。回执事件表是一条只增不改的流水用来解决乱序问题。create table reach_intent ( id bigserial primary key, intent_key text not null, -- 业务侧幂等键 scene text not null, user_id text not null, template_id text not null, variables jsonb not null default {}, priority smallint not null default 5, expire_at timestamptz, created_at timestamptz not null default now(), unique (intent_key) ); create table reach_task ( id bigserial primary key, intent_id bigint not null references reach_intent(id), user_id text not null, channel text not null, idem_key text not null, status text not null default PENDING, attempt smallint not null default 0, scheduled_at timestamptz not null, channel_msg_id text, rendered_body text, last_error text, updated_at timestamptz not null default now(), unique (idem_key) ); create index on reach_task (status, scheduled_at); create index on reach_task (user_id, updated_at desc); create table reach_receipt_event ( id bigserial primary key, task_id bigint not null, event_type text not null, event_time timestamptz not null, payload jsonb, unique (task_id, event_type, event_time) );unique (idem_key)这一行是整张表的灵魂。堵重复不靠应用层的先查再插那有竞态靠数据库唯一约束插入失败就是重复捕获异常直接返回已有任务。这个过程不需要事务也不需要分布式锁成本极低。rendered_body存渲染后的最终内容不是模板 ID。有人会问这是不是浪费空间。我的答案是当你三个月后需要回答用户当时到底收到了什么这个问题时模板可能已经改了七八版不复存快照你永远说不清。我们遇到过一起投诉靠这个字段五分钟就自证清白省下的沟通成本远超那点存储。3.2 状态机只允许前进不允许回退任务状态我定义了八个PENDING、SCHEDULED、SENDING、SENT、DELIVERED、FAILED、SUPPRESSED、EXPIRED。核心规则只有一条状态只能沿着一个有向图前进收到任何更早的状态事件都直接忽略。这条规则的价值在处理回执乱序时体现得淋漓尽致。通道的回执可能先到已送达再到已发送因为两条回执走了不同的网络路径如果状态是简单的覆盖写任务会从DELIVERED退回SENT后续的统计全错。加了序号映射之后DELIVERED的序号大于SENT低序号事件直接被丢弃。STATE_ORDER { PENDING: 0, SCHEDULED: 1, SENDING: 2, SENT: 3, DELIVERED: 4, SUPPRESSED: 10, FAILED: 10, EXPIRED: 10, } def advance(current: str, incoming: str) - str: # 终态不再接受任何变更 if STATE_ORDER[current] 10: return current if STATE_ORDER[incoming] STATE_ORDER[current]: return current return incoming失败态和抑制态我给了 10和送达一样都是终态。为什么不给失败态留重试空间因为重试是在SENDING阶段判断的重试预算用尽后才落到FAILED落到终态就意味着这条任务的生命周期结束了。这样设计的好处是任何一条处于终态的任务都可以安全归档不需要担心它哪天又活过来。3.3 频次控制令牌桶加多级配额频控我用的是令牌桶加多级配额四道闸门从粗到细依次是全局通道配额、租户配额、用户日上限、用户通道日上限。全局通道配额必须设因为通道服务商给你的额度是有上限的。超了会怎样轻则限流重则封号。我们的做法是给通道配一个略低于服务商上限的阈值比如服务商给 1000 QPS我们设 800留出余量应对突发的回执查询和重试流量。令牌桶的原子性必须靠 Redis 脚本保证普通的GET再SET会有竞态高并发下会超发。-- KEYS[1] 桶的键 -- ARGV: capacity, rate_per_sec, now_ms, ttl_sec, cost local key KEYS[1] local cap tonumber(ARGV[1]) local rate tonumber(ARGV[2]) local now tonumber(ARGV[3]) local ttl tonumber(ARGV[4]) local cost tonumber(ARGV[5]) local d redis.call(HMGET, key, tokens, ts) local tokens tonumber(d[1]) local ts tonumber(d[2]) if tokens nil then tokens cap ts now end local delta math.max(0, now - ts) / 1000.0 * rate tokens math.min(cap, tokens delta) local allowed 0 if tokens cost then tokens tokens - cost allowed 1 end redis.call(HMSET, key, tokens, tokens, ts, now) redis.call(EXPIRE, key, ttl) return allowed这里有个容易忽略的细节ts要用毫秒并且每次请求都更新否则在长时间没有流量的情况下桶会一直按上一次请求时间到这次请求时间的间隔来补令牌而实际上桶早就该满了。上面这段用min(cap, ...)兜住了但如果你把ts写成只在有令牌消耗时才更新就会出现长时间空闲后突然爆发的问题。用户日上限的计数不适合用令牌桶因为它是固定窗口语义。我用的是按天分桶的计数器键里带上日期和用户 ID过期时间设成 26 小时多两小时防止跨时区的边界问题。这里不要用当天 0 点做过期对齐因为服务器时区和用户时区不一样容易在凌晨出现计数重置的抽风现象。3.4 重试区分可重试和不可重试比退避算法更重要退避算法的代码谁都会写真正容易搞错的是什么错误该重试。我的判断准则很粗暴只有传输层的错误才重试业务层的拒绝一律不重试。具体说连接超时、读写超时、通道返回 5xx、通道返回明确的限流错误码这些重试通道返回参数错误、模板不存在、手机号格式非法、用户已退订一次都不重试直接落终态。import random RETRYABLE_CODES {TIMEOUT, CONN_RESET, RATE_LIMIT, SERVER_ERROR} def should_retry(error_code: str | None, attempt: int, max_attempts: int) - bool: if attempt max_attempts: return False return error_code in RETRYABLE_CODES def next_delay(attempt: int, base: float 3.0, cap: float 300.0, jitter: float 0.3) - float: raw min(cap, base * (2 ** (attempt - 1))) return raw * (1 random.uniform(-jitter, jitter))抖动是必须加的。不加抖动的话一次通道抖动会导致几千条任务在同一毫秒集体重试本来通道只是轻微过载被你的整齐重试直接打穿形成雪崩。我们加的是正负 30% 的均匀抖动实测下来重试流量被摊平得很自然。max_attempts我建议设 4配合上面的参数大约覆盖了 3 秒、6 秒、12 秒、24 秒这几个重试点累计跨度一分钟左右。超过一分钟还没发出去对通知类场景已经没有意义了用户可能已经通过别的渠道知道了。营销类可以放宽到 8 次但一定要给任务设expire_at过了时间窗直接标EXPIRED别再发。注意重试次数的统计必须和attempt字段绑定不能靠日志数数。我见过有人在重试队列里再投一次消息、然后用消费次数当重试次数一旦队列本身重投就全乱了。3.5 静默期与内容渲染两个看起来简单实际上最容易出事的模块静默期看着简单实现起来有一堆边界。我们的规则是用户本地时间 22:30 到次日 08:00 不发营销和普通通知验证类和告警类不受限制。被静默的任务不做丢弃而是重排到静默期结束后的第一个可用时间点并且重排时要做一次抖动在 08:00 到 08:30 之间随机否则所有被静默的任务会在 08:00:00 集体涌出。时区怎么取优先用用户资料里显式设置的时区没有就用最近 30 天最后一次活跃行为的时区再没有就按服务默认时区。这个降级顺序是有讲究的一个用户注册时填的时区可能早就变了而最近活跃时区更能反映他现在在看手机的时间点。旅行场景下这个策略效果明显更好。内容渲染我吃过一次大亏。当时用模板引擎渲染一个变量名写错了渲染引擎的行为是抛异常导致整批两千条任务全部失败。后来改成缺失变量渲染成空串并打一条告警同时在渲染前做一次变量校验模板声明的变量列表和传入的变量字典做差集缺哪个直接记日志。from jinja2 import Environment, StrictUndefined, Undefined TEMPLATE_VARS { order_shipped: {user_name, order_no, express_company}, recall_v2: {user_name, last_active_desc, benefit_amount}, } def validate_vars(template_id: str, variables: dict) - list[str]: need TEMPLATE_VARS.get(template_id) if need is None: return [UNKNOWN_TEMPLATE] return sorted(need - set(variables.keys()))渲染结果的空串问题更要小心。有一次变量缺失渲染成了空串最终短信内容变成尊敬的您的订单已发出被用户截图发到社交平台。所以我在渲染之后加了一道内容检查如果必填变量渲染后为空整条任务直接失败并告警绝不发出去。宁可少发一条也不要发一条会让用户觉得你不专业的消息。4. 从零搭起来的具体过程4.1 环境准备与目录结构本地跑起来需要 Python 3.11 以上、PostgreSQL 14 以上、Redis 7 以上。用 Docker Compose 起依赖别在本地裸装数据库环境问题会消耗掉你一半的调试时间。docker run -d --name reach-pg -p 5432:5432 \ -e POSTGRES_PASSWORDreach -e POSTGRES_DBreach postgres:16 docker run -d --name reach-redis -p 6379:6379 redis:7-alpine python -m venv .venv source .venv/bin/activate pip install fastapi uvicorn asyncpg redis jinja2 pydantic-settings目录结构我按职责切不按技术层切。每个通道一个文件策略、调度、幂等各自独立回执单独一个包。agent-reach/ ├── reach/ │ ├── api/ # 意图提交、任务查询接口 │ ├── core/ │ │ ├── policy.py # 四道闸门 │ │ ├── scheduler.py # 延迟调度与时间轮 │ │ ├── idempotency.py # 幂等键与去重 │ │ └── state.py # 状态机推进 │ ├── channels/ │ │ ├── base.py │ │ ├── email.py │ │ ├── sms.py │ │ └── webhook.py │ ├── worker/ # 消费与执行 │ └── receipts/ # 回执接收与状态回流 ├── config/policies.yaml ├── migrations/ └── tests/4.2 策略配置文件长什么样策略我放在 YAML 里进程启动时加载支持发信号热加载。放到 Git 里管每次改动都有评审记录出问题能立刻回滚到上一版。version: 3 policies: default: quiet_hours: [22:30, 08:00] quiet_exempt_scenes: [verify_code, alert] user_daily_cap: 3 user_channel_daily_cap: sms: 2 email: 5 retry: max_attempts: 4 base_seconds: 3 cap_seconds: 300 jitter_ratio: 0.3 task_ttl_seconds: 86400 channels: sms: burst: 200 rate_per_sec: 200 daily_quota: 500000 fallback: email email: burst: 500 rate_per_sec: 300 daily_quota: 2000000 webhook: burst: 100 rate_per_sec: 100 daily_quota: 1000000fallback字段是降级链。短信失败或者配额用尽时自动落到邮件但要经过 capability 校验超长内容要走短模板。这里有个细节降级不能无限链式往下传我限制最多降一级否则一条消息可能在三个通道之间来回跳排查问题时根本说不清它到底经历过什么。4.3 工作进程的核心循环工作进程的逻辑其实就是取任务、判断、执行、写回四步但每一步都要小心。async def consume_once(ctx) - None: task await pop_task(ctx) # 从执行队列取出 if task is None: return if task[expire_at] and task[expire_at] now(): await mark(ctx, task[id], EXPIRED, past_expire) return if not await take_token(ctx, task[channel], cost1): # 通道没额度放回去稍后再来注意不要立刻重投 await requeue_later(ctx, task, delaynext_delay(task[attempt] 1)) return await mark(ctx, task[id], SENDING) channel ctx.channels[task[channel]] result await channel.send(task) if result.ok: await mark(ctx, task[id], SENT, channel_msg_idresult.channel_msg_id) elif result.retryable and should_retry(result.error_code, task[attempt] 1, ctx.policy.retry.max_attempts): await requeue_later(ctx, bump_attempt(task), delaynext_delay(task[attempt] 1)) else: await mark(ctx, task[id], FAILED, result.error_code)几个容易被忽略的点。第一取令牌失败时不要立刻重投要走退避否则通道打满时你的进程会疯狂空转把 Redis 也打满。第二mark到SENDING必须在调用通道之前这样进程万一崩了这条任务会留在SENDING状态被一个兜底扫描任务发现超过 5 分钟没变化就重置为PENDING并加一次 attempt不会永久卡死。第三SENT不等于DELIVERED别混用前者只代表通道收下了后者要等回执。4.4 回执接口与状态回流回执接口是最容易被低估的部分。它要处理乱序、重复、以及来源不可信三个问题。async def handle_receipt(ctx, payload: dict) - None: if not verify_signature(payload, ctx.receipt_secret): return # 静默丢弃不要返回错误细节 event_key (payload[task_id], payload[event_type], payload[event_time]) inserted await insert_event_if_absent(ctx, event_key, payload) if not inserted: return # 重复回执直接忽略 row await load_task(ctx, payload[task_id]) target map_event_to_state(payload[event_type]) new_state advance(row[status], target) if new_state ! row[status]: await update_state(ctx, row[id], new_state)去重靠数据库唯一约束和幂等键是同一套思路。签名校验失败时不要返回签名错误这种明确信息直接返回 200 静默丢弃避免给外部探测者提供反馈。回执写入只增不改即使状态推进失败事件流水也完整保留了事后可以用它重算任何一条任务的真实历史这个特性在排查疑难问题时救过我好几次。4.5 灰度与压测怎么做才有效上线不要一次全量。我的做法是按场景灰度先接验证码这类容忍度最高、链路最短的场景跑一周看指标再接通知类最后才接营销类因为营销类触发量大、用户敏感度高翻车代价最大。压测必须打出真实的并发形态。我见过太多人压测时用均匀的 100 QPS 打十分钟一切正常上线后真实流量是每十分钟来一个尖峰尖峰持续八秒量级五万条直接把通道打穿。所以压测脚本要能造尖峰而且要按通道服务商给你声明的限流阈值来造——你压出的上限绝对不能高于服务商的限流值否则压测本身就是一次事故。上线前必看的四个指标任务成功率应该 99% 以上、降级率超过 5% 要查通道、抑制率突然飙升通常是误伤、回执到达延迟的 P99超过十分钟说明回执链路有问题。这四个指标我专门做了个看板每次发版后盯半小时。5. 踩坑实录与问题速查5.1 那些只有真跑起来才会遇到的怪问题任务永远卡在SENDING。根因是进程在调用通道时被 kill状态没回写。解决办法是加一个兜底扫描每两分钟扫一次超过五分钟未变化的SENDING任务重置成PENDING并把 attempt 加一。注意重置时不要清空 attempt否则一个必然会崩的任务会无限循环。同一个用户短时间内收到两条内容几乎一样的消息。不是因为幂等失效而是因为 Agent 提交了两个不同的intent_key语义上是同一件事。纯粹靠幂等键解决不了这个问题必须靠内容指纹去重把渲染后的正文做一次归一化去掉空白、变量值替换成占位符再算哈希一小时内相同指纹只放一条。这个逻辑我加在渲染之后、入库之前。回执里的task_id找不到对应任务。两种可能一是回执里带的是通道自己的消息 ID 而非我们的 task_id需要建一张映射表二是任务被归档清理了。我们最早为了省空间把 30 天前的任务删掉结果 45 天后回来的延迟回执全部成了孤儿。后来改成任务表不删只把rendered_body和payload转冷存储。静默期误伤了时区识别的边界用户。有个用户在时区切换的当天被连续静默了两天。根因是我们在静默期判断时用了用户本地时间而时区来源在两天内变了两次。解决办法是把时区选择的结果缓存两小时并且对静默重排的任务加一个最多顺延一次的限制——如果顺延后再次落入静默期说明时区判断有抖动就直接放行并打告警宁可发得早一点也不要无限制地顺延下去。压测时一切正常真实流量下成功率掉到 85%。最后定位到是连接池。通道适配器用的是异步 HTTP 客户端默认连接池上限很低压测机上的并发远低于真实生产环境的并发所以没暴露。把连接池上限调到和通道配额匹配之后就恢复了。这个教训是压测环境要尽量逼近生产的实例数和连接配置。5.2 常见问题速查表现象最可能的根因排查动作处理方式用户收到重复消息幂等键输入缺了模板版本或窗口桶对比两条任务的 idem_key补齐键输入历史任务手动标记消息发不出去且无错误任务被抑制状态为 SUPPRESSED查抑制原因字段和当日计数确认是否误伤必要时清理计数通道整体限流报错令牌桶容量配置高于服务商上限对比配置与通道文档下调 rate 与 burst留 20% 余量回执大量丢失回执地址未加白名单或签名校验过严看接入网关日志修正校验逻辑补对账任务状态机出现回退回执乱序且用了覆盖写查回执事件表时间序引入状态序号比较次日 0 点后频控计数不重置计数器过期时间与用户时区不对齐查 Redis 键的 TTL改为 26 小时过期 按日期分桶降级后内容被截断未做 capability 校验对比降级前后内容长度加降级前的长度与类型校验5.3 几条不打折扣的经验第一任何跟用户感知相关的逻辑都要有这条消息为什么被发出/被拦下的完整解释。我们的任务表里有一个decision_trace字段记录四道闸门各自的判断结果。看起来是冗余但每次有人问为什么这个用户没收到我都能在十秒内给出答案省下的时间远超维护成本。第二日志脱敏要在写入日志的那一刻做不能靠事后过滤。手机号、邮箱、姓名这些字段在渲染之后就会被拼进正文一旦进了日志系统就再也收不回来了。我在日志中间件里做了统一的正则替换并且规定rendered_body只允许在调试模式下打印生产环境打印哈希。第三退订和抑制名单必须是强一致的不能走缓存。我们最早为了性能把抑制名单放在本地缓存里五分钟刷新一次结果一个刚退订的用户在两分钟内又收到了一条营销消息直接引发投诉。后来改成每次发送前实时查一次 Redis 的抑制集合命中就直接抑制这个查询的成本完全值得。6. 上线之后的影响范围这个项目上线后变化最大的其实不是技术指标而是团队协作方式。研发侧最明显的收益是 Agent 的开发不再需要关心发送细节。以前接一个新通道要改三四个地方现在只要实现一个适配器接口注册进配置就完事。同时因为策略层独立运营想调整频次上限不需要发版改 YAML 走一次评审就生效响应速度从排期到下周变成当天就能上。运营侧的收益是终于能回答为什么没发出去。以前被问到这个问题只能翻代码猜现在有一个明确的抑制原因字段是静默期、是超频次、还是退订名单一目了然。这直接改变了运营的优化方式——从多试几次看看变成看着原因去修。用户侧的收益更隐性但更重要。静默期和频次上限这两个约束本质上是把尊重用户从一句口号变成了系统强制执行的规则。任何一条路径上想绕过它都必须改代码并且过测试而不是靠某个人的自觉。需要清醒的是这套东西带来的成本。多了一层服务就多了一层运维对象、多了一套要监控的指标、多了一份要维护的配置。如果你的日均触达量只有几百条接入一个成熟的消息平台可能比自建更划算自建的意义在量上来之后、在策略复杂度上来之后、在你需要把触达效果回流给智能体做自我评估之后才会显现。我个人在实际操作中的体会是这类项目最难的部分从来不是写代码而是把什么情况该发、什么情况不该发这些原本模糊的、散落在各个人脑子里的规则一条条整理成可以执行的判断。这个过程很枯燥会反复和产品、运营、法务来回确认但整理完之后你会发现它才是这套系统真正的价值所在——代码只是把这些规则固化了下来。后续如果继续往下做我会把回执数据接回智能体的评估链路让它知道自己上一次触达到底有没有效果这样整个闭环才算真正合上。
返回列表