You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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();

代码说明

  1. 模式切换:修改CONFIG.readMode即可在实时/批量模式间切换
  2. 实时模式:通过WaitTimeSeconds实现长轮询,数据到达立即处理,同时减少无效API调用
  3. 批量模式:缓存数据,满足条数或时间条件时批量输出,同时控制每次拉取的最大条数
  4. 错误处理:添加简单的重试逻辑,解析错误时保留原始数据便于排查
  5. 优雅退出:捕获中断信号,清理定时器后安全退出

内容的提问来源于stack exchange,提问作者anonymous_33008899

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.22 18:24:54