连接
美国境内使用 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));
}
});