原文发布于 Medium,2023 年 11 月 11 日。

在本指南中,我们将讨论使用 AWS Lambda 和 NestJS 在分布式系统中实现通信的方法。

基于我们之前的探索,我们将 NestJS 与 AWS Lambda 集成,并在 Docker 中设置调试和热重载,本指南将迈出下一步。我们将使用相同的基础设置,并对其进行增强,以实现 Lambda 之间的无缝通信。

本指南中的方法灵感来源于我过去使用 Dapr 与 C# 微服务的一个项目。Dapr 简化了构建微服务的挑战。例如,Dapr 提供了简单的方法来设置服务间通信,而不是手动设置。这让开发者能够专注于编写实际功能,而不必陷入服务间通信的技术细节。这种 Lambda + NestJS 通信策略深受 Dapr 的简洁性和有效性的启发。

如果你对了解更多关于 Dapr 和 C# 的内容感兴趣,请查看我的 其他仓库

本指南重点介绍同步请求,但我稍后会在其他指南中介绍使用 SQS 的异步请求。

你可以在此仓库中找到演示应用https://github.com/VenyaBrodetskiy/Lambda-NestJS-Demo

该演示包括:

  • 一个按照上一个指南配置为在本地运行的应用。
  • 实现 Lambda 之间的通信,使其既能在本地运行,也能在部署后运行(本指南的重点)
  • 如何在解决方案中添加队列(SQS)

在接下来的部分中,我们将分解创建 Lambda 之间通信所涉及的步骤:

  1. 了解简化 Lambda 通信的服务设计
  2. 构建 Lambda 通信服务
  3. 通过重试增强 Lambda 通信服务
  4. 实现 Lambda Factory 以支持本地和云环境

让我们开始吧!

第一部分。了解简化 Lambda 通信的服务设计

概念:目标是创建一个服务,封装复杂的逻辑,使其像 Dapr 一样用户友好和简单。本质上,此服务充当桥梁,简化我们系统中不同组件之间的通信。

实际示例:让我们通过一个简单的 NestJS 控制器来说明这个服务应该如何工作。此控制器处理外部请求(例如,来自前端应用或 Postman),并与另一个 Lambda 函数通信以从数据库中获取数据。

以下是我们的 PlanController 的代码:

import { LambdaCommunicationService } from 'src/core/modules/communication';

@Controller('plan')
export class PlanController {
  private readonly logger = new Logger(PlanController.name);
  constructor(
    private readonly lambdaService: LambdaCommunicationService,
  ) {}

  @Get('/:id')
  public async getPlanById(@Param('id') id: string): Promise<PlanRes> {
    this.logger.log(`Inside ${this.getPlanById.name}, id: ${id}`);

    // service to call other lambda
    const result: PlanRes = await this.lambdaService.invoke<PlanRes>(
      Accessor.Plan,          // name of lambda to be called. `Accessor` is an enum representing different lambda functions
      `/planaccessor/${id}`,  // path of request
      HttpMethod.Get,         // the HTTP method for the request
    );
    return result;
  }

Enter fullscreen mode Exit fullscreen mode

lambdaService.invoke 方法的签名:

LambdaCommunicationService.invoke<TResponse>(
  service: string, 
  path: string, 
  httpMethod?: HttpMethod, 
  payload?: object): Promise<TResponse>

Enter fullscreen mode Exit fullscreen mode

  • service: 标识要调用的 Lambda 函数的名称。
  • path: 指定请求的路径。
  • httpMethod: 请求的 HTTP 方法,如果未指定则默认为 GET。
  • payload: 可选对象,包含要发送到 Lambda 的数据。

优势:通过这样的服务,调用其他 Lambda 变得简单明了。你只需要函数名称、调用路径、方法以及必要时的任何有效负载。

这种设计显著减少了开发时间,因为不需要编写和维护用于 Lambda 交互的复杂代码。它还通过标准化服务通信的方式来最小化错误,确保一致性和可靠性。此外,这样的服务提高了代码的可读性和可维护性,使团队更容易理解和修改代码库。

关于错误处理的说明: 你可能会注意到缺少 try-catch 块。这是因为 NestJS 具有内置的异常过滤器,可以优雅地处理错误。你可以在 NestJS 异常过滤器文档中了解有关此功能的更多信息。此外,我的 Lambda-NestJS Demo 仓库 包含异常过滤器的自定义实现,尽管这超出了本指南的范围。

第二部分。构建 Lambda 通信服务

在本部分中,让我们实现第一部分中描述的服务:

@Injectable()
export class LambdaCommunicationService {
  private readonly logger = new Logger(LambdaCommunicationService.name);
  constructor(private lambdaFactory: LambdaFactory) {}

  public async invoke<TResponse>(
    service: string,
    path: string,
    httpMethod: HttpMethod = HttpMethod.Get,
    payload?: object,
  ): Promise<TResponse> {
    try {
      // get instance of lambda object
      const { lambda, functionName } = this.lambdaFactory.getLambda(service);

      // prepare payload
      const lambdaPayload: ICommunicationPayload = {
        httpMethod: httpMethod,
        path: path,
        body: payload ?? undefined,
        headers: {
          'Content-Type': 'application/json',
        },
      };

      const params: InvokeCommandInput = {
        FunctionName: functionName,
        Payload: JSON.stringify(lambdaPayload),
      };

      this.logger.debug(
        `Inside ${this.invoke.name}. Invoking function: ${params.FunctionName} with payload: ${params.Payload}`,
      );

      // call other lambda using aws-sdk
      const response: InvokeCommandOutput = await lambda.invoke(params);

      // handle the response
      const responsePayload = JSON.parse(Buffer.from(response.Payload).toString());
      if (RetriableStatusCodes.includes(responsePayload.statusCode)) {
        throw new CommunicationException(
          JSON.parse(responsePayload.body),
          responsePayload.statusCode,
        );
      }

      this.logger.debug(
        `Inside ${this.invoke.name}. Lambda invoke response payload: ${JSON.stringify(
          responsePayload,
          null,
          ' ',
        )}`,
      );

      // parse the response to expected type
      if (typeof responsePayload.body === 'string' && this.isJsonString(responsePayload.body)) {
        const result = JSON.parse(responsePayload.body) as TResponse;
        return result;
      }

      return responsePayload.body as TResponse;
    } catch (error: any) {
      if (error.code === 'ECONNREFUSED')
        throw new CommunicationException(
          `Failed to invoke lambda: ${service}`,
          HttpStatus.INTERNAL_SERVER_ERROR,
        );
      throw error;
    }
  }

  private isJsonString(str: string): boolean {
    try {
      JSON.parse(str);
      return true;
    } catch (e) {
      return false;
    }
  }
}

Enter fullscreen mode Exit fullscreen mode

现在我将分解上面的代码并逐一解释。

  1. 对于第一步,我们使用 AWS SDK 来实例化一个 Lambda 对象。此对象负责调用其他 Lambda 函数。除此之外,我们还检索要调用的特定函数名称。此 Lambda 对象的创建和管理由 LambdaFactory 高效处理。了解 LambdaFactory 的工作原理是我们实现的关键,我们将在后续部分(第 4 部分)中更详细地探讨它:
const { lambda, functionName } = this.lambdaFactory.getLambda(service);

Enter fullscreen mode Exit fullscreen mode

2. 下一步是构建 Lambda Payload 和调用参数,调用 Lambda:

const lambdaPayload: ICommunicationPayload = {
  ...
};

const params: InvokeCommandInput = {
  ...
};

const response: InvokeCommandOutput = await lambda.invoke(params);

Enter fullscreen mode Exit fullscreen mode

3. 处理响应

这是一个关键步骤。我们需要了解 responsePayload 中的内部状态码,以确定调用是否成功:

const responsePayload = JSON.parse(Buffer.from(response.Payload).toString());

if (RetriableStatusCodes.includes(responsePayload.statusCode)) {
  throw new CommunicationException(
    ...
  );
}

Enter fullscreen mode Exit fullscreen mode

4. 解析响应

来自 Lambda 的响应可以是 TResponse 类型的对象(由开发者定义)或字符串。以下是我们处理它的方式:

if (typeof responsePayload.body === 'string' && this.isJsonString(responsePayload.body)) {
  const result = JSON.parse(responsePayload.body) as TResponse;
  return result;
}

return responsePayload.body as TResponse;

Enter fullscreen mode Exit fullscreen mode

这部分代码检查响应中 body 的类型并相应地解析它。TypeScript 的使用允许我们将响应转换为开发者指定的 TResponse 类型。

第三部分。通过重试增强 Lambda 通信服务

在本部分中,我们将向 Lambda 通信服务添加重试功能。重试失败请求的能力是一个关键特性,特别是用于处理瞬态网络问题或临时服务不可用。

注意。但是,由于其幂等性特性,必须仔细选择 HTTP 方法。GET 和 DELETE 主要针对重试,因为它们不会以导致重复时产生副作用的方式更改状态。PUT 和 PATCH 也可以考虑用于重试,因为它们被设计为幂等的,确保重复请求导致相同的状态。但是,建议对 POST 请求谨慎,因为它们通常修改状态或创建资源,如果没有适当的幂等性处理而重试,可能会导致意外后果。

应针对指示瞬态问题或服务器错误的重试错误代码应用重试,在这些情况下,重复请求可能会成功。通常,5xx 系列错误(例如,500 Internal Server Error、502 Bad Gateway、503 Service Unavailable、504 Gateway Timeout)是适合重试的候选者,表明服务器端存在临时问题,但不建议对 501 Not Implemented 进行重试,因为此错误表明服务器存在永久性限制,并且后续重试不太可能成功。至于 4xx 系列错误(例如,400 Bad Request、401 Unauthorized、404 Not Found),它们通常表明客户端问题,不太可能通过在不更改请求的情况下重试来解决。但是,瞬态错误如 408 Request Timeout、423 Locked 和 429 Too Many Requests 可能是例外,在这些情况下,使用退避策略重试可能会成功。

为了集成此功能,我们将使用 async-retry npm 包,它提供了一种实现重试逻辑的简单方法。Lambda 通信服务中的 invoke 方法增强如下:

  1. 添加重试参数:invoke 方法现在包括一个额外的参数 retries,默认值为 3。此参数确定请求的最大重试次数。
public async invoke<TResponse>(
    service: string,
    path: string,
    httpMethod: HttpMethod = HttpMethod.Get,
    payload?: object,
    retries: number = 3,  // new parameter
  ): Promise<TResponse> {

Enter fullscreen mode Exit fullscreen mode

2. 条件重试逻辑:我们引入一个检查,仅对幂等请求启用重试。对于其他 HTTP 方法,有效重试次数设置为零。

// enable retries only for Retriable Http Methods requests
let effectiveRetries;
if (RetriableHttpMethods.includes(httpMethod)) {
  effectiveRetries = retries;
} else {
  effectiveRetries = 0;
}

Enter fullscreen mode Exit fullscreen mode

3. 重试机制:使用 async-retry 包中的 retry 函数,我们包装实际的 Lambda 调用逻辑。此函数将根据提供的重试条件自动重试调用。

let responsePayload;
await retry(
  async () => {
    this.logger.debug(
      `Inside ${this.invoke.name}. Invoking function: ${params.FunctionName} with payload: ${params.Payload}`,
    );

    const response = await lambda.invoke(params);
    responsePayload = JSON.parse(Buffer.from(response.Payload).toString());

    if (RetriableStatusCodes.includes(responsePayload.statusCode)) {
      throw new CommunicationException(
        JSON.parse(responsePayload.body),
        responsePayload.statusCode,
      );
    }
  },
  // retry configuration
  {
    retries: effectiveRetries,
    onRetry: (error) => {
      error &&
        this.logger.warn(`Error while calling lambda: ${params.FunctionName} with payload: ${
          params.Payload
        }, retrying...
        Error: ${JSON.stringify(error, null, ' ')}`);
    },
  },
);

Enter fullscreen mode Exit fullscreen mode

重试期间的错误处理:如果在调用尝试期间发生错误,将生成错误日志。这有助于监控和调试与失败的 Lambda 调用相关的问题。

通过实现此重试机制,我们确保我们的 Lambda 通信服务更加健壮,并且可以更优雅地处理瞬态故障,从而提高系统的整体可靠性。

第四部分。实现 Lambda Factory 以支持本地和云环境

在上一节中,我们在 Lambda 通信服务中引入了 LambdaFactory。

const { lambda, functionName } = this.lambdaFactory.getLambda(service);

Enter fullscreen mode Exit fullscreen mode

现在,让我们探讨为什么它很重要以及如何实现它以确保在本地和云中顺利运行。

LambdaFactory 的作用是抽象 AWS Lambda 实例的创建和配置。这种抽象至关重要,因为它允许我们的应用程序动态适应不同的环境——本地开发或云部署——而无需更改核心业务逻辑。

1. 了解 LambdaFactory 的实现

让我们检查 LambdaFactory 服务的关键组件。

lambda-factory.service.ts:

...
import { Configuration } from 'src/config';

interface ILambdaClient {
  lambda: Lambda;
  functionName: string;
}

@Injectable()
export class LambdaFactory {
  private lambdaInstance: Lambda;

  constructor(private config: Configuration) {}

  public getLambda(service: string): ILambdaClient {
    // get function name and endpoint from configuration
    const { name: functionName, endpoint: endpoint } = this.config.getService(service);

    // for cloud
    if (!this.config.IsOffline) {
      this.lambdaInstance = this.lambdaInstance ?? new Lambda({});

      return {
        lambda: this.lambdaInstance,
        functionName: functionName,
      };
    }

    // for local development
    this.lambdaInstance = new Lambda({
      endpoint: endpoint,
    });

    return {
      lambda: this.lambdaInstance,
      functionName: functionName,
    };
  }
}

Enter fullscreen mode Exit fullscreen mode

LambdaFactory 服务中的 getLambda 方法是根据 Configuration 服务确定的,为本地开发或云部署配置 Lambda 客户端实例的关键。这至关重要,原因有两个:

  1. 本地开发:对于本地测试,特别是使用 Serverless 框架,每个 Lambda 函数通常需要一个唯一的端点。因此,该方法为每次调用创建一个新的 Lambda 对象,确保准确的本地模拟。
  2. 云部署:在云中,Lambda 函数由其名称标识,而不是端点。在这里,该方法通过为每次调用重用相同的 Lambda 实例来优化性能,遵循单例模式。

这种方法确保了不同环境之间的灵活性和一致性,简化了开发和部署过程。

2. 了解 Configuration 服务

在 LambdaFactory 中,我们使用 Configuration 服务的 getService 方法来获取 LambdaCommunicationService 的特定设置,基于定义的枚举。

configuration.service.ts:

@Injectable()
export class Configuration {
  constructor(private configService: ConfigService) {
    this.validateConfig();
  }

  get IsOffline(): boolean {
    return Boolean(this.configService.get<boolean>('IS_OFFLINE'));
  }

  public getService(service: string): IAcccessorConfig {
    try {
      switch (service) {
        case Accessor.Plan:
          return {
            name: this.configService.getOrThrow<string>('PLANACCESSOR_NAME'),
            endpoint: this.configService.getOrThrow<string>('PLANACCESSOR_ENDPOINT'),
          };
        // case Accessor.OtherAccessor:
        // ...
        default:
          throw new Error(
            `Unknown accessor type. Configuration.service misses accessor: ${service}`,
          );
      }
    } catch (e) {
      throw new Error(e);
    }
  }
  private validateConfig(): void {
    // when running application, this function checks that developer didn't forget to add necessary configs to configuration.service (mostly for enums)

    // Validate accessor configurations
    for (const accessor of Object.values(Accessor)) {
      this.getService(accessor as Accessor);
    }

    // Validate queue configuration not missed
    for (const queue of Object.values(Queue)) {
      this.getQueue(queue as Queue);
    }
  }
}

Enter fullscreen mode Exit fullscreen mode

你需要执行一些操作才能从使用 Configuration 服务中获益:

  • 定义枚举:开发者需要为不同的服务创建枚举,以简化配置检索。
  • 维护 .env 文件:他们必须使用所有必要的变量设置和更新 .env 文件。

错误处理和验证:该服务在 getService 方法中包含错误处理,并包含 validateConfig 方法,以确保所有配置都正确且完整。

3. 使用 LambdaFactory 的主要优势

总之,LambdaFactory 是一种强大的模式,它简化了特定于环境配置的管理,并增强了我们 Lambda 通信服务的健壮性,使其能够适应本地和云环境。让我们总结一下这种方法的优势:

  1. 与环境无关:LambdaFactory 在本地和云配置之间无缝切换,提高了开发人员的生产力并减少了特定于环境的错误。
  2. 配置灵活性:通过集中 Lambda 配置,它允许在不改变核心逻辑的情况下轻松更新和维护服务配置。
  3. 增强的可读性:清晰的关注点分离和特定于环境细节的抽象使代码更具可读性和可维护性。

总结

当我们完成本指南关于使用 NestJS 改进 AWS Lambda 通信时,我们已经探索了创建既适用于本地开发又适用于云中使用的通信系统的细节。

从了解服务设计和构建健壮的 Lambda 通信服务到实现用于环境适应性的 Lambda Factory,本指南提供了一条全面的途径,以掌握分布式系统中的 Lambda 通信。

🔗 探索演示仓库,其中包含所讨论概念的实际实现。

🔍 如果你想了解这一切的起源,请重温我文章中的第一部分:“使用 AWS Lambda 和 NestJS 进行本地开发:Docker、调试和热重载”

🤝 你的反馈非常宝贵!欢迎发表评论、提出问题或分享你的见解和优化。每一份贡献都有助于增强我们的集体知识并建立一个资源丰富的开发者社区。

祝你编码愉快,期待在即将到来的指南中探索使用 SQS 的异步请求!🚀