亚马逊数据 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/2TLS 指纹2026 年的状态
fetch(Node 内置)Node 20 默认不走 h2,需显式 agentNode 自己的 OpenSSL 握手零依赖,但两层都要自己配
undiciallowH2: 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 与原生边界的开销而非网络延迟。数值会随版本变化,引用前请按当时版本自行复测。

引擎可模拟的最新 ChromeHTTP/2 指纹是否正确req/s冷启动
wreq-jsRust wreq + BoringSSL,进程内149正确128427 ms
imperscurl-impersonate,进程内146正确843916 ms
node-wreq同上 Rust 内核149正确650010 ms
impitRust reqwest + 打补丁的 rustls124不正确671037 ms
CycleTLSGo 子进程 + 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 与 platformUA 写 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 时要认名字而不是只看状态码:

错误来源典型表现处理
UndiciErrorUND_ERR_CONNECT_TIMEOUT连接阶段超时可重试
UndiciErrorUND_ERR_HEADERS_TIMEOUT响应头迟迟不来可重试
UndiciErrorUND_ERR_BODY_TIMEOUT响应体流式读取超时可重试,但要记录 body 大小
DOMExceptionTimeoutErrorAbortSignal.timeout 触发可重试
ENOTFOUND / EAI_AGAINDNS 解析失败可重试,连续出现要查出口网络
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 stopkubectl 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。换来的是失败之后敢重跑,以及三个月后能画出价格曲线。并发控制与令牌桶等清单上到四位数再加。

微信扫一扫
与我们联系

QR Code
快速测试

联系我们,您的问题,我们随时倾听

无论您在使用 Pangolin 产品的过程中遇到任何问题,或有任何需求与建议,我们都在这里为您提供支持。请填写以下信息,我们的团队将尽快与您联系,确保您获得最佳的产品体验。

Talk to our team

If you encounter any issues while using Pangolin products, please fill out the following information, and our team will contact you as soon as possible to ensure you have the best product experience.