如何用CDK将DynamoDB表设为开启增强扇出的Kinesis数据流消费者
如何用AWS CDK实现DynamoDB作为Kinesis Data Stream增强扇出的消费者
首先明确:DynamoDB本身不能直接作为Kinesis增强扇出的消费者——增强扇出模式要求消费者主动调用SubscribeToShard API拉取流数据,而DynamoDB没有内置这个能力。你必须通过一个中间计算层(比如Lambda函数)来充当增强扇出的消费者,再由这个中间层把数据写入DynamoDB表。
下面是完整的AWS CDK实现代码,基于你提供的片段扩展:
import * as cdk from 'aws-cdk-lib'; import { Construct } from 'constructs'; import * as kinesis from 'aws-cdk-lib/aws-kinesis'; import * as dynamodb from 'aws-cdk-lib/aws-dynamodb'; import * as lambda from 'aws-cdk-lib/aws-lambda'; import * as lambdaEventSources from 'aws-cdk-lib/aws-lambda-event-sources'; import { PolicyStatement } from 'aws-cdk-lib/aws-iam'; export class KinesisDynamoStack extends cdk.Stack { constructor(scope: Construct, id: string, props?: cdk.StackProps) { super(scope, id, props); // 创建Kinesis数据流 const eventStream = new kinesis.Stream(this, 'EventStream', { streamName: 'event-stream', shardCount: 1, }); // 创建目标DynamoDB表 const targetTable = new dynamodb.Table(this, 'EventLeaseTable', { tableName: `event-lease`, billingMode: dynamodb.BillingMode.PROVISIONED, readCapacity: 10, writeCapacity: 10, partitionKey: { name: 'leaseKey', type: dynamodb.AttributeType.STRING } }); // 创建增强扇出消费者 const streamConsumer = new kinesis.CfnStreamConsumer(this, 'StreamConsumer', { consumerName: `event-lease-consumer`, streamArn: eventStream.streamArn }); // 创建处理Kinesis数据并写入DynamoDB的Lambda函数 const kinesisToDynamoLambda = new lambda.Function(this, 'KinesisToDynamoLambda', { runtime: lambda.Runtime.NODEJS_20_X, handler: 'index.handler', code: lambda.Code.fromInline(` const AWS = require('aws-sdk'); const dynamodb = new AWS.DynamoDB.DocumentClient(); const TABLE_NAME = process.env.TABLE_NAME; exports.handler = async (event) => { // 批量处理Kinesis记录 const putRequests = event.Records.map(record => { const payload = Buffer.from(record.kinesis.data, 'base64').toString('utf-8'); const data = JSON.parse(payload); return { PutRequest: { Item: { leaseKey: data.leaseKey, // 确保和DynamoDB分区键匹配 ...data // 其他字段按需添加 } } }; }); // 批量写入DynamoDB if (putRequests.length > 0) { await dynamodb.batchWrite({ RequestItems: { [TABLE_NAME]: putRequests } }).promise(); } return { processedCount: putRequests.length }; }; `), environment: { TABLE_NAME: targetTable.tableName } }); // 给Lambda添加权限:订阅增强扇出消费者 kinesisToDynamoLambda.addToRolePolicy(new PolicyStatement({ resources: [streamConsumer.attrConsumerArn], actions: ['kinesis:SubscribeToShard'] })); // 给Lambda添加权限:写入DynamoDB表 targetTable.grantWriteData(kinesisToDynamoLambda); // 将Lambda关联到增强扇出消费者作为事件源 kinesisToDynamoLambda.addEventSource(new lambdaEventSources.KinesisEventSource(eventStream, { startingPosition: lambda.StartingPosition.LATEST, consumerArn: streamConsumer.attrConsumerArn, // 关键:指定增强扇出消费者ARN batchSize: 100, enabled: true })); } }
关键说明:
- 增强扇出关联:通过
KinesisEventSource的consumerArn参数,把Lambda和你创建的CfnStreamConsumer绑定,这样Lambda就会使用增强扇出模式拉取数据,而非普通的共享吞吐量模式。 - 权限配置:必须给Lambda授予
kinesis:SubscribeToShard权限(针对消费者ARN),同时授予DynamoDB的写入权限。 - 数据处理:Lambda负责解析Kinesis的Base64编码数据,转换成DynamoDB接受的格式,这里用批量写入来提升效率。
如果你不想用Lambda,也可以用ECS/EKS部署的自定义应用来作为增强扇出消费者,核心逻辑都是:调用SubscribeToShard拉取数据,然后写入DynamoDB。
内容的提问来源于stack exchange,提问作者CKT
相关产品推荐
相关产品推荐

