让 HubSpot 同步真正可靠:事务性 Outbox、QStash 与可修复的一致性

2026-08-20 43 预计阅读时间: 1 分钟
来源: postgr.es AI 摘要 Original link

Disclaimer: This article is an AI-assisted summary. Read it together with the original source when precision matters. The summary may omit context, version differences, or edge cases and is not official documentation.

预计阅读时间:13 分钟

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 限流,以及 502503504521524:指数退避并加入抖动;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 成功但响应丢失、429423207、验证失败、过期 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 和事件账本:当任意一个外部系统不可用时,团队仍能重建状态、定位故障并继续修复。


相关推荐