如何用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);
补充权限与配置说明
- Step Functions角色权限:如果自动生成的权限不足,手动给Step Functions执行角色添加Firehose调用权限:
putEventToFirehose.addToRolePolicy(new iam.PolicyStatement({ actions: ['firehose:PutRecords'], resources: [firehoseeDeliveryStream.attrArn], }));
- Firehose投递流角色权限:你当前创建的
firehoseRole需要具备写入目标S3桶的权限,补充以下代码:
firehoseRole.addToPolicy(new iam.PolicyStatement({ actions: ['s3:PutObject', 's3:PutObjectAcl'], resources: [`${eventLogBucket.bucketArn}/*`], }));
- 多条记录处理:如果需要批量发送多条事件,用
JsonPath.map遍历事件数组生成多条Record:
Records: JsonPath.map( JsonPath.objectAt('$.events'), JsonPath.object({ Data: JsonPath.base64Encode(JsonPath.stringify(JsonPath.current())) }) )
内容的提问来源于stack exchange,提问作者gillJ
相关产品推荐
相关产品推荐

