Node 后端实战 · Serverless 导出 CSV 总超时?用 Queue + R2 异步任务彻底解决

各位看官,今天聊一个在 Serverless / 边缘运行时里几乎绕不开、但又特别容易写错的需求:大批量数据导出。

我这个后端是跑在 Cloudflare Workers 上的,技术栈 Hono + D1(边缘 SQLite) + R2 + Queue。某天产品说"能不能让用户把几万条线索导成 CSV"。我第一反应很朴素:查出来,拼成字符串,直接返回。结果一上线就发现,这个想法在边缘环境里根本行不通。这篇文章就把这次踩坑和最终落地的异步导出方案讲透。

在这里插入图片描述

一、为什么不能同步导出

在自有服务器上,你 SELECT * 然后慢慢拼 CSV,顶多慢点,用户等得久而已。但在 Workers 这种边缘运行时,同步导出有三道硬墙:

约束具体限制同步导出会怎样
CPU 时间单次请求 CPU 约 50ms(免费档)/有限档拼几万行直接超时,请求被掐
内存单实例约 128MB把全量行塞进一个字符串,OOM 回收
实例生命周期请求结束即可能被回收长任务中途被杀,前功尽弃
响应体走单一 HTTP 响应边查边等,用户干瞪眼还占着连接

核心矛盾就一句:边缘实例是"短命、廉价、可被随时回收"的,你没法指望它一口气干完一个重活。

所以正确姿势是"接单即返回,后台慢慢干,干完了通知你来拿"——也就是异步导出。整个链路长这样:

客户端 POST /export/tasks  ──→  写 pending 任务 + 投递 Queue  ──→  立即返回 {id, status}
                                                                      │
                                          Queue consumer 反复消费同一 taskId  │
                                                                      ▼
                                          游标分页查询 → 分片写 R2 → 续跑
                                                                      │
                                                                      ▼
                                          标记 done + 写入 r2Key/expiresAt
客户端 GET /tasks 轮询状态  ──→  status: processing(带进度) / done
客户端 GET /tasks/:id/download ──→  R2 预签名直链 或 流式代理

二、幂等 + 游标续跑:at-least-once 的必然代价

Queue 的投递语义是 at-least-once——同一条消息可能被消费两次。这意味着你的消费函数必须幂等,否则重复导出会把 R2 文件和数据库状态搞乱。

在这里插入图片描述

我的做法很直接:终态任务直接跳过。

export const processExportTask = async (env: Bindings, taskId: string): Promise<void> => {
  const task = await db.query.exportTasks.findFirst({ where: eq(exportTasks.id, taskId) });
  if (!task) { console.warn("[export] task not found, skip:", taskId); return; } // ack,不重试
  // at-least-once 幂等:终态任务直接跳过
  if (task.status === "done" || task.status === "expired" || task.status === "failed") return;
  // ...
};

第二个坑是"续跑"。Worker 实例随时可能被回收,一个几万行的任务不可能一次消费跑完。我用游标分页而不是 OFFSET

const baseWhere = (lastId: string): SQL => {
  const conds = [eq(table.tenantId, tid), isNull(table.deletedAt), ...conditions];
  if (lastId) conds.push(gt(table.id, lastId)); // 游标:只取 id 比上次大的
  return and(...conds)!;
};
// 每批查 EXPORT.CHUNK 行,按 id 升序,把最后一行 id 记进 lastId
const rows = await db.select().from(table).where(baseWhere(lastId)).orderBy(asc(table.id)).limit(EXPORT.CHUNK);
// ...
const finished = rows.length < EXPORT.CHUNK;

为什么不用 OFFSET?深分页时 OFFSET 50000 数据库要从头数五万行再跳过,越往后越慢,而且续跑期间数据有增删会导致错位。游标用 gt(id, lastId) 走主键索引,复杂度稳定 O(分页大小),还能安全断点续跑。

每跑完一批,把进度落库:

await db.update(exportTasks).set({
  processedRows: processed,
  lastId: lastId || null,   // 断点游标
  r2Parts: JSON.stringify(parts),
  updatedAt: nowSec(),
}).where(eq(exportTasks.id, taskId));

if (!finished) {
  if (env.EXPORT_QUEUE) await env.EXPORT_QUEUE.send({ taskId }); // 没跑完 → 再投一条给自己
  return;
}

lastId 持久化到 D1,下次消费从断点接着查。任务状态机如下:

状态含义进入条件
pending待处理创建任务初始状态
processing处理中首次进入消费,且未完成
done完成所有分片上传并 complete 成功
failed失败过滤条件非法 / 处理抛错(记 error)
expired已过期清理done 且 expiresAt 到期,R2 文件已删

三、R2 分片续跑:跨调用不重头来

R2 的 multipart upload 有个绕不过去的规则:除最后一片外,每个 part 必须 ≥ 5MB。几万行 CSV 累到 5MB 才传一片,期间如果实例被回收,uploadId 和已传 parts 就没了,下次得从头传。

在这里插入图片描述

解法同样是把中间状态持久化。我让 r2UploadIdr2Parts 都进 D1:

let uploadId = task.r2UploadId;
let parts: UploadPart[] = task.r2Parts ? JSON.parse(task.r2Parts) : [];
let mpu: R2MultipartUpload;
if (!uploadId) {
  mpu = await env.BUCKET.createMultipartUpload(finalKey, { /* httpMetadata */ });
  uploadId = mpu.uploadId;
  await db.update(exportTasks).set({ r2UploadId: uploadId, updatedAt: nowSec() }).where(eq(exportTasks.id, taskId));
} else {
  mpu = env.BUCKET.resumeMultipartUpload(finalKey, uploadId); // 断点恢复
}
持久化字段作用不持久化的后果
lastId查询游标断点重复导出已处理的行
r2UploadId复用未完成的 multipart每次从头开新上传,浪费流量
r2Parts已传分片的 ETagcomplete 时缺片,文件损坏
processedRows进度分子前端进度条无法续算
totalRows进度分母进度百分比失真

buffer 攒够 5MB 就传一片,末片允许小于 5MB,最后按 partNumber 升序 complete

parts.sort((a, b) => a.partNumber - b.partNumber);
await mpu.complete(parts);

失败时要记得 abort 掉未完成的 multipart,否则 R2 会一直留着半截文件占成本:

const markFailed = async (env, task, error) => {
  if (task.r2UploadId) {
    try { await env.BUCKET.resumeMultipartUpload(`exports/${task.tenantId}/${task.id}.csv`, task.r2UploadId).abort(); } catch {}
  }
  await db.update(exportTasks).set({ status: "failed", error: error.slice(0, 500), updatedAt: nowSec() }).where(eq(exportTasks.id, task.id));
};

顺带一提,文件 key 是 exports/${tenantId}/${taskId}.csv——把租户 ID 焊进路径前缀,和上一篇讲的租户隔离一脉相承,连导出文件都按租户分目录,下载时天然隔离。

四、下载双模式:预签名直连 vs 流式代理

文件在 R2 上,用户怎么拿?我给了两条路:

模式触发条件客户端行为Worker 成本
presigned配置了 R2 API 凭证拿到直链,浏览器直连 R2 边缘下载零带宽、零 CPU
inline未配置凭证Worker 拉 R2 对象流式转发走 Worker 出口带宽

预签名直连是首选——客户端直连 R2 边缘,Worker 完全不碰文件流,带宽压力为零。但 Cloudflare 的 R2 绑定不提供 createPresignedUrl,官方推荐改用 S3 兼容 API。我直接用纯 Web Crypto 手搓了 SigV4 签名,零第三方 SDK:

// 派生签名密钥:AWS4<secret> → date → region(auto) → service(s3) → aws4_request
const kDate = await hmac(new TextEncoder().encode(`AWS4${o.secretAccessKey}`), dateStamp);
const kRegion = await hmac(kDate, "auto");
const kService = await hmac(kRegion, "s3");
const kSigning = await hmac(kService, "aws4_request");
const signature = toHex(await hmac(kSigning, stringToSign));
return `https://${host}${path}?${query}&X-Amz-Signature=${signature}`;

这里有个真实踩坑:R2 要求把 response-content-disposition 这类"响应头覆盖"参数也签进 canonical request,而且参数名要按 ASCII 排序。漏签或排错序,下载直接 403 SignatureDoesNotMatch。我第一次就栽在这,回传的 CanonicalRequest 对照着改完才通。

没配凭证时回落为流式代理——Worker 拉对象原样转发 ReadableStream,不落内存,立即可用。两条路响应形态不同,前端按 mode 字段分流即可。

五、防滥用:限流 + 去重

导出是重活,必须挡住滥用。我没用 KV,直接用 D1 计数,简单够用:

export const checkExportRateLimit = async (db, createdBy) => {
  const since = nowSec() - EXPORT.RATE_WINDOW_SEC;
  const rows = await db.select({ n: count() }).from(exportTasks)
    .where(and(eq(exportTasks.createdBy, createdBy), gte(exportTasks.createdAt, since)));
  if (Number(rows[0]?.n ?? 0) >= EXPORT.RATE_LIMIT_PER_HOUR) {
    throw err("RATE_LIMIT", `export limit exceeded: ${EXPORT.RATE_LIMIT_PER_HOUR} per hour`);
  }
};

创建任务时还做了去重:同一用户、同租户、同类型、同格式、同筛选条件且还在 pending/processing 的,直接返回既有任务,不重复开干:

const existing = await db.query.exportTasks.findFirst({
  where: and(
    eq(exportTasks.tenantId, input.tenantId),
    eq(exportTasks.createdBy, input.createdBy),
    eq(exportTasks.type, input.type),
    eq(exportTasks.format, input.format),
    input.filters ? eq(exportTasks.filters, input.filters) : isNull(exportTasks.filters),
    inArray(exportTasks.status, ["pending", "processing"]),
  ),
  orderBy: [desc(exportTasks.createdAt)],
});
if (existing) return existing; // 去重:同条件复用

六、CSV 生成的两个细节

最后说说 CSV 本身,两个容易翻车的小点:

一是 UTF-8 BOM。不写 BOM,Excel 打开中文表头就是乱码。我在表头最前面拼一个 \uFEFF

export const csvHeader = (columns) => "" + columns.map(c => csvEscape(c.key)).join(",") + "\r\n";

二是 反显别搞出 N+1。导出要把 projectIdownerId 这种外键翻成中文名。我在 consumer 端一次性把维度表全查出来建 Map,遍历行时 O(1) 反查:

const [projRows, catRows, ownerRows] = await Promise.all([
  db.select({ id: projects.id, name: projects.name }).from(projects).where(...),
  db.select({ id: leadCategories.id, name: leadCategories.name }).from(leadCategories).where(...),
  db.select({ id: users.id, name: users.name }).from(users).where(eq(users.tenantId, tid)),
]);
return { projects: new Map(...), categories: new Map(...), owners: new Map(...) };

字段还做了白名单和转义——含逗号、引号、换行的字段整体双引号包裹、内部引号翻倍,避免 CSV 结构被脏数据冲垮。

小结

边缘 Serverless 下做大批量导出,关键是承认"实例会死",把一切中间状态外置:任务状态进 D1,上传进度进 D1,查询游标进 D1。Queue 只负责反复"叫醒"同一个 taskId,真正的进度靠持久化的 lastId / r2UploadId / r2Parts 续命。幂等挡住 at-least-once 的重复投递,TTL 清理挡住 R2 成本膨胀。这套下来,几万行导出稳稳当当,用户还能实时看到进度条。

各位看官,下一篇我打算聊聊限流和审计日志里"敏感字段自动脱敏"那点事——毕竟导出明文手机号已经够刺激了,日志里再漏一份就更热闹了。


相关阅读

本文由 FungLeo 主导,Deepseek 优化校阅,转发请注明首发地址,谢谢大家!

更多推荐