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

如何用Node.js/CDK向Direct PUT模式的Kinesis Firehose投递数据?

问题解答

1. 能否用Node.js/CDK调用Kinesis Firehose的PutRecord/PutRecords API?

完全可以,AWS官方提供的JavaScript SDK(v2或v3版本)原生支持调用Firehose的API;在CDK场景下,既可以通过Step Functions的CallAwsService直接调用,也能在Lambda中嵌入SDK代码实现。

Node.js SDK v3 示例代码

先安装依赖包@aws-sdk/client-firehose,然后用以下代码调用:

import { FirehoseClient, PutRecordCommand, PutRecordsCommand } from "@aws-sdk/client-firehose";

// 初始化Firehose客户端
const firehoseClient = new FirehoseClient({ region: "你的AWS区域" });

// 发送单条记录
const sendSingleRecord = async () => {
  const params = {
    DeliveryStreamName: "test-firehose-delivery-stream",
    Record: {
      Data: Buffer.from(JSON.stringify({ event: "test-event" })) // JSON对象转Buffer,SDK自动完成Base64编码
    }
  };
  const command = new PutRecordCommand(params);
  const response = await firehoseClient.send(command);
  console.log("单条记录发送结果:", response);
};

// 发送多条记录
const sendMultipleRecords = async () => {
  const params = {
    DeliveryStreamName: "test-firehose-delivery-stream",
    Records: [
      { Data: Buffer.from(JSON.stringify({ event: "event-1" })) },
      { Data: Buffer.from(JSON.stringify({ event: "event-2" })) }
    ]
  };
  const command = new PutRecordsCommand(params);
  const response = await firehoseClient.send(command);
  console.log("多条记录发送结果:", response);
};

2. 将Step Functions的S3 PutObject替换为Firehose PutRecords

直接修改CallAwsService的配置即可,同时要注意Firehose的数据格式要求和权限配置:

修改后的Step Functions代码

const putEventToFirehose = new CallAwsService(this, 'Firehose Put Records Step', {
  service: 'firehose',
  action: 'putRecords',
  parameters: {
    DeliveryStreamName: firehoseeDeliveryStream.ref, // 直接引用CDK创建的投递流实例
    Records: JsonPath.array(
      JsonPath.object({
        // Firehose要求Data是Base64编码的字符串,所以先把JSON对象转字符串再编码
        Data: JsonPath.base64Encode(JsonPath.stringify(JsonPath.objectAt('$.S3Event')))
      })
    )
  },
  iamResources: [firehoseeDeliveryStream.attrArn],
});

// 替换原有的工作流链
const sendEventToFirehose = Chain.start(eventTransformer).next(putEventToFirehose);

补充权限与配置说明

  1. Step Functions角色权限:如果自动生成的权限不足,手动给Step Functions执行角色添加Firehose调用权限:
putEventToFirehose.addToRolePolicy(new iam.PolicyStatement({
  actions: ['firehose:PutRecords'],
  resources: [firehoseeDeliveryStream.attrArn],
}));
  1. Firehose投递流角色权限:你当前创建的firehoseRole需要具备写入目标S3桶的权限,补充以下代码:
firehoseRole.addToPolicy(new iam.PolicyStatement({
  actions: ['s3:PutObject', 's3:PutObjectAcl'],
  resources: [`${eventLogBucket.bucketArn}/*`],
}));
  1. 多条记录处理:如果需要批量发送多条事件,用JsonPath.map遍历事件数组生成多条Record:
Records: JsonPath.map(
  JsonPath.objectAt('$.events'),
  JsonPath.object({
    Data: JsonPath.base64Encode(JsonPath.stringify(JsonPath.current()))
  })
)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 02:25:35