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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.15 09:08:03