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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 15:55:20