Lambda调用Kinesis PutRecord时频繁出现ETIMEDOUT错误(偶现正常,修改PartitionKey部署后有时恢复)
Lambda调用Kinesis PutRecord时频繁出现ETIMEDOUT错误(偶现正常,修改PartitionKey部署后有时恢复)
我最近写了个简单的Node.js Lambda,功能很明确——从手动触发的SQS里拿消息,然后调用Kinesis的PutRecord接口把消息推流。权限配置我都检查过了,PutRecord的权限是给足的,但大部分时候都会报超时错误,偶尔又能正常跑,尤其是我改个随机的PartitionKey重新部署之后,成功率会高一些,这实在搞不懂哪里出问题了。
错误日志
{ "errorType": "Runtime.UnhandledPromiseRejection", "errorMessage": "Error [ERR_HTTP2_STREAM_CANCEL]: The pending stream has been canceled (caused by: connect ETIMEDOUT 34.223.45.15:443)", "trace": [ "Runtime.UnhandledPromiseRejection: Error [ERR_HTTP2_STREAM_CANCEL]: The pending stream has been canceled (caused by: connect ETIMEDOUT 34.223.45.15:443)", " at process.<anonymous> (file:///var/runtime/index.mjs:1276:17)", " at process.emit (node:events:517:28)", " at emit (node:internal/process/promises:149:20)", " at processPromiseRejections (node:internal/process/promises:283:27)", " at process.processTicksAndRejections (node:internal/process/task_queues:96:32)" ] }
我的Lambda代码
import { KinesisClient, PutRecordCommand } from "@aws-sdk/client-kinesis"; export const KINESIS_CLIENT = new KinesisClient([ { httpOptions: { connectTimeout: 10000, }, maxRetries: 1, region: 'us-west-2', }, ]); export const handler = async (event) => { const promises = []; for (const {messageId, body} of event.Records) { processEvent(body) promises.push(processEvent(body)); } const responses = await Promise.allSettled(promises); responses.forEach((response) => { if (response.status !== "fulfilled") { throw Error(JSON.stringify(responses)); } }); }; export async function processEvent(body) { const newBody = JSON.parse(body); newBody['field'] = 'Random'; await KINESIS_CLIENT.send( new PutRecordCommand({ Data: new TextEncoder().encode(JSON.stringify(newBody)), StreamName: 'InputEventStream', PartitionKey: '3', // <--- 我改这个值重新部署后有时会正常 }), ); }
问题分析与修复方案
1. 核心问题推测
从错误ERR_HTTP2_STREAM_CANCEL和connect ETIMEDOUT来看,主要问题集中在网络连接+SDK配置,结合改PartitionKey生效的现象,也可能关联到Kinesis分片的网络路由问题:
- AWS SDK v3默认用HTTP2协议,HTTP2的连接复用机制在Lambda的冷热启动场景下可能出现失效连接未被清理的情况,导致超时
- 你设置的
maxRetries:1太低,遇到网络波动直接失败;仅配置connectTimeout不够,缺少覆盖整个请求生命周期的超时参数 - 改PartitionKey后路由到不同Kinesis分片,如果某分片所在AZ和Lambda的AZ网络有问题,换分片就临时恢复,这指向VPC/跨AZ网络可能存在异常
2. 具体修复步骤
(1)优化KinesisClient配置
调整超时、重试参数,甚至暂时禁用HTTP2改用HTTP1.1(避免HTTP2连接复用的坑):
import { KinesisClient, PutRecordCommand } from "@aws-sdk/client-kinesis"; import { NodeHttpHandler } from "@aws-sdk/node-http-handler"; // 改用更稳定的HTTP1.1,同时补全超时与重试配置 export const KINESIS_CLIENT = new KinesisClient({ region: 'us-west-2', maxRetries: 3, // 提高重试次数应对网络波动 retryMode: "adaptive", // 启用自适应重试,AWS SDK会根据错误类型自动调整策略 requestHandler: new NodeHttpHandler({ connectionTimeout: 10000, socketTimeout: 30000, // 覆盖整个请求的超时时间,不止是连接阶段 http2: false, // 禁用HTTP2,改用HTTP1.1 }), });
(2)修复代码里的重复调用bug
你的handler里processEvent(body)被调用了两次,会导致同一条消息重复发送,还会浪费连接资源,去掉多余的调用:
export const handler = async (event) => { const promises = []; for (const {messageId, body} of event.Records) { // 删掉这行重复调用 // processEvent(body) promises.push(processEvent(body)); } const responses = await Promise.allSettled(promises); responses.forEach((response) => { if (response.status !== "fulfilled") { // 只抛出失败的具体原因,方便排查 throw Error(`Send failed: ${JSON.stringify(response.reason)}`); } }); };
(3)检查VPC/网络配置(如果Lambda在VPC内)
- 确保Lambda所在安全组允许出站访问Kinesis的443端口
- 如果用VPC端点访问Kinesis,检查端点的安全组是否允许Lambda的访问,且端点状态正常
- 排查NAT网关(如果用了)的连接数是否耗尽,或者有没有带宽瓶颈
(4)增加日志排查分片问题
在processEvent里加日志,记录PartitionKey和对应的分片ID,确认是不是特定分片的网络问题:
export async function processEvent(body) { const newBody = JSON.parse(body); newBody['field'] = 'Random'; const partitionKey = '3'; // 可以动态生成或者传参 try { const response = await KINESIS_CLIENT.send( new PutRecordCommand({ Data: new TextEncoder().encode(JSON.stringify(newBody)), StreamName: 'InputEventStream', PartitionKey: partitionKey, }), ); console.log(`消息发送成功: 分片ID=${response.ShardId}, PartitionKey=${partitionKey}`); return response; } catch (e) { console.error(`消息发送失败: PartitionKey=${partitionKey}, 错误=${e.message}`); throw e; } }
备注:内容来源于stack exchange,提问作者Saad
相关产品推荐
相关产品推荐

