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

如何通过AWS CDK为已有DynamoDB表创建Kinesis流消费更新?

用AWS CDK实现现有DynamoDB表的Kinesis流消费方案

1. 先为现有DynamoDB表开启流功能

DynamoDB流确实需要开启才能捕获表的更新操作,针对现有表的情况,分两种场景处理:

场景1:表是通过CDK创建的

直接修改原CDK代码,给表添加stream配置并指定流视图类型,重新部署即可:

import { Table, StreamViewType } from 'aws-cdk-lib/aws-dynamodb';

// 原表定义代码,添加stream属性
const myTable = new Table(this, 'MyTable', {
  // 原有配置(表名、主键等)
  stream: StreamViewType.NEW_IMAGE, // 根据需求选择视图类型
});

场景2:表不是CDK创建的(纯现有资源)

你可以先通过AWS CLI开启流:

aws dynamodb update-table --table-name your-table-name --stream-specification StreamEnabled=true,StreamViewType=NEW_IMAGE

或者在CDK中通过CfnTable导入并修改配置:

import { CfnTable, StreamViewType } from 'aws-cdk-lib/aws-dynamodb';

const existingTable = CfnTable.fromCfnTableAttributes(this, 'ExistingTable', {
  tableArn: 'arn:aws:dynamodb:your-region:your-account-id:table/your-table-name',
});

existingTable.streamSpecification = {
  streamViewType: StreamViewType.NEW_IMAGE,
};

注意:用L2的Table.fromTableName/fromTableArn导入的资源是只读的,无法直接修改流配置,必须用CfnTable方式。

2. 用CDK创建Kinesis流并关联DynamoDB流消费

DynamoDB流无法直接对接Kinesis流,通常用Lambda作为中间转发器搭建链路,以下是完整实现代码:

完整代码示例(TypeScript)

import { Stack, StackProps } from 'aws-cdk-lib';
import { Construct } from 'constructs';
import { Table } from 'aws-cdk-lib/aws-dynamodb';
import { Stream } from 'aws-cdk-lib/aws-kinesis';
import { Function, Runtime, StartingPosition } from 'aws-cdk-lib/aws-lambda';
import { DynamoEventSource } from 'aws-cdk-lib/aws-lambda-event-sources';

export class DynamoToKinesisStack extends Stack {
  constructor(scope: Construct, id: string, props?: StackProps) {
    super(scope, id, props);

    // 1. 导入已开启流的现有DynamoDB表
    const existingTable = Table.fromTableName(this, 'ExistingTable', 'your-table-name');

    // 2. 创建目标Kinesis流
    const kinesisStream = new Stream(this, 'DynamoUpdatesStream', {
      streamName: 'dynamodb-table-updates',
      shardCount: 1, // 根据表的吞吐量调整分片数
    });

    // 3. 创建Lambda转发函数:读取DynamoDB流记录并写入Kinesis流
    const forwarderLambda = new Function(this, 'DynamoToKinesisForwarder', {
      runtime: Runtime.NODEJS_18_X,
      handler: 'index.handler',
      code: Function.fromInline(`
        const AWS = require('aws-sdk');
        const kinesis = new AWS.Kinesis();

        exports.handler = async (event) => {
          // 转换DynamoDB流记录为Kinesis可接受的格式
          const kinesisRecords = event.Records.map(record => ({
            Data: JSON.stringify(record.dynamodb),
            PartitionKey: record.dynamodb.Keys.id.S // 建议用表的主键作为分区键,避免热点
          }));

          await kinesis.putRecords({
            Records: kinesisRecords,
            StreamName: process.env.KINESIS_STREAM_NAME
          }).promise();

          return { processedCount: kinesisRecords.length };
        };
      `),
      environment: {
        KINESIS_STREAM_NAME: kinesisStream.streamName,
      },
    });

    // 4. 给Lambda添加DynamoDB流触发源
    forwarderLambda.addEventSource(new DynamoEventSource(existingTable, {
      startingPosition: StartingPosition.LATEST, // 从最新记录开始消费
      batchSize: 100, // 每次批量处理的记录数
      enabled: true,
    }));

    // 5. 授予Lambda写入Kinesis流的权限
    kinesisStream.grantWrite(forwarderLambda);
  }
}

关键配置说明

  • 流视图类型:StreamViewType有四种可选值,根据业务需求选择:
    • NEW_IMAGE:仅捕获更新后的新数据
    • OLD_IMAGE:仅捕获更新前的旧数据
    • NEW_AND_OLD_IMAGES:同时捕获新旧数据
    • KEYS_ONLY:仅捕获主键数据
  • Kinesis分片数:1个分片支持每秒1MB写入/2MB读取,需匹配DynamoDB表的读写吞吐量,避免限流。
  • Lambda批量大小:调整批量大小可平衡处理效率和延迟,批量越大,单次处理记录越多,但延迟可能越高。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 21:05:16