HubSpot 同步最棘手的故障,不是请求明显失败,而是 HubSpot 已经提交变更,Worker 却在保存结果前丢失响应。此时重试可能重复写入,拒绝重试又会让本地事件永久悬挂。可靠方案必须把数据库、发布器、Worker 和远端变更都设计成可识别、可重试、可修复的状态机。
先接受“至少一次”
队列并不能保证消息只处理一次。发布确认可能丢失,QStash 也可能在 Worker 返回非 2xx 时再次投递,调度器重启后还可能重新发布同一条 Outbox 记录。正确目标不是消灭重复,而是让重复处理安全,并最终收敛到期望状态。
应用先在本地接受用户想要的状态,异步管道再让 HubSpot 逐步达到该状态。这样网站请求不会被 CRM 延迟阻塞,同时本地数据库仍然是可重建、可审计的事实来源。
用事务性 Outbox 消除双写窗口
“写 PostgreSQL,然后发布队列消息”包含一个无法忽略的崩溃窗口:数据库已提交但进程尚未发布,事件就不会被调度。反过来先发消息也有问题,Worker 可能处理一个随后回滚的业务事务。
事务性 Outbox 把业务状态和不可变事件放进同一个数据库事务。独立 Dispatcher 负责发布已经提交的 Outbox 行;即使它发布多次,幂等 Worker 也能安全处理。
CREATE TABLE subscription_requests (
id uuid PRIMARY KEY,
brand_id text NOT NULL,
contact_key text NOT NULL,
product_id text NOT NULL,
desired_state text NOT NULL CHECK (desired_state IN ('subscribed', 'unsubscribed')),
mapping_version integer NOT NULL,
idempotency_key text NOT NULL,
request_hash text NOT NULL,
sync_status text NOT NULL DEFAULT 'pending',
created_at timestamptz NOT NULL DEFAULT now(),
UNIQUE (brand_id, product_id, idempotency_key)
);
CREATE TABLE hubspot_outbox_events (
id uuid PRIMARY KEY,
request_id uuid NOT NULL REFERENCES subscription_requests(id),
event_type text NOT NULL,
payload jsonb NOT NULL,
status text NOT NULL DEFAULT 'pending',
publish_attempts integer NOT NULL DEFAULT 0,
process_attempts integer NOT NULL DEFAULT 0,
next_attempt_at timestamptz NOT NULL DEFAULT now(),
lease_expires_at timestamptz,
qstash_message_id text,
last_error_code text,
last_error_summary text,
created_at timestamptz NOT NULL DEFAULT now(),
completed_at timestamptz,
UNIQUE (request_id, event_type)
);
接收请求时,先规范化输入并计算 request_hash。同一幂等键和同一请求再次到达时返回原始 receipt;同一幂等键却对应不同 hash 时返回冲突。这样客户端即使因为网络超时而重发,也不会创建第二个逻辑订阅。
事件 payload 应保存 mapping snapshot,而不是要求未来的 Worker 重新读取当前配置。普通重试因此具有确定性;需要迁移路由时,创建新的 mapping version 和新的事件。
用租约并发领取事件
多个 Dispatcher 可以使用 FOR UPDATE SKIP LOCKED 领取不同批次,并设置过期租约:
WITH candidates AS (
SELECT id
FROM hubspot_outbox_events
WHERE (status IN ('pending', 'publish_retry') AND next_attempt_at <= now())
OR (status = 'publishing' AND lease_expires_at < now())
ORDER BY next_attempt_at, created_at
FOR UPDATE SKIP LOCKED
LIMIT 100
)
UPDATE hubspot_outbox_events AS event
SET status = 'publishing',
publish_attempts = publish_attempts + 1,
lease_expires_at = now() + interval '2 minutes'
FROM candidates
WHERE event.id = candidates.id
RETURNING event.*;
SKIP LOCKED 适合队列表,因为暂时跳过一条记录比让所有 Dispatcher 排队等待更有吞吐量。它不是 Reconciler 的替代品:过期租约、漏发布记录和丢失调度的重试事件仍需要后续扫描。
发布到 QStash 时只携带事件 ID,必要时加上 mapping version:
const result = await qstash.publishJSON({
url: process.env.HUBSPOT_WORKER_URL!,
body: { eventId: event.id },
retries: 5,
deduplicationId: event.id,
flowControl: {
key: 'hubspot-private-app-primary',
rate: 150,
period: '10s',
parallelism: 12,
},
});
QStash 去重窗口只能作为短期优化,不能代替 PostgreSQL 的持久幂等性。客户邮箱、授权信息和同意证据也不应复制到队列消息、死信工具或日志中。
限制 HubSpot API 调用,而不是消息数量
QStash Flow Control 限制的是 Worker 启动次数,不知道一个 Worker 会发出一次还是十次 HubSpot 请求。如果不同事件的调用计划不同,消息速率仍可能轻易超过 API 配额。
可以在执行计划前按调用数预留共享令牌:
const plan = await buildHubSpotPlan(event);
const estimatedCalls = estimateHubSpotCalls(plan);
const reservation = await hubspotBudget.reserve({
key: 'portal:12345:private-app:primary',
tokens: estimatedCalls,
limit: 150,
windowMs: 10_000,
});
if (!reservation.granted) {
return new Response('HubSpot budget unavailable', {
status: 503,
headers: { 'Retry-After': String(reservation.retryAfterSeconds) },
});
}
共享 Redis token bucket 是一种实现;另一种是把每个操作拆成独立消息,再让 QStash 的流控直接近似调用限流。选择取决于业务是否要求计划原子执行,以及能够接受多少延迟。
Worker 必须验证签名并实现收敛
公开的 Next.js 路由不能因为“理论上只有 QStash 会调用”就被信任。Worker 应在读取和处理事件前验证 QStash 签名,并同时接受当前和下一把签名密钥,以支持密钥轮换。发布 token 只放在服务端环境变量中。
import { verifySignatureAppRouter } from '@upstash/qstash/nextjs';
async function handler(request: Request) {
const { eventId } = await request.json() as { eventId: string };
const event = await claimEvent(eventId);
if (!event || event.status === 'synced') {
return new Response('OK');
}
try {
const plan = await buildPlanFromSnapshot(event);
await applyHubSpotPlan(event, plan);
await markEventSynced(event.id);
return new Response('OK');
} catch (error) {
return handleWorkerError(event, error);
}
}
export const POST = verifySignatureAppRouter(handler);
HubSpot 操作应使用稳定的自定义唯一属性 upsert 联系人;批量 API 支持时,写入稳定的 objectWriteTraceId。通信偏好设置为目标状态,分群操作按实际需要计算 add/remove 差异。“已经处于目标状态”应视为成功。
本地 operation receipt 可以减少重复请求,却无法解决“远端成功、本地保存前崩溃”的最后窗口。要处理这个不确定性,需要远端幂等机制、稳定 upsert 身份,或者在不确定写入后执行 read-after-write reconciliation。
按错误类型决定重试
最大尝试次数不是第一判断条件。建议把错误分成几类:
423锁定、429限流,以及502、503、504、521、524:指数退避并加入抖动;429优先遵守Retry-After。401:暂停任务并刷新或轮换凭证。400、缺失 scope 的403、无效 segment 或 Brand ID:标记为不可重试,进入死信和配置修复流程。207 Multi-Status:逐项检查部分成功结果,不要把整批操作简单地当成成功或失败。
对已知不可重试错误,可以持久化规范化失败信息后返回 QStash 的 489,并设置 Upstash-NonRetryable-Error: true,让消息直接进入死信队列。数据库中的事件状态必须先标记为 dead,保证本地事实不会落后于队列。
同一联系人和产品的顺序通常是业务不变量,但不应把 40 个站点放入一个全局 FIFO 队列。为每个 contact/product stream 保存 sequence,并用数据库版本检查拒绝旧事件:
UPDATE contact_subscription_state
SET applied_sequence = $new_sequence,
desired_state = $desired_state,
updated_at = now()
WHERE contact_key = $contact_key
AND brand_id = $brand_id
AND product_id = $product_id
AND applied_sequence < $new_sequence
RETURNING applied_sequence;
如果没有返回行,说明更高版本已经应用,旧事件可以直接完成而不修改 HubSpot。这样延迟到达的 subscribe 不会覆盖更新的 unsubscribe。
Cron 是修复器,不是主队列
QStash 负责正常投递和重试,定时任务则是独立安全网。一个受保护的 Vercel Cron 路由可以寻找从未发布的 pending 事件、过期租约、丢失调度的 retryable 事件、配置修复后的 dead 事件,以及本地状态和 HubSpot 不一致的记录:
export async function GET(request: Request) {
if (request.headers.get('authorization') !== `Bearer ${process.env.CRON_SECRET}`) {
return new Response('Unauthorized', { status: 401 });
}
const batch = await findReconciliationCandidates({
limit: 250,
pendingOlderThanMinutes: 10,
expiredLeases: true,
retryableBefore: new Date(),
});
const result = await scheduleCandidatesIdempotently(batch);
return Response.json({ checked: batch.length, scheduled: result.length });
}
每次 reconciliation 都应限制行数和执行时间,并使用租约防止重叠运行。成熟的 reconciler 不只重试旧事件,还会批量读取 HubSpot 的权威状态,与本地 desired state 对比,再按最新有效 mapping 生成修正事件。
观测与故障演练
Sentry 标签应能回答 integration、operation、brand、状态和失败类别;上下文可以记录 event ID、request ID、mapping version、尝试次数和 HubSpot correlation ID。不要记录完整事件、邮箱、Authorization header、点击 token 或同意证据。
建议监控 accepted-to-synced 延迟、pending 年龄、重试率、dead 事件数量、QStash 等待队列长度、HubSpot 响应码、剩余限额和 reconciliation 修复数。告警应聚焦 backlog 年龄和永久失败,而不是每一次临时重试。
测试不能只覆盖 HTTP 200。至少应模拟:本地事务提交后 Dispatcher 尚未运行、发布成功后未保存 message ID、两个 Worker 并发处理、HubSpot 成功但响应丢失、429、423、207、验证失败、过期 token、临时 5xx、租约过期回收、旧 subscribe 晚于新 unsubscribe 到达,以及两个 reconciler 同时运行。还要验证修复 mapping 后如何以新 mapping version 重放 dead 事件。
落地检查清单
- 在同一事务中提交 desired state 和 Outbox event。
- 用
SKIP LOCKED与过期租约领取工作。 - 发布持久事件 ID,减少个人数据复制。
- 验证 QStash 签名,并把去重视为短期优化。
- 按 HubSpot 调用数控制下游配额。
- 分离临时、凭证、配置和部分成功错误。
- 只在业务确实要求的 stream 内保证顺序。
- 用受保护的 Cron 做独立 reconciliation。
- 向 Sentry 发送标识符和脱敏失败上下文。
数据库说明应该发生什么,QStash 提供投递,HubSpot 提供 CRM 投影,Sentry 提供可见性。真正让系统可恢复的,是本地 desired state 和事件账本:当任意一个外部系统不可用时,团队仍能重建状态、定位故障并继续修复。