写入大量记录到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

