ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

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

Node 后端实战 · Serverless 导出 CSV 总超时?用 Queue + R2 异步任务彻底解决 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 文件和数据库状态搞乱。我的做法很直接终态任务直接跳过。exportconstprocessExportTaskasync(env:Bindings,taskId:string):Promisevoid{consttaskawaitdb.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.statusdone||task.statusexpired||task.statusfailed)return;// ...};第二个坑是续跑。Worker 实例随时可能被回收一个几万行的任务不可能一次消费跑完。我用游标分页而不是OFFSETconstbaseWhere(lastId:string):SQL{constconds[eq(table.tenantId,tid),isNull(table.deletedAt),...conditions];if(lastId)conds.push(gt(table.id,lastId));// 游标只取 id 比上次大的returnand(...conds)!;};// 每批查 EXPORT.CHUNK 行按 id 升序把最后一行 id 记进 lastIdconstrowsawaitdb.select().from(table).where(baseWhere(lastId)).orderBy(asc(table.id)).limit(EXPORT.CHUNK);// ...constfinishedrows.lengthEXPORT.CHUNK;为什么不用OFFSET深分页时OFFSET 50000数据库要从头数五万行再跳过越往后越慢而且续跑期间数据有增删会导致错位。游标用gt(id, lastId)走主键索引复杂度稳定 O(分页大小)还能安全断点续跑。每跑完一批把进度落库awaitdb.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)awaitenv.EXPORT_QUEUE.send({taskId});// 没跑完 → 再投一条给自己return;}lastId持久化到 D1下次消费从断点接着查。任务状态机如下状态含义进入条件pending待处理创建任务初始状态processing处理中首次进入消费且未完成done完成所有分片上传并complete成功failed失败过滤条件非法 / 处理抛错记 errorexpired已过期清理done 且expiresAt到期R2 文件已删三、R2 分片续跑跨调用不重头来R2 的multipart upload有个绕不过去的规则除最后一片外每个 part 必须 ≥ 5MB。几万行 CSV 累到 5MB 才传一片期间如果实例被回收uploadId 和已传 parts 就没了下次得从头传。解法同样是把中间状态持久化。我让r2UploadId和r2Parts都进 D1letuploadIdtask.r2UploadId;letparts:UploadPart[]task.r2Parts?JSON.parse(task.r2Parts):[];letmpu:R2MultipartUpload;if(!uploadId){mpuawaitenv.BUCKET.createMultipartUpload(finalKey,{/* httpMetadata */});uploadIdmpu.uploadId;awaitdb.update(exportTasks).set({r2UploadId:uploadId,updatedAt:nowSec()}).where(eq(exportTasks.id,taskId));}else{mpuenv.BUCKET.resumeMultipartUpload(finalKey,uploadId);// 断点恢复}持久化字段作用不持久化的后果lastId查询游标断点重复导出已处理的行r2UploadId复用未完成的 multipart每次从头开新上传浪费流量r2Parts已传分片的 ETagcomplete时缺片文件损坏processedRows进度分子前端进度条无法续算totalRows进度分母进度百分比失真buffer 攒够 5MB 就传一片末片允许小于 5MB最后按partNumber升序completeparts.sort((a,b)a.partNumber-b.partNumber);awaitmpu.complete(parts);失败时要记得 abort 掉未完成的 multipart否则 R2 会一直留着半截文件占成本constmarkFailedasync(env,task,error){if(task.r2UploadId){try{awaitenv.BUCKET.resumeMultipartUpload(exports/${task.tenantId}/${task.id}.csv,task.r2UploadId).abort();}catch{}}awaitdb.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 边缘下载零带宽、零 CPUinline未配置凭证Worker 拉 R2 对象流式转发走 Worker 出口带宽预签名直连是首选——客户端直连 R2 边缘Worker 完全不碰文件流带宽压力为零。但 Cloudflare 的 R2 绑定不提供createPresignedUrl官方推荐改用 S3 兼容 API。我直接用纯 Web Crypto 手搓了 SigV4 签名零第三方 SDK// 派生签名密钥AWS4secret → date → region(auto) → service(s3) → aws4_requestconstkDateawaithmac(newTextEncoder().encode(AWS4${o.secretAccessKey}),dateStamp);constkRegionawaithmac(kDate,auto);constkServiceawaithmac(kRegion,s3);constkSigningawaithmac(kService,aws4_request);constsignaturetoHex(awaithmac(kSigning,stringToSign));returnhttps://${host}${path}?${query}X-Amz-Signature${signature};这里有个真实踩坑R2 要求把response-content-disposition这类响应头覆盖参数也签进 canonical request而且参数名要按 ASCII 排序。漏签或排错序下载直接 403SignatureDoesNotMatch。我第一次就栽在这回传的CanonicalRequest对照着改完才通。没配凭证时回落为流式代理——Worker 拉对象原样转发ReadableStream不落内存立即可用。两条路响应形态不同前端按mode字段分流即可。五、防滥用限流 去重导出是重活必须挡住滥用。我没用 KV直接用 D1 计数简单够用exportconstcheckExportRateLimitasync(db,createdBy){constsincenowSec()-EXPORT.RATE_WINDOW_SEC;constrowsawaitdb.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){throwerr(RATE_LIMIT,export limit exceeded:${EXPORT.RATE_LIMIT_PER_HOUR}per hour);}};创建任务时还做了去重同一用户、同租户、同类型、同格式、同筛选条件且还在pending/processing的直接返回既有任务不重复开干constexistingawaitdb.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)returnexisting;// 去重同条件复用六、CSV 生成的两个细节最后说说 CSV 本身两个容易翻车的小点一是UTF-8 BOM。不写 BOMExcel 打开中文表头就是乱码。我在表头最前面拼一个\uFEFFexportconstcsvHeader(columns)columns.map(ccsvEscape(c.key)).join(,)\r\n;二是反显别搞出 N1。导出要把projectId、ownerId这种外键翻成中文名。我在 consumer 端一次性把维度表全查出来建Map遍历行时 O(1) 反查const[projRows,catRows,ownerRows]awaitPromise.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:newMap(...),categories:newMap(...),owners:newMap(...)};字段还做了白名单和转义——含逗号、引号、换行的字段整体双引号包裹、内部引号翻倍避免 CSV 结构被脏数据冲垮。小结边缘 Serverless 下做大批量导出关键是承认实例会死把一切中间状态外置任务状态进 D1上传进度进 D1查询游标进 D1。Queue 只负责反复叫醒同一个 taskId真正的进度靠持久化的lastId/r2UploadId/r2Parts续命。幂等挡住 at-least-once 的重复投递TTL 清理挡住 R2 成本膨胀。这套下来几万行导出稳稳当当用户还能实时看到进度条。各位看官下一篇我打算聊聊限流和审计日志里敏感字段自动脱敏那点事——毕竟导出明文手机号已经够刺激了日志里再漏一份就更热闹了。相关阅读Node 后端实战 · 多租户 SaaS 的数据隔离Node 后端实战 · JWT 双密钥轮转与 token 版本号Node 后端实战 · D1 那些坑Node 后端实战 · Workers 踩坑Node 后端实战 · 架构决策全景一次 Web 服务雪崩的完整排查手记Koa 怎么做 JWT 会话与鉴权MySQL 数据备份与恢复的实战方案RSA 非对称加密在后端服务中的落地本文由 FungLeo 主导Deepseek 优化校阅转发请注明首发地址谢谢大家
返回列表