跳到文档正文

构建生产环境消费者

分发已知操作,通过退避策略重连,并确保每个副作用都可安全重复。

运行五分钟快速开始
在聊天中打开
按产品代码校对:2026年9月2日

连接

美国境内使用 wss://ws-iad.tweetstream.io/ws,其他地区使用 wss://ws-global.tweetstream.io/ws。服务端返回 tweetstream.v1,不会回显携带认证 token 的协议。示例通过指数退避重连。短连接会继续增加等待时间,稳定连接 30 秒后重置。

使用 Node.js 连接

typescript
import WebSocket from "ws";
 
type StreamEvent = {
  t?: string;
  op?: string;
  d?: {
    author?: { handle?: string };
    detected?: unknown;
    text?: string;
  };
};
 
const apiKey = process.env.TWEETSTREAM_API_KEY;
if (!apiKey) {
  throw new Error("Missing TWEETSTREAM_API_KEY");
}
 
let retry = 0;
let healthyTimer: ReturnType<typeof setTimeout> | undefined;
let reconnectTimer: ReturnType<typeof setTimeout> | undefined;
 
function scheduleReconnect(reason: string) {
  if (reconnectTimer) return;
 
  const delayMs = Math.min(30_000, 1_000 * 2 ** retry) + Math.floor(Math.random() * 500);
  retry += 1;
  console.warn(`Reconnecting in ${delayMs}ms: ${reason}`);
  reconnectTimer = setTimeout(() => {
    reconnectTimer = undefined;
    connect();
  }, delayMs);
}
 
function connect() {
  const ws = new WebSocket("wss://ws-global.tweetstream.io/ws", [
    "tweetstream.v1",
    `tweetstream.auth.token.${apiKey}`,
  ]);
 
  ws.on("open", () => {
    if (reconnectTimer) clearTimeout(reconnectTimer);
    reconnectTimer = undefined;
    healthyTimer = setTimeout(() => {
      retry = 0;
      healthyTimer = undefined;
    }, 30_000);
    console.log("TweetStream connected");
  });
 
  ws.on("message", (raw) => {
    const event = JSON.parse(raw.toString()) as StreamEvent;
 
    if (event.t === "tweet" && event.op === "content") {
      const tweet = event.d;
      console.log(tweet?.author?.handle, tweet?.text);
    }
 
    if (event.t === "tweet" && event.op === "meta") {
      console.log("enrichment", event.d?.detected);
    }
  });
 
  ws.on("close", (code, reason) => {
    if (healthyTimer) clearTimeout(healthyTimer);
    healthyTimer = undefined;
    scheduleReconnect(`close ${code}: ${reason.toString()}`);
  });
 
  ws.on("unexpected-response", (_request, response) => {
    console.error("Connection rejected", response.statusCode, response.statusMessage);
    response.resume();
    if (response.statusCode === 429 || response.statusCode === 503) {
      scheduleReconnect(`HTTP ${response.statusCode}`);
      return;
    }
    process.exitCode = 1;
  });
 
  ws.on("error", (error) => {
    console.error("WebSocket error", error.message);
  });
}
 
connect();

准备生产环境

  • 连接关闭或发生网络错误后,使用退避策略重连。
  • 只连接文档列出的 TweetStream 端点,并使用上文说明的认证方式。路由由 TweetStream 处理,不需要基础设施专用 Header。
  • 使用 (platform, tweetId) 作为帖子键,而不是消息信封 id。content 从 author.platform 读取 platform,后续操作从 d.platform 读取。
  • 旧版操作缺少 d.platform 时,使用其中的 author platform;如果没有,则按 twitter 处理。
  • 幂等应用每个 tweet/update。更新包含 ref 时,用该快照替换已存引用,并沿 ref.subtweet 读取下一条被引用的帖子。
  • 只有对完整消息信封计算指纹后,才丢弃完全相同的重放。
  • meta 视为某条已路由推文的迟到富化信息。
  • 将回执和投影存入带索引的持久化存储。根据回放窗口和存储预算制定保留或归档策略;不要删除待处理任务。
  • 记录操作失败时不要写入 payload,并通过退避策略重试。使用相同的回执 ID 作为幂等键。
  • 跟踪活跃 WebSocket 连接数,避免超过套餐限制。

测量信号延迟

在消息回调开始时记录本地接收时间。对于 X/Twitter 事件,解码 tweet snowflake 时间戳,测量从发布到接收的时间。主机时钟必须与 NTP 同步,wall-clock 结果才有意义。接收后的处理阶段使用单调时钟。

  • 报告监控列表、客户端区域、UTC 测试时段、样本量、p50 和 p95。
  • 记录重连、缺失事件、时钟同步状态和每条排除规则。
  • 单独报告冷启动和重连后的样本,不要将它们混入热路径结果。
  • 只有发布点、接收点、地理位置、样本和百分位边界一致时,才比较结果。
阶段起点终点
发布到接收X snowflake 时间戳WebSocket 回调开始时记录的本地时间
接收到决策本地 WebSocket 接收策略与风控决策就绪
决策到交易场所确认通过风控的决策单独的交易场所响应或拒绝

测量 Snowflake 延迟

typescript
type ContentEvent = {
  t?: string;
  op?: string;
  d?: {
    tweetId?: string;
  };
};
 
const TWITTER_EPOCH_MS = 1_288_834_974_657n;
 
function tweetIdToTimestampMs(tweetId: string) {
  const id = BigInt(tweetId);
  return Number((id >> 22n) + TWITTER_EPOCH_MS);
}
 
function measureSnowflakeLatency(tweetId: string, arrivedAtMs: number) {
  const tweetedAtMs = tweetIdToTimestampMs(tweetId);
  return arrivedAtMs - tweetedAtMs;
}
 
ws.on("message", (raw) => {
  const arrivedAtMs = Date.now();
  const event = JSON.parse(raw.toString()) as ContentEvent;
  const tweetId = event.d?.tweetId;
 
  if (event.t === "tweet" && event.op === "content" && tweetId) {
    console.log("publication-to-receipt ms", measureSnowflakeLatency(tweetId, arrivedAtMs));
  }
});