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

写入大量记录到Kinesis时触发ProvisionedThroughputExceededException的排查

问题描述

我正在用Lambda结合Kinesis做批量记录的验证与插入,作为Kinesis新手,搞不清是不是配置错了,也觉得简单场景下它太复杂。

我把一个27.5万行的大文件拆成275个各含1000行的小文件,这些文件在30秒内上传到S3后触发生产者Lambda,生产者向Kinesis写记录,再触发消费者Lambda以10条为批量写入PostgreSQL。

测试100、1000甚至2万条记录都正常,但处理27.5万条时出现错误:

ProvisionedThroughputExceededException: Rate exceeded for shard

最终只成功处理了约13万条记录。每个小文件都小于1MB,包含1000条JSON记录,每条有40个字段,数据量不大。

我推测是消费者触发时的读取请求超出了流的处理能力,是不是要启用Enhanced Fan-Out?另外调整事件源的批处理大小会出现重复记录的问题。

基础设施代码(CDK)

Kinesis 流配置

export class StreamStack extends Stack {
    public readonly myStream: Stream;

    constructor(scope: Construct, id: string, props?: StackProps) {
        super(scope, id, props);

        this.myStream = new Stream(this, "MyStream", {
            streamName: "my-stream",
            streamMode: StreamMode.ON_DEMAND
        });
    }
}

事件源映射配置

interface EventSourceStackProps extends StackProps {
    myStream: Stream,
    myConsumerLambda: NodejsFunction
}

export class EventSourceStack extends Stack {
    constructor(scope: Construct, id: string, props: EventSourceStackProps) {
        super(scope, id, props);

        new EventSourceMapping(this, "ConsumerFunctionEvent", {
            target: props.myConsumerLambda,
            batchSize: 10,
            startingPosition: StartingPosition.LATEST,
            eventSourceArn: props.myStream.streamArn
        });
        props.myStream.grantRead(props.myConsumerLambda);
    }
}

Lambda 函数配置

// 生产者Lambda
this.myProducerLambda = new NodejsFunction(this, "MyProducerLambda", {
    functionName: "myProducerHandler",
    runtime: Runtime.NODEJS_18_X,
    handler: "myProducerHandler",
    entry: "./lambda/myProducerHandler.ts",
    memorySize: 256,
    timeout: Duration.minutes(5)
});
props.myStream.grantWrite(this.myProducerLambda);

// 消费者Lambda
this.myConsumerLambda = new NodejsFunction(this, "MyConsumerLambda", {
    functionName: "myConsumerHandler",
    runtime: Runtime.NODEJS_18_X,
    handler: "myConsumerHandler",
    entry: "./lambda/myConsumerHandler.ts",
    memorySize: 1024,
    timeout: Duration.minutes(15)
});

业务代码

生产者发送记录

export const putStreamRecord = async (item: StreamItem): Promise<void> => {
    try {
        const client = new KinesisClient();
        const params = {
            Data: Buffer.from(JSON.stringify(item)),
            PartitionKey: item.key,
            StreamName: process.env.STREAM_NAME,
        };
        await client.send(new PutRecordCommand(params));
    } catch (e) {
        console.error(e);
        throw e;
    }
};

消费者处理逻辑

export const myConsumerHandler = async (event: KinesisStreamEvent): Promise<boolean> => {
    let batchSaveRes = false;
    try {
        // 提取批量记录
        const batchData = event.Records
            .map(rec => Buffer.from(rec.kinesis.data, "base64").toString())
            .map(buff => JSON.parse(buff))
            .map(item => item.payload);
        
        batchSaveRes = await consumeBatchStream(batchData);
        return batchSaveRes;
    } catch (e) {
        console.error("ERROR:", e);
        return batchSaveRes;
    }
};

PostgreSQL批量插入

export const consumeBatchStream = async (batchData: Record<string, unknown>[]): Promise<boolean> => {
    let client: Client | null = null;
    try {
        client = await getDbClient("my-secret");
        await client?.connect();
        const query = // 构建插入SQL语句
        await client.query(query);
        console.log("INSERTED:", batchData);
        return true;
    } catch (e) {
        console.error("Error consuming batch stream:", e)
        throw e;
    } finally {
        await client?.end();
    }
};

额外疑问

我只需要一种简单的方式通知多个消费者特定类型的新记录,不关心顺序,SQS/SNS能不能替代Kinesis解决当前的限制问题?之前用Kafka处理类似场景简单多了。


解决方案

一、Kinesis 当前问题分析与优化

1. 吞吐量超限原因

你用的是On-Demand模式的Kinesis流,它的默认吞吐量是每个分片每秒2MB写入/读取,或者每秒1000条写入/读取。但你的生产者在30秒内要写入27.5万条记录,平均每秒近9000条,远超单个分片的处理能力——On-Demand模式会自动扩容分片,但扩容有延迟,短时间内突发的高写入量会导致分片跟不上,进而触发ProvisionedThroughputExceededException。

另外,消费者的批处理大小设为10,会导致Lambda频繁调用,每个调用只处理10条记录,增加了读取请求的频次,进一步加剧了分片的读取压力。

2. 优化方案

(1)优化生产者写入方式

不要用PutRecord逐条写入,改用**PutRecords批量写入**,每次批量提交500条记录(Kinesis的批量上限),这样能大幅减少API调用次数,降低写入压力,同时提升效率。示例代码:

export const putStreamRecords = async (items: StreamItem[]): Promise<void> => {
    try {
        const client = new KinesisClient();
        const records = items.map(item => ({
            Data: Buffer.from(JSON.stringify(item)),
            PartitionKey: item.key
        }));
        const params = {
            Records: records,
            StreamName: process.env.STREAM_NAME,
        };
        const response = await client.send(new PutRecordsCommand(params));
        // 处理失败的记录(如果有)
        const failedRecords = response.FailedRecordCount || 0;
        if (failedRecords > 0) {
            console.warn(`Failed to write ${failedRecords} records`);
            // 可实现重试逻辑
        }
    } catch (e) {
        console.error(e);
        throw e;
    }
};

(2)调整消费者批处理大小

把batchSize从10调大到100-500(根据你的PostgreSQL插入性能调整),这样单个Lambda调用能处理更多记录,减少读取请求频次,降低分片压力。

关于重复记录问题:Kinesis本身就可能出现至少一次交付,所以你的消费者必须实现幂等性——比如在PostgreSQL表中给每条记录的唯一键加唯一约束,插入时用ON CONFLICT DO NOTHING或者ON CONFLICT UPDATE,这样即使重复处理也不会导致数据重复。

(3)启用Enhanced Fan-Out(可选)

如果你的消费者需要独立的读取吞吐量(比如多个消费者同时读取),启用Enhanced Fan-Out能给每个消费者分配独立的分片读取配额,不会和其他消费者抢占吞吐量。但如果只有一个消费者,这个优化的收益不大,优先调整批处理大小和生产者写入方式即可。

二、SQS/SNS 替代方案分析

如果你的需求是通知多个消费者、不关心顺序,SQS/SNS完全可以替代Kinesis,而且更简单:

1. 架构选择

  • 如果是多个消费者需要接收相同的记录:用SNS主题+多个SQS队列订阅,每个消费者对应一个SQS队列,SNS把消息推送到所有订阅的队列。
  • 如果是单个消费者或多个消费者分摊处理:直接用SQS标准队列,多个消费者可以同时拉取消息。

2. 优势

  • 无需关心分片、吞吐量配置,SQS会自动扩容,处理突发流量更省心。
  • 消息重试、死信队列(DLQ)的配置更简单,能轻松处理失败的消息。
  • 成本更低:SQS按请求次数和存储量收费,对于批量处理场景,比Kinesis更划算。

3. 注意事项

  • SQS标准队列可能会出现重复消息,同样需要实现幂等性。
  • 如果需要保留消息超过14天(SQS标准队列的最大保留期),那Kinesis更合适,但你的场景应该不需要。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 06:14:58