亚马逊数据 API Node.js 的教程大多停在「await fetch 一行拿到 JSON」。从这一刻到「每天按点交付干净数据」,中间隔着三个 TypeScript 编译器帮不上忙的地方:网络指纹、并发模型、幂等重跑。网络指纹是 Node 没有 Python 那样的默认答案;并发模型会在开启 HTTP/2 之后整段改变,让常见的并发上限写法失去意义;幂等重跑决定一次失败之后是重跑一遍还是重复落库。本文按这个顺序走完,每段给可运行的 TypeScript,并说清难点在哪、什么时候该买现成的。
一、为什么 Node 团队更容易在这件事上翻车
把一篇标准 Node 教程里的脚本放进 cron,前三天的日报通常很好看。问题集中在第三到第五周浮出水面,而且都不以报错的形式出现。
1.1 三个不报错的失效方式
| 现象 | 账面上的表现 | 缺口在哪一层 |
|---|---|---|
| 返回的是拦截页 | res.status 仍是 200,JSON 里没有你要的字段 | 缺一层「这份响应是不是真数据」的判断 |
| 字段静默变空 | 类型通过编译,运行时字段是 null | 缺一层「这个字段必须有值」的运行时断言 |
| 任务中断后重跑 | 行数增长,业务侧开始投诉同一天两条价格 | 缺主键设计与幂等写入 |
三项的共同点是监控不会叫。HTTP 状态码、异常计数、进程退出码全都是正常的,只有报表里的数字在偏离。指望告警发现它们,等于让监控去猜业务语义。
1.2 编译器通过,不代表数据对
TypeScript 团队特别容易在这里产生错觉。定义好 interface,编译器把字段名、是否可选、嵌套结构都检查了一遍,看起来数据已经校验过了。实际情况是:类型标注只约束你自己的代码,不约束网络上送来的那一包 JSON。上游把 price 从数字改成字符串,或者干脆不再返回这个 key,编译器一个字都不会说,因为它在编译时看不到运行时的数据。
网络那一侧的缺口同理,而且 Node 更麻烦。改动不是上游发起的,是你选的 HTTP 客户端决定的,很多团队直到被拦才意识到有这一层。
1.3 还有一笔返工成本
本地脚本错了,改一行再跑一次,代价是三十秒。生产任务里混了三周的重复记录,代价是重新追回所有下游报表,而且你未必知道哪些报表用过这批数据。结构上省下的几个小时,会在数据事故里连本带利还回去。
下面的划分依据就是这个边界:前三节解决「请求能不能被当成正常流量」,中间五节解决「拿到的数据能不能每天对齐」,最后三节解决「跑起来的任务能不能长期不惹事」。
二、Node.js 侧有哪些 HTTP 客户端可选
选型不是挑 API 好不好用,是挑「还要自己写多少代码」。先看一张对照表,再解释为什么这张表里有些行是陷阱。
| 客户端 | HTTP/2 | TLS 指纹 | 2026 年的状态 |
|---|---|---|---|
fetch(Node 内置) | Node 20 默认不走 h2,需显式 agent | Node 自己的 OpenSSL 握手 | 零依赖,但两层都要自己配 |
undici | 需 allowH2: true | 同上 | Node 官方维护,连接池可控 |
axios | 不支持 | 同上 | 生态最大,多路复用缺失 |
got | 支持 | 同上 | 用途通用,指纹层没有覆盖 |
got-scraping | — | 只改写请求头 | 已停止维护 |
impers | 支持 | curl-impersonate 内核 | Node 绑定,接近 curl_cffi 的用法 |
wreq-js | 支持 | Rust + BoringSSL,进程内 | 吞吐最高,需原生二进制 |
各家方案在字段层面的对比方法,见亚马逊数据 API 选型对比。这里有一个容易被忽略的历史包袱:got-scraping 曾经是 Node 圈做采集的默认答案,它已经停止维护。更重要的是,它从设计上只改写请求头,从不碰 TLS 握手。也就是说,凡是卡在 JA3 或 JA4 上的拦截,它本来就解决不了——这个缺口不是版本问题,是能力边界问题。如果你的方案建立在这类库上,换版本不会有变化。
2.1 一张实测表里的取舍
指纹类库之间差别不小。下面这组数据取自 wreq-js 仓库公开的基准测试,测量日期 2026-08-06,环境是 M 系列 Mac,300 次串行请求对一个本地服务,测的是跨 JS 与原生边界的开销而非网络延迟。数值会随版本变化,引用前请按当时版本自行复测。
| 库 | 引擎 | 可模拟的最新 Chrome | HTTP/2 指纹是否正确 | req/s | 冷启动 |
|---|---|---|---|---|---|
wreq-js | Rust wreq + BoringSSL,进程内 | 149 | 正确 | 12842 | 7 ms |
impers | curl-impersonate,进程内 | 146 | 正确 | 8439 | 16 ms |
node-wreq | 同上 Rust 内核 | 149 | 正确 | 6500 | 10 ms |
impit | Rust reqwest + 打补丁的 rustls | 124 | 不正确 | 6710 | 37 ms |
CycleTLS | Go 子进程 + IPC | 未参与 | 未参与 | 未参与 | 每次请求额外 IPC 开销 |
这张表里要看的不是吞吐,而是「HTTP/2 指纹是否正确」这一列。
2.2 「声称 Chrome,说着 Rust」
一个客户端可以在 TLS 层伪装成 Chrome,同时在 HTTP/2 层露馅。impit 那一格的成因是:它的 HTTP/2 SETTINGS 用的是底层 Rust HTTP 库的默认值,而不是 Chrome 会送的那组值——缺少 HEADER_TABLE_SIZE,并且发送了一个 Chrome 从不发送的 MAX_FRAME_SIZE。结果是:任何一个对 SETTINGS 帧做哈希的风控,都会看到一个自称 Chrome、却在用 Rust 的语法说话的客户端。
这类不一致比纯粹的落后更可疑。风控系统抓的是矛盾,不是版本旧。同理,如果 UA 说 Windows 而 HTTP/2 的窗口大小像 Linux 客户端,或者 Accept-Language 写着 de-DE 而出口 IP 在美国,这些都是同样的错——信号之间彼此否认。
这也是 Node 比 Python 难的地方。Python 生态里 curl_cffi 基本是唯一答案,伪装 OpenSSL 与浏览器 profile 是一套东西一起解决的;Node 生态里没有一个同等级的默认选择。亚马逊数据 API Python 那一篇里,curl_cffi 基本是唯一答案;Node 侧上述几个库的实现路径各不相同(Rust 原生、Go 子进程、curl 绑定),各自的覆盖程度也不一样,需要自己验证而不是默认信任。
2.3 会话与非会话的成本差
还有一个和语言无关、但 Node 团队常踩的点:单次调用 fetch 每次都开新连接,等于每次都要完成一次完整 TLS 握手;用会话(或连接池)的话,握手可以被复用。公开基准里给出的对比量级是每次约 53 ms 对 15 ms。批量任务如果每次请求都新建连接,这一段握手成本会直接换算成任务时长。
// 不推荐:每次调用都新建连接
for (const asin of asins) {
await fetch(buildUrl(asin)); // 每次完整 TLS 握手
}
// 推荐:会话/连接池复用 TLS 会话与 cookie
const session = await createSession({ browser: "chrome_149" });
try {
for (const asin of asins) {
await session.fetch(buildUrl(asin));
}
} finally {
await session.close();
}
三、TLS 与 HTTP/2 指纹:怎么验证你改对了
上一节说的不一致,在 Node 里只能靠自己验证,没有编译器或类型帮你兜底。验证分两步。
3.1 第一步:打到公开指纹端点
不要用「能返回 200」当作通过标准。拦截页也返回 200。去公开指纹检测端点(例如 TLS 指纹查询类服务)打一次,把返回的 JA3 / JA4 与你声称要模拟的浏览器版本对比。
const res = await session.fetch("https://tls.peet.ws/api/all", {
headers: { "accept-language": "en-US,en;q=0.9" },
});
const fp = await res.json();
console.log(fp.ja4, fp.akamai_h2); // 与真实浏览器的值比对
这里顺带解释 JA3 与 JA4 的关系:JA3 是把 ClientHello 五个字段按顺序拼起来取 MD5。Chrome 110 之后扩展顺序会做随机置换,JA3 跟着变,所以主流风控早就换成了先排序再哈希的 JA4。这也解释了「我明明固定了 JA3 还是被拦」——你修的是一个已经不被主要使用的特征。
3.2 第二步:检查五件套是否同源
TLS 指纹对了还不够,还要看这几个信号是否来自同一个 profile:
| 信号 | 常见错配 |
|---|---|
| User-Agent 与 platform | UA 写 Windows,HTTP/2 窗口像 Linux 客户端 |
| Accept-Language 与站点 | amazon.de 却发 en-US |
| 出口国家与站点 | 德国站点从美国出口 |
| 时区与 IP 归属 | UTC 时区配欧洲 IP |
| TLS 与 HTTP/2 两套指纹 | TLS 像 Chrome 149,HTTP/2 SETTINGS 像 Rust 默认值 |
提醒一句:原生二进制的依赖在容器里最容易出事。这类库通常需要预编译二进制或 Rust 构建工具,Alpine(musl)与常见的 Debian(glibc)镜像要分别验证。上线前先跑一次冷启动,确认容器里能加载,而不要等定时任务第一次触发才发现。
四、开启 HTTP/2 之后,并发模型整段改变了
这一节是 Node 侧最容易被跳过、也最容易把人绕进去的部分。结论先说:你在大多数教程里看到的并发上限写法,是在 HTTP/1.1 的假设下成立的;一旦连接协商成 HTTP/2,那个上限的含义就变了。
4.1 Node 内置 fetch 默认不协商 HTTP/2
Node 的 fetch 建立在 undici 之上,但它默认不在 ALPN 协商里选 h2。官方 issue 里的答复是 HTTP/2 支持属实验性质且未默认开启,需要显式传 agent:
import { Agent, setGlobalDispatcher, fetch } from "undici";
setGlobalDispatcher(new Agent({ allowH2: true })); // 显式开启才会走 HTTP/2
const res = await fetch("https://api.example.com/v1/product");
而如果你直接用 undici 的 Client,它的 allowH2 默认值是 true。同一个 undici,两种入口的默认值不同,这本身就是一类事故来源——同一份业务代码,换一层调用方式,连接行为就不一样了。
4.2 并发天花板从 pipelining 换成了 maxConcurrentStreams
关键在这一句:一旦协商成 HTTP/2,限制单条连接上同时在飞请求数的,不再是 pipelining,而是 maxConcurrentStreams,默认值 100。
含义是:HTTP/1.1 时代一条连接上一次只能处理一个请求(严格说是 pipelining 默认关闭),所以「并发数」几乎等价于「连接数」;换成 HTTP/2 之后,一条连接上可以跑上百条流,此时的瓶颈变成了三处:
| 层级 | HTTP/1.1 时的含义 | 协商成 HTTP/2 后的含义 |
|---|---|---|
应用层并发(p-limit 等) | ≈ 同时在飞的请求数 | 仍在限制你的代码一次发起多少请求 |
流数量 maxConcurrentStreams | 不存在 | 单连接同时流的硬上限,可被服务端 SETTINGS 覆盖 |
流控窗口 initialWindowSize | 不存在 | 默认值 262144,窗口耗尽时吞吐会被卡住 |
所以「把并发从 8 提到 64」在两种协议下的效果并不相同。HTTP/1.1 下这会直接放大连接压力;HTTP/2 下如果 maxConcurrentStreams 或窗口先到顶,提高应用层并发只是让排队时间变长,吞吐不涨。
4.3 正确的两层夹心
因此并发控制要写两层:应用层的队列上限(决定你的代码同时放多少请求出去)+ 传输层的连接池与 HTTP/2 参数(决定这些请求怎么落到连接上)。
import { Agent } from "undici";
import pLimit from "p-limit";
// 传输层:连接池 + HTTP/2 参数
const dispatcher = new Agent({
allowH2: true,
connections: 8, // 到单个 origin 的连接数
pipelining: 0,
maxConcurrentStreams: 100, // HTTP/2 的默认上限,可按服务端 SETTINGS 调整
bodyTimeout: 30_000,
headersTimeout: 15_000,
connect: { timeout: 5_000 },
});
// 应用层:同时在飞的请求数
const limit = pLimit(16);
async function run(jobs: Job[]): Promise<Result[]> {
return Promise.all(jobs.map((job) => limit(() => fetchOne(dispatcher, job))));
}
调参顺序:固定应用层并发不动,先调 connections;观察 429 占比与端到端时延往上加;加到某个点吞吐不再涨而排队时间开始变长,就退回上一档。天花板多数时候在上游给你的额度,而不是你的机器核数。
五、请求层封装:超时、凭据、重试收进一个类
第一件工程化的事,是把散在各处的参数收进一个对象,让「一次调用」有唯一定义。三条硬规则:超时按阶段分别给;凭据来自环境变量;重试在类内完成并把次数返回。
5.1 为什么要按阶段给超时
Node 里最常见的写法是 AbortSignal.timeout(30_000) 一个值包打天下。问题是这个值同时承担了三件不同的事:建连要多久、响应头要多久返回、响应体要多久流式读完。一个值意味着最宽松的那个决定了全部,于是慢连接会把 worker 一直占住。
import { Agent, request } from "undici";
type Marketplace = "US" | "DE" | "JP" | "UK";
interface FetchOutcome<T> {
ok: boolean;
status: number;
body: T | null;
attempts: number;
error?: string;
}
export class AmazonClient {
private readonly dispatcher: Agent;
constructor(private readonly baseUrl: string, opts: Partial<ClientOpts> = {}) {
this.dispatcher = new Agent({
allowH2: true,
connections: opts.connections ?? 8,
connect: { timeout: 5_000 }, // 建连
headersTimeout: 15_000, // 等响应头
bodyTimeout: 30_000, // 等响应体
});
}
}
三个值分别卡三个阶段,最容易被长尾卡住的阶段有自己的上限。返回结构里的 attempts 必须带上——它是成本公式里的变量,也是判断上游是否在退化的信号。
5.2 重试的边界
重试只给两类失败:瞬时错误(超时、连接重置、502/504)和限流(429 或带 Retry-After 的 503)。第三类该重试的对象最容易被搞错:200 但字段缺失不属于失败,它属于契约错误,重试一百次只会得到同样的空字段,同时把成本乘一百。
const LANG: Record<Marketplace, string> = {
US: "en-US,en;q=0.9",
DE: "de-DE,de;q=0.9",
JP: "ja-JP,ja;q=0.9",
UK: "en-GB,en;q=0.8",
};
async function withRetry<T>(
fn: () => Promise<T>,
opts: { maxAttempts?: number; baseDelayMs?: number } = {},
): Promise<{ value: T; attempts: number }> {
const maxAttempts = opts.maxAttempts ?? 3;
const base = opts.baseDelayMs ?? 1000;
let lastError: unknown;
for (let attempt = 1; attempt <= maxAttempts; attempt++) {
try {
return { value: await fn(), attempts: attempt };
} catch (err) {
lastError = err;
if (!isTransient(err) || attempt === maxAttempts) throw err;
// 指数退避 + 抖动,避免一批失败请求在同一秒集体重发
const wait = base * 2 ** (attempt - 1) + Math.random() * 400;
await new Promise((r) => setTimeout(r, Math.min(wait, 30_000)));
}
}
throw lastError;
}
抖动那一段不要省。不带抖动的指数退避会让同一批失败请求在下一次重试时同时发出,正好制造一个新的突发。
5.3 Node 特有的错误名字
Node 的错误模型和其他语言不一样,判定 retryable 时要认名字而不是只看状态码:
| 错误来源 | 典型表现 | 处理 |
|---|---|---|
UndiciError:UND_ERR_CONNECT_TIMEOUT | 连接阶段超时 | 可重试 |
UndiciError:UND_ERR_HEADERS_TIMEOUT | 响应头迟迟不来 | 可重试 |
UndiciError:UND_ERR_BODY_TIMEOUT | 响应体流式读取超时 | 可重试,但要记录 body 大小 |
DOMException:TimeoutError | AbortSignal.timeout 触发 | 可重试 |
ENOTFOUND / EAI_AGAIN | DNS 解析失败 | 可重试,连续出现要查出口网络 |
| HTTP 403 / 429 | 拦截或限流 | 429 按 Retry-After 退避;403 换档位并降载 |
| JSON 解析失败 | 响应体不是 JSON | 多为拦截页,不能按瞬时错误重试 |
最后一行值得单独说。解析失败常常是「返回了 HTML 拦截页」而非「网络抖动」,把它归到瞬时错误里重试,等于用一百分力气反复撞同一道门。
5.4 凭据放哪
三条线:本地开发用 .env 但不要进版本库;容器环境由编排平台的密钥管理注入环境变量,不要写进镜像层——镜像层能被任何人 docker history 翻出来;轮换要有明确的生效路径。启动时读一次意味着换钥匙要重启,批量任务能接受这个前提就这么做,不能接受的话改成按需读取配短缓存。
还有一条容易被忽略的:不要把 token 拼进 URL 查询串。查询串会完整留在你的访问日志、对方的访问日志以及中间任何一层代理的记录里,而放在 Authorization 头里的同一个 token 不会。看到 ?api_key= 这类拼法,改掉它。
六、TypeScript 类型保不住的那部分:Zod 运行时契约
这一节回应标题里的第一个失败。先把边界划清楚。
6.1 类型标注约束的是你的代码,不是网络
写下 interface ProductRow { price?: number },编译器保证的是「你自己写的代码里不会把 price 当字符串用」。它不知道:
- 上游这次会不会返回
price; - 返回的是
19.99还是"19.99"; currency从"USD"变成null时意味着不该入库。
校验必须发生在运行时,在数据进入你的存储之前。用 Zod 做这件事的好处是 schema 与类型同源——一份 schema 可以同时给出静态类型和运行时校验器,不需要维护两份定义。
import { z } from "zod";
export const ProductSnapshot = z.object({
asin: z.string().regex(/^[A-Z0-9]{10}$/),
marketplace: z.enum(["US", "DE", "JP", "UK"]),
capturedAt: z.string().regex(/^\d{4}-\d{2}-\d{2}$/), // 采集日期,进主键
contractVersion: z.string().default("2026-09-01"),
title: z.string().min(1),
brand: z.string().nullish(),
price: z.number().positive().nullish(), // Buy Box 价格
currency: z.string().length(3).nullish(),
rating: z.number().min(0).max(5).nullish(),
reviewCount: z.number().int().nonnegative().nullish(),
bsrMain: z.number().int().positive().nullish(),
bsrCategory: z.string().nullish(),
});
export type ProductSnapshotT = z.infer<typeof ProductSnapshot>;
三条纪律:price 这类可空字段显式用 nullish() 标注,表示「允许缺失」而不是「一定能拿到」;capturedAt 用采集日期而不是写入时间戳,同一天重跑必须落到同一天;字段语义一旦改动就升 contractVersion,让新旧数据可以共存而不是互相覆盖。
6.2 P0 字段硬断言
不是所有字段都同等重要。先挑出「缺了这条记录就不能用」的那几个,校验不到直接拦下来:
const P0_FIELDS = ["title", "price", "currency"] as const;
function assertP0(row: unknown): row is ProductSnapshotT {
const parsed = ProductSnapshot.safeParse(row);
if (!parsed.success) return false;
return P0_FIELDS.every((f) => parsed.data[f] !== null && parsed.data[f] !== undefined);
}
注意用了 safeParse 而不是 parse。后者会抛异常,在一个 Promise.allSettled 的循环里,一次畸形响应会把整批任务的错误处理路径炸成一个异常,反而更难定位。
6.3 为什么要同时留着 contractVersion
数据事故里最难查的一类不是「采不到」,是「采到了但含义变了」——今天的 price 是 Buy Box 价格,上个月是列表价,上游不会提前通知。把版本号写进主键之后,两个口径的数据可以同时存在,对齐时按版本取,历史不会被悄悄改写。
七、质量门:覆盖率与填充率为什么要分开算
这两个指标一合并成「成功率」,你就再也发现不了降质。它们的分母不同:
- 覆盖率:我要的那批 ASIN 里,实际返回了多少条?分母是任务清单。
- 填充率:返回来的记录里,多少条过了 P0 校验?分母是返回记录。
分开之后,coverage 0.98 / fillRate 0.61 这种组合才有意义:清单基本取到了,但返回的记录里有近四成不可用。合并成一个数看,它只是「六成成功」,你会以为清单取漏了,往「覆盖率」的方向排查;实际是多半被拦截面或字段缺失吃掉了,属于填充率的问题。
export interface GateResult {
coverage: number;
fillRate: number;
usable: number;
blocked: number;
ok: boolean;
}
export function qualityGate(
wanted: ReadonlySet<string>,
rows: Array<{ raw: unknown; asin?: string }>,
minCoverage = 0.98,
minFill = 0.95,
): GateResult {
const got = new Set(rows.map((r) => r.asin).filter(Boolean) as string[]);
const intersected = [...got].filter((a) => wanted.has(a));
const usableRows = rows.filter((r) => assertP0(r.raw));
const coverage = wanted.size === 0 ? 0 : intersected.length / wanted.size;
const fillRate = rows.length === 0 ? 0 : usableRows.length / rows.length;
return {
coverage,
fillRate,
usable: usableRows.length,
blocked: rows.length - usableRows.length,
ok: coverage >= minCoverage && fillRate >= minFill,
};
}
ok: false 的时候拦下整批,不要往库里写。半批脏数据比没有数据更贵,因为它看起来像成功了,下游会照常用它出报表。阈值的具体标定方法见下面这一小节。
7.1 阈值怎么定,别拍脑袋
阈值按业务场景分,不要全站一个值。定法是从一段干净的历史倒推:取过去两周你认为「可以用」的那些批次,算它们的覆盖率与填充率分布,把阈值定在第 5 百分位附近,而不是平均值。这样日常批次能够通过,异常批次会被拦下。
几个场景的量级参考:价格监控要求对齐度高,覆盖率不低于 98%、P0 填充率不低于 95%;竞品评论采集看重趋势而非精确值,可放宽到覆盖率 95%、填充率 90%;新品榜与类目榜这类结构会变的数据,覆盖率的权重低于单批能否完整拉取。定完之后写进配置,随契约版本一起管理。
最后留一条人工通道:阈值触发拦截时,要能让人一目了然地看到是哪一项没过、样本长什么样。否则同事只会把阈值往上调,直到这道门形同虚设。
八、翻页、去重与幂等:主键怎么定
主键定错,翻页做得越漂亮,库里的重复行越多。主键四要素:asin + marketplace + capturedAt + contractVersion。前两项定位对象,第三项定位时点,第四项定位口径。
8.1 幂等写入
用 SQLite 时用 better-sqlite3 的预编译语句加事务:
const upsert = db.prepare(`
INSERT INTO product_snapshot
(asin, marketplace, captured_at, contract_version, title, brand, price, currency, rating, review_count)
VALUES
(@asin, @marketplace, @capturedAt, @contractVersion, @title, @brand, @price, @currency, @rating, @reviewCount)
ON CONFLICT(asin, marketplace, captured_at, contract_version) DO UPDATE SET
title = excluded.title, brand = excluded.brand,
price = excluded.price, currency = excluded.currency,
rating = excluded.rating, review_count = excluded.review_count;
`);
const writeBatch = db.transaction((rows: ProductSnapshotT[]) => {
for (const row of rows) upsert.run(row);
});
writeBatch(usableRows); // 同一天重跑,覆盖而不是追加
Postgres 侧是同一套语义,语句换成 ON CONFLICT (asin, marketplace, captured_at, contract_version) DO UPDATE SET ...。要点是冲突目标必须写成复合唯一约束,只认 asin 的话同一天采两次会覆盖掉历史。
8.2 翻页的两个坑
- 页数上限:超出上限的请求会静默返回首页或空页,表现是「数据在增长但增速不对」,日志里全是 200。
- 结果集会漂移:两次翻同一页,顺序可能不同。靠翻页去「发现」集合本身就不稳。
所以任务要由 ASIN 清单驱动,翻页只用来补齐清单内部的明细。清单从哪里来、多久刷新一次,是另一个可以单独讨论的问题,但「任务必须有明确的集合」这一点没有例外。
function dedupe<T extends { asin: string; marketplace: string; capturedAt: string }>(
rows: T[],
): T[] {
const seen = new Set<string>();
const out: T[] = [];
for (const r of rows) {
const key = `${r.asin}|${r.marketplace}|${r.capturedAt}`;
if (seen.has(key)) continue;
seen.add(key);
out.push(r);
}
return out;
}
数据管道这一侧的边界怎么划,亚马逊数据管道里写过。跨页重复率本身是一个值得监控的指标:同一批任务里出现大量重复 ASIN,通常意味着翻页键失效或上游开始回滚游标,这在单个请求的日志里看不出来。
九、批处理与背压:清单大了之后会发生什么
几十个 ASIN 的时候,Promise.all(asins.map(...)) 很好用。清单到四位数之后,同一个写法会变成一次性的资源事故。
9.1 一次性 map 出去的代价
Promise.all(asins.map(fetchOne)) 会立刻为全部元素发起请求。四千个 ASIN 等于四千个同时在飞的 HTTP 请求和四千份未解析的响应体,内存先于带宽出问题,而且失败时你拿不到「已经稳妥落库的部分」——因为它要么全成功,要么抛出第一万个异常。
9.2 分批 + 检查点
做法是把大批次切成小块,每块结束就写一次库并落一条进度:这样重跑时可以从断点继续,而不是从头再来。
import { chunk } from "./util";
async function runBatch(sourceAsins: string[], size = 200) {
for (const [i, group] of chunk(sourceAsins, size).entries()) {
const results = await Promise.allSettled(group.map(fetchOne));
const rows = results
.filter((r): r is PromiseFulfilledResult<Row> => r.status === "fulfilled")
.map((r) => r.value);
const gate = qualityGate(new Set(group), rows);
if (!gate.ok) {
await checkpoint(i, "blocked", gate);
continue; // 这块不入库,留给下一轮
}
writeBatch(rows);
await checkpoint(i, "ok", gate);
}
}
checkpoint 那一行是这段的核心。它让「跑到第 17 块断掉」变成一个可以续跑的状态,而不是一个要重新花两小时的任务。
9.3 p-limit 限的是并发,不是速率
p-limit(8) 保证的是「同时最多 8 个请求在飞」,它不限速率。同一个配置在两种网络状况下的实际速率可以差几百倍:每个请求 10 ms 返回时约 800 req/s,每个请求 5 s 返回时只剩约 1.6 req/s。而上游给你的额度通常按每秒请求数计,所以只用 p-limit 的结果是——响应变快时你超限被限流,响应变慢时你又根本没跑满。
import PQueue from "p-queue";
// 并发上限与速率上限各管一边
const queue = new PQueue({
concurrency: 8, // 同时在飞的请求数
interval: 1_000, // 每 1 秒
intervalCap: 4, // 最多放行 4 个
});
await queue.addAll(jobs.map((job) => () => fetchOne(job)));
两个参数分别对齐两个约束:concurrency 对应你自己的连接资源与对方的连接容忍;intervalCap 对应上游额度。调优顺序是先贴近额度定速率,再调并发去消化它——反过来做会出现「并发很高但速率被卡住,排队时间一直在涨」。
还有一层是消费端的反压。如果写库的速度跟不上请求速度,队列会无限增长。Node 在这里的优势是原生异步、单请求内存开销小;代价恰恰也是它——内存慢慢涨上去时不会有任何报错,直到某次超时才暴露。所以生产任务里我倾向于给队列一个显式上限,超限时放慢生产,而不是让它在内存里膨胀。
十、优雅停机:跑到一半收到 SIGTERM 怎么办
这一条在容器环境里比在传统服务器上重要得多。docker stop 与 kubectl rollout 都是先发 SIGTERM,等待一段时间再 SIGKILL。不做处理的话,你可能会得到半写的一批记录,或者丢失还没有入库的内存数据。
10.1 三段式处理
let shuttingDown = false;
process.on("SIGTERM", () => {
shuttingDown = true;
log.info("SIGTERM received, draining");
});
async function loop(jobs: Job[]) {
for (const job of jobs) {
if (shuttingDown) {
await flush(); // 把已取到但未入库的部分写完
log.info("stopped cleanly, resume from checkpoint next run");
return;
}
await handle(job);
}
}
关键在于 flush():已经过了质量门但还没写库的数据要先落地。配合上一节的 checkpoint,下一次启动可以从断点继续,两次运行之间不会漏也不会重。
10.2 超时与信号的组合
Node 17.3 之后可以用 AbortSignal.timeout(ms) 给单次请求设上限,也可以用 AbortSignal.any([signalA, signalB]) 把「停机信号」和「超时」合成一个取消源——这在 Node 22 之后可用。组合的价值是:停机的时候,正在飞的请求会被主动取消,而不是等它们自己超时。
const stopController = new AbortController();
process.on("SIGTERM", () => stopController.abort());
const res = await fetch(url, {
signal: AbortSignal.any([AbortSignal.timeout(30_000), stopController.signal]),
});
部署侧注意等待时长:容器编排默认的优雅停机宽限期通常是 30 秒左右。如果你的单块任务平均要跑 40 秒,要么把块调小,要么把宽限期调长,否则永远靠不大可能在 30 秒内完成刷写。
十一、结构化日志:出事那天靠什么复盘
日志这一项容易被写成 console.log(url),到排查时才发现能看到的只有一堆 200。判断标准很简单:出事之后能不能只靠日志,在不重新跑一遍的前提下定位到是哪个市场、哪一批、哪一层。
11.1 一条记录要带哪些字段
| 字段 | 作用 |
|---|---|
traceId | 串起一次任务的所有请求,用来还原现场 |
asin / marketplace | 定位对象,没有这两个字段的日志无法用来复盘 |
attempt | 第几次成功,成本公式里的变量 |
status / errorName | 失败分类的依据,而不是靠消息文本猜 |
durationMs | 端到端时延的样本,p95 从这里算 |
gate | 质量门的输出,记录当时为什么被拦下 |
11.2 一份能用的样本
用结构化日志库(如 pino)输出 JSON,一行一条,便于后续聚合:
import pino from "pino";
const log = pino({ base: { service: "amazon-collector" } });
// 一次调用结束后的落点
log.info({
traceId, asin, marketplace, capturedAt,
attempt: outcome.attempts,
status: outcome.status,
errorName: outcome.error ?? null,
durationMs: Date.now() - startedAt,
gate: gateResult, // { coverage, fillRate, usable, blocked, ok }
});
实际落到日志里是这一行:
{"level":30,"time":1788940800000,"service":"amazon-collector",
"traceId":"f47ac10b","asin":"B08N5WRWNW","marketplace":"DE",
"capturedAt":"2026-09-09","attempt":2,"status":200,"errorName":null,
"durationMs":1843,"gate":{"coverage":0.98,"fillRate":0.61,
"usable":610,"blocked":390,"ok":false}}
最后那个 gate.ok: false 是这份日志的价值所在。三周之后业务方问起「德国站的价格为什么对不上」,你能在这一行里看到当时那一批只有六成通过 P0 校验,而不是先去猜。
11.3 从日志反推那三个数
有了这些字段,成本与质量指标可以不另起一套系统:
-- 近七天,按天看质量与重跑成本
SELECT captured_at,
AVG(attempt)::numeric(4,2) AS avg_attempts,
COUNT(*) FILTER (WHERE status = 200) * 1.0 / COUNT(*) AS http_ok_rate,
AVG((gate->>'fillRate')::float) AS fill_rate
FROM collector_log
WHERE captured_at >= current_date - 7
GROUP BY captured_at
ORDER BY captured_at;
盯三条曲线就够了:avg_attempts 抬升说明上游在抖动或你在撞拦截,http_ok_rate 单独看没意义(拦截页也返回 200),fill_rate 下滑几乎总是字段契约或渲染层出了问题。三者同时异常通常是出口 IP 或指纹那一层被识别。
十二、成本埋点:每千条可用记录到底多少钱
这一层必须自己埋点,因为分子在供应商那里,分母只在你自己的日志里。
每千条可用记录成本 = 1000 × 单价 × 平均尝试次数 × (1 + 渲染倍率)
────────────────────────────────────────────
成功率 × 填充率 × (1 - 拦截率)
+ 每千条的代理流量费 + 每千条的工程维护摊销
分母三个数必须取自你自己的日志,不能用供应商给出的平均值——那些数字不会替你区分失败率与被拦率。
分母三个数必须取自你自己的日志,不能用供应商给出的平均值。代理流量按 GB 计,失败请求同样计入;工程维护摊销是最容易漏的一项,一次风控升级通常是两三个人日。
class MeteredClient extends AmazonClient {
credits = 0;
attempts = 0;
override async fetch<T>(kind: EndpointKind, params: Record<string, string>) {
const res = await super.fetch<T>(kind, params);
this.credits += 1;
this.attempts += res.attempts;
return res;
}
costPer1kUsable(usable: number, unitPrice: number): number {
if (usable === 0) return Infinity;
return (this.credits * unitPrice * 1000) / usable;
}
}
把报价口径换算成真实账单这件事,我们在另一篇里把六种计费单位拆开对比过,见亚马逊数据 API 的定价口径拆解。
十三、四个指标,两条告警
指标四个就够,多了没人看:
| 指标 | 含义 | 为什么需要单独看 |
|---|---|---|
| 覆盖率 | 清单里实际返回的比例 | 区分「没取到」与「取到了但不可用」 |
| P0 填充率 | 返回记录里通过校验的比例 | 取到了但不可用 |
| p95 端到端时延 | 单条记录从发到落库 | 先于超时报错出现退化 |
| 每千条可用记录成本 | 上面那个公式的输出 | 唯一能把账单和业务量对齐的数 |
告警两条:连续两个周期覆盖率或 P0 填充率跌破阈值(数据质量事故);每千条可用记录成本环比涨超 30%(预算事故)。两者都会在业务侧投诉之前先出现在日志里。
十四、生产部署:PM2、Docker 与 cron 的重启语义
部署方式决定失败怎么被处理,这一点常被低估。
| 运行方式 | 失败后的行为 | 需要注意 |
|---|---|---|
cron 直接拉进程 | 下次触发才重试,期间无感知 | 要自己写锁,防止上一轮没结束就叠下一轮 |
| PM2 常驻 + 内部定时器 | 崩溃后自动重启 | 重启会丢掉内存状态,检查点必须落盘 |
| Docker + 编排平台 | 发 SIGTERM,宽限期后强杀 | 宽限期要大于单块任务耗时 |
| Serverless 定时触发 | 有硬超时上限 | 清单要在超时前切完,靠分批而不是靠加并发 |
共同的一条:不管用哪种方式,都要把「上一次跑到哪一块」写到外部存储里。依赖进程内存的进度在重启之后会归零,表现为每次崩溃都从头跑一遍,账单翻倍而数据没多。
十五、哪些可以不自己写
把前面各节做完之后,你会发现属于业务的只有三件事:定义字段契约、定义质量阈值、定义数据怎么用。其余都是「为了拿到干净数据而不得不建的基础设施」。
| 层 | 自托管要做什么 | Pangolinfo 侧 |
|---|---|---|
| TLS / HTTP/2 指纹 | 选客户端、锁定 profile、跟踪浏览器版本 | 包含,随浏览器版本更新 |
| 浏览器指纹一致性 | 维护设备 profile,对齐渲染串 | 包含 |
| 住宅与移动 IP 出口 | 买流量、维护池子、盯信誉衰减 | 包含,不单独计费 |
| JS 渲染 | 维护浏览器集群,判断哪些字段需要渲染 | 按端点在服务端完成,无渲染倍率 |
| 拦截与验证码 | 识别拦截页、换出口、重试 | 服务端处理与重试 |
| 地域与邮区对齐 | 按市场准备出口节点 | 包含 |
换句话说,一次请求的费用覆盖了一切:Pangolinfo Amazon Scraper API 把这些收在服务端,住宅 IP、指纹伪装、JS 渲染、验证码处理都不拆成加价项。你这一侧仍然是普通的 fetch,区别是省掉了第六到第十节里那部分与业务无关的工作,守住 P0 契约、质量门和幂等落地。
这套结构可以先用免费额度跑一遍:注册后前 60 次请求免费,不需要信用卡。拿 20 个真实 ASIN 跑通请求层与质量门,再决定要不要自己扛指纹与 IP 那几层。
十六、亚马逊数据 API Node.js 任务上线前的自查清单
五行,少任何一行都会在第三周之后慢慢返工:
- 超时是否按建连 / 响应头 / 响应体分别设置,而不是一个值包打天下;
- 是否有一处对 P0 字段的运行时断言,而不只是 TypeScript 类型;
- 主键是否包含采集日期与契约版本,写入是否为 upsert;
- 错误是否按瞬时 / 限流 / 拦截 / 契约四类分开处理,而不是一个
catch兜住; - 每次调用是否记了数与额,能不能算出每千条可用记录成本。
前面几节解决「能不能拿到」,后面几节解决「能不能每天对齐」。两者之间的边界,就是「完成任务」和「完成任务并且可以放心离开」的距离。
如果你的团队正处在选型阶段,亚马逊数据 API 全景把五条路线的授权边界与成本结构逐条列过;从零自建和调用数据 API 的分界线,见自建抓取与调用数据 API 的成本分界。
十七、常见问题
Node 内置的 fetch 可以直接调亚马逊数据 API 吗?
可以发通,但两层要自己配:一是 Node 内置 fetch 默认不协商 HTTP/2,需要显式设置 allowH2 的 agent;二是它的 TLS 握手来自 Node 自带的 OpenSSL,与浏览器指纹不一致,强保护页面上会被识别。做原型够快,放到日均任务上要补一层指纹。
换了 UA 和代理还是被拦,问题在哪?
九成在两处:TLS 与 HTTP/2 指纹仍是 Node 客户端的;或者信号互相矛盾——UA 说 Windows 而 HTTP/2 窗口像 Linux,IP 在德国却发 en-US。用公开指纹端点打一次确认 JA4,再核对出口国家、站点、语言是否同源。
undici 的 allowH2 要不要开?
按 endpoint 决定。开了之后并发模型会从「连接数」变成「单连接内的流数量」,受 maxConcurrentStreams 与流控窗口约束,此时单纯调高应用层并发不一定涨吞吐。批量场景建议开,先定连接池大小再调应用层队列。
有了 TypeScript 类型还需要 Zod 吗?
需要。类型标注只在编译期约束你自己的代码,不知道运行时送来的 JSON 里有没有这个字段、类型和长度对不对。缺少运行时校验时,上游把 price 从数字改成字符串不会被发现,只会让下游报表安静地偏离。
并发设多大合适?p-limit 和连接池怎么配合?
两层叠加:应用层 p-limit 决定同时放多少请求出去,传输层的连接池与 HTTP/2 参数决定这些请求怎么落到连接上。调参时固定应用层,先调连接数,盯 429 占比与端到端时延,吞吐不再涨就退回上一档。
任务中断后重跑,会不会产生重复数据?
取决于主键与写入方式。主键要包含 asin、marketplace、采集日期和契约版本,写入用 ON CONFLICT DO UPDATE 而不是追加。配合每块结束后写一次检查点,重跑是覆盖同日数据并从断点续跑,不会重复。
Node 的错误应该怎么分类重试?
瞬时错误(连接超时、body 超时、DNS 失败)与限流(429、带 Retry-After 的 503)可重试;403 与拦截页不能按瞬时错误重试;200 但字段缺失属契约错误,重试几次都是同样的空字段。JSON 解析失败通常是返回了 HTML 拦截页。
部署到 Docker 或 K8s 要注意什么?
一是优雅停机,编排平台发 SIGTERM 后有宽限期,要在这个窗口里把已取到但未入库的数据刷写落盘;二是把「跑到第几块」写到外部存储,别放在内存里,重启后进度归零会导致每次崩溃都从头跑,账单翻倍而数据没多。
Pangolinfo 的费用里包含住宅 IP 吗?要不要另外买代理?
包含,而且不止住宅 IP。一次请求的费用覆盖了一切:住宅与移动 IP 的出口与轮换、TLS 与 HTTP/2 指纹、浏览器指纹一致性、需要时的 JS 渲染、拦截与验证页的服务端处理与重试,都不作为独立加价项出现。你发一个请求,收一份结构化的实时 JSON。
浏览器指纹、JS 渲染、验证码处理会不会单独计费?
不会。指纹伪装、需要时的 JS 渲染、验证页的服务端处理与重试,都包含在一次请求的计费口径内,没有渲染倍率、端点难度倍率或住宅 IP 加价。这也是我们与按加价项叠加计费的方案在账单结构上的主要差别。
只有几十个 ASIN,值得写这么多结构吗?
值得,但只用其中三段:请求层的超时与重试、P0 字段校验、带日期主键的快照表,大约 60 行 TypeScript。换来的是失败之后敢重跑,以及三个月后能画出价格曲线。并发控制与令牌桶等清单上到四位数再加。
