Node.js读取Kinesis数据流:实时/批量定时读取及代码优化问询
Node.js读取Kinesis数据流:实时/定时批量读取实现与代码优化
核心需求
- 实时读取:数据到达Kinesis后立即处理
- 定时批量读取:每隔指定秒数读取15-20条数据并批量处理
- 优化现有代码的健壮性与可维护性
关键优化点与实现思路
1. 实时读取优化
原代码递归调用会频繁触发无意义的API请求(无数据时也会立即重试),通过给GetRecordsCommand添加WaitTimeSeconds参数(最大5秒),让Kinesis在无数据时等待一段时间再返回,既减少无效API调用,又保证数据到达时能及时响应。
2. 定时批量读取实现
- 维护一个数据缓存队列,每次拉取数据后先存入队列
- 每隔指定时间检查队列,若数据量达到15条以上或到时间阈值,就批量处理并清空队列
- 控制
GetRecordsCommand的Limit参数为20,确保每次拉取不超过目标批量数
3. 通用代码优化
- 抽离配置常量,便于后续修改
- 提取数据解析与格式化的公共函数,避免重复代码
- 替换递归调用为定时器循环,避免栈溢出风险
- 增强错误处理,添加基础重试逻辑(针对Kinesis节流或临时错误)
完整实现代码
const { KinesisClient, GetRecordsCommand, GetShardIteratorCommand, } = require("@aws-sdk/client-kinesis"); // 配置常量抽离,可根据需求修改 const CONFIG = { region: "ap-south-1", streamName: "my-data-stream", shardId: "my-shardId", readMode: "real-time", // 可选值: "real-time" | "batch" batchInterval: 5000, // 批量读取间隔(毫秒) batchMinSize: 15, // 批量处理最小条数 batchMaxSize: 20, // 批量处理最大条数 kinesisWaitTime: 1, // 实时读取时的等待时间(秒,0-5) }; const kinesis = new KinesisClient({ region: CONFIG.region }); let currentShardIterator = null; let batchCache = []; let batchTimer = null; // 格式化时间戳 const formatTimestamp = (epoch) => { try { const date = new Date(epoch); const options = { day: "2-digit", month: "short", year: "numeric", hour: "2-digit", minute: "2-digit", second: "2-digit", hour12: false, timeZone: "Asia/Kolkata", }; return date.toLocaleString("en-IN", options); } catch (err) { console.error("时间格式化失败:", err); return epoch; } }; // 解析单条Kinesis记录 const parseRecord = (record) => { try { const parsedData = Buffer.from(record.Data).toString("utf-8"); const payload = JSON.parse(parsedData); return { ...payload, time: formatTimestamp(payload.time), }; } catch (err) { console.error("记录解析失败:", err, "原始数据:", record.Data.toString("base64")); return null; } }; // 处理批量数据 const processBatchData = (dataList) => { if (dataList.length === 0) return; console.log("批量处理数据:", JSON.stringify(dataList)); }; // 实时读取循环 const realTimeReadLoop = async () => { if (!currentShardIterator) return; try { const command = new GetRecordsCommand({ ShardIterator: currentShardIterator, Limit: CONFIG.batchMaxSize, WaitTimeSeconds: CONFIG.kinesisWaitTime, }); const response = await kinesis.send(command); currentShardIterator = response.NextShardIterator; if (!currentShardIterator) { console.log("无更多记录,终止读取"); return; } // 处理实时获取的记录 const validRecords = response.Records.map(parseRecord).filter(Boolean); validRecords.forEach(record => console.log("实时处理记录:", JSON.stringify(record))); // 继续循环 setTimeout(realTimeReadLoop, 0); } catch (err) { console.error("实时读取错误:", err); // 遇到错误重试 setTimeout(realTimeReadLoop, 1000); } }; // 批量读取循环 const batchReadLoop = async () => { if (!currentShardIterator) return; try { const command = new GetRecordsCommand({ ShardIterator: currentShardIterator, Limit: CONFIG.batchMaxSize, }); const response = await kinesis.send(command); currentShardIterator = response.NextShardIterator; if (!currentShardIterator) { console.log("无更多记录,终止读取"); // 处理剩余缓存 processBatchData(batchCache); return; } // 将有效记录加入缓存 const validRecords = response.Records.map(parseRecord).filter(Boolean); batchCache = [...batchCache, ...validRecords]; // 检查是否达到批量处理条件 if (batchCache.length >= CONFIG.batchMinSize) { processBatchData(batchCache.splice(0, CONFIG.batchMaxSize)); } // 继续循环 setTimeout(batchReadLoop, CONFIG.batchInterval); } catch (err) { console.error("批量读取错误:", err); setTimeout(batchReadLoop, CONFIG.batchInterval); } }; // 初始化分片迭代器 const initShardIterator = async () => { try { const command = new GetShardIteratorCommand({ StreamName: CONFIG.streamName, ShardId: CONFIG.shardId, ShardIteratorType: "LATEST", }); const response = await kinesis.send(command); currentShardIterator = response.ShardIterator; if (!currentShardIterator) { throw new Error("无法获取分片迭代器"); } } catch (err) { console.error("初始化分片迭代器失败:", err); throw err; } }; // 启动读取流程 const startReading = async () => { await initShardIterator(); if (CONFIG.readMode === "real-time") { console.log("启动实时读取模式"); realTimeReadLoop(); } else if (CONFIG.readMode === "batch") { console.log("启动定时批量读取模式"); // 启动定时器,到时间处理缓存 batchTimer = setInterval(() => { if (batchCache.length > 0) { processBatchData(batchCache.splice(0, CONFIG.batchMaxSize)); } }, CONFIG.batchInterval); batchReadLoop(); } else { throw new Error("不支持的读取模式,请配置real-time或batch"); } }; // 优雅退出处理 process.on("SIGINT", () => { console.log("正在终止读取流程..."); if (batchTimer) clearInterval(batchTimer); process.exit(0); }); startReading();
代码说明
- 模式切换:修改
CONFIG.readMode即可在实时/批量模式间切换 - 实时模式:通过
WaitTimeSeconds实现长轮询,数据到达立即处理,同时减少无效API调用 - 批量模式:缓存数据,满足条数或时间条件时批量输出,同时控制每次拉取的最大条数
- 错误处理:添加简单的重试逻辑,解析错误时保留原始数据便于排查
- 优雅退出:捕获中断信号,清理定时器后安全退出
内容的提问来源于stack exchange,提问作者anonymous_33008899
相关产品推荐
相关产品推荐

