发布于 2025-12-10 0 阅读
0

AWS 无服务器入门 - Step Functions

AWS 无服务器入门 - Step Functions

TL;DR

在本系列中,我将尝试讲解 AWS 上无服务器的基础知识,以便您构建自己的无服务器应用程序。在上一篇文章中,我们学习了如何使用 Cognito 创建受身份验证保护的 REST API。在本文中,我们将深入探讨 Step Functions,这项服务允许您通过构建与其他 AWS 服务交互的状态机来编排您的无服务器应用程序。

⬇️ 我会定期发布无服务器内容,如果你想了解更多 ⬇️

在 Twitter 上关注我🚀

快速公告:我还在开发一个名为🛡 sls-mentor 🛡的库。它汇集了 30 条无服务器最佳实践,这些实践会在您的 AWS 无服务器项目中自动检查(无论使用哪种框架)。它是免费开源的,欢迎随时查看!

在 Github 上查找 sls-mentor⭐️

介绍

在构建无服务器应用程序时,您希望最小化每个组件(尤其是 Lambda 函数)的职责。您希望它们只做一件事,并且做好它。

为了避免构建具有大量副作用且难以调试的单体函数,一个简单的方法是使用 AWS Step Functions。您无需使用单个 Lambda 来处理数据检索、处理、存储以及消息传递等副作用,而是可以将这些职责拆分为多个 Lambda,并使用 Step Functions 进行编排。您甚至可以摆脱 Lambda 函数,直接与其他 AWS 服务(例如 DynamoDB、S3、SQS、SNS 等)集成:这称为无函数编程。

今天,我们将一起构建一个简单的购物应用程序。我们将创建一个由产品组成的数据库,每种产品都有库存。然后,我们将创建一个状态机,用于根据购物车创建订单。首先,它将检查所有产品是否有库存,然后创建订单并同时更新每件产品的库存。此流程将使用 Step Functions 完成,并允许在产品缺货时避免创建订单。

我们可以像这样描述我们的应用程序的架构:

建筑学

数据库将存储产品和订单,并/create-store-item允许创建新产品。

状态机将由带有路由的 API 网关触发/create-order。首先,借助Map状态,它将检查所有产品是否有货。

然后,如果成功,它将并行执行:

  • 创建订单
  • 根据另一个状态更新每种产品的库存Map

映射状态允许状态机遍历一个数组项,并为每个项执行一个任务。如果其中一个任务失败,则整个状态机失败。如果所有任务都成功,状态机将继续进入下一个状态。

创建数据库、API 和所有必要的 Lambda 函数

让我们从创建应用程序的构建块开始。我们将创建一个 DynamoDB 表来存储产品和订单,一个 API 网关来触发状态机,以及状态机将使用的所有 Lambda 函数。

为了配置这些资源,我将使用 AWS CDK 和 TypeScript。如果您不熟悉此方法,请随时查看本系列的先前文章,我将在其中解释如何使用它。本节没有什么新内容,我在之前的文章中已经介绍过 Lambda 函数、DynamoDB 表和 API 网关。

调配资源

首先创建一个新的 CDK 堆栈:



import * as cdk from 'aws-cdk-lib';
import { Construct } from 'constructs';
import path from 'path';

export class LearnServerlessStack extends cdk.Stack {
  constructor(scope: Construct, id: string, props?: cdk.StackProps) {
    super(scope, id, props);

    // Provision a new REST API Gateway
    const myFirstApi = new cdk.aws_apigateway.RestApi(this, 'myFirstApi', {});

    // Provision a new DynamoDB table
    const storeDB = new cdk.aws_dynamodb.Table(this, 'storeDB', {
      partitionKey: {
        name: 'PK',
        type: cdk.aws_dynamodb.AttributeType.STRING,
      },
      sortKey: {
        name: 'SK',
        type: cdk.aws_dynamodb.AttributeType.STRING,
      },
      billingMode: cdk.aws_dynamodb.BillingMode.PAY_PER_REQUEST,
    });

    // Provision a new Lambda function, and grant it read access to the DynamoDB table
    const isItemInStock = new cdk.aws_lambda_nodejs.NodejsFunction(this, 'isItemInStock', {
      entry: path.join(__dirname, 'isItemInStock', 'handler.ts'),
      handler: 'handler',
      environment: {
        TABLE_NAME: storeDB.tableName,
      },
    });
    storeDB.grantReadData(isItemInStock);

    // Provision a new Lambda function, and grant it write access to the DynamoDB table
    const updateItemStock = new cdk.aws_lambda_nodejs.NodejsFunction(this, 'updateItemStock', {
      entry: path.join(__dirname, 'updateItemStock', 'handler.ts'),
      handler: 'handler',
      environment: {
        TABLE_NAME: storeDB.tableName,
      },
    });
    storeDB.grantWriteData(updateItemStock);

    // Provision a new Lambda function, and grant it write access to the DynamoDB table
    const createOrder = new cdk.aws_lambda_nodejs.NodejsFunction(this, 'createOrder', {
      entry: path.join(__dirname, 'createOrder', 'handler.ts'),
      handler: 'handler',
      environment: {
        TABLE_NAME: storeDB.tableName,
      },
    });
    storeDB.grantWriteData(createOrder);

    // Provision a new Lambda function, and grant it write access to the DynamoDB table
    const createStoreItem = new cdk.aws_lambda_nodejs.NodejsFunction(this, 'createStoreItem', {
      entry: path.join(__dirname, 'createStoreItem', 'handler.ts'),
      handler: 'handler',
      environment: {
        TABLE_NAME: storeDB.tableName,
      },
    });
    storeDB.grantWriteData(createStoreItem);

    // Add a new POST route to the REST API Gateway, and link it to the createStoreItem Lambda function
    const createStoreItemResource = myFirstApi.root.addResource('create-store-item');
    createStoreItemResource.addMethod('POST', new cdk.aws_apigateway.LambdaIntegration(createStoreItem));
  }
}


Enter fullscreen mode Exit fullscreen mode

这里没有什么新内容,我创建必要的资源,将表名作为环境变量传递给 Lambda 函数,并授予它们必要的权限。

向 Lambda 函数添​​加代码

资源现已配置完毕,唯一缺少的是 4 个 Lambda 函数的代码。让我们从createStoreItem函数开始:



import { DynamoDBClient, PutItemCommand } from '@aws-sdk/client-dynamodb';

const client = new DynamoDBClient({});

export const handler = async ({ body }: { body: string }): Promise<{ statusCode: number; body: string }> => {
  const tableName = process.env.TABLE_NAME;

  const { itemId, quantity } = JSON.parse(body) as { itemId?: string; quantity?: number };

  if (itemId === undefined || quantity === undefined) {
    return {
      statusCode: 200,
      body: JSON.stringify({ message: 'itemId or quantity is undefined' }),
    };
  }

  await client.send(
    new PutItemCommand({
      TableName: tableName,
      Item: {
        PK: { S: 'StoreItem' },
        SK: { S: itemId },
        stock: { N: quantity.toString() },
      },
    }),
  );

  return {
    statusCode: 200,
    body: JSON.stringify({ message: 'Store item created' }),
  };
};


Enter fullscreen mode Exit fullscreen mode

这个 lambda 函数将被插入到 API 中,因此它必须遵循特定的类型。我快速检查了输入,然后使用 DynamoDB SDK 在表中创建一个新项。

在这里,我选择了数据库中 StoreItems 的数据结构:PK是一个常量字符串StoreItem, 是SKitemId我还添加了一个stock属性,它是一个数字。

如果您对 DynamoDB 不太熟悉,请随时查看我的dynamoDB文章,我在其中解释了如何使用它。

然后,让我们创建isItemInStock函数:



import { DynamoDBClient, GetItemCommand } from '@aws-sdk/client-dynamodb';

const client = new DynamoDBClient({});

export const handler = async ({
  item: { itemId, quantity },
}: {
  item: { itemId: string; quantity: number };
}): Promise<void> => {
  const tableName = process.env.TABLE_NAME;

  const { Item } = await client.send(
    new GetItemCommand({
      TableName: tableName,
      Key: {
        PK: { S: 'StoreItem' },
        SK: { S: itemId },
      },
    }),
  );

  const stock = Item?.stock.N;

  if (stock === undefined || +stock < quantity) {
    throw new Error('Item not in stock');
  }
};


Enter fullscreen mode Exit fullscreen mode

因为它是未来状态机的构建块,所以输入更简单。此函数接收一个由itemId和组成的单个商品quantity。它会在商店数据库中检查该商品是否有货,如果没有,则会抛出错误。

然后,让我们创建updateItemStock函数:



import { DynamoDBClient, UpdateItemCommand } from '@aws-sdk/client-dynamodb';

const client = new DynamoDBClient({});

export const handler = async ({
  item: { itemId, quantity },
}: {
  item: { itemId: string; quantity: number };
}): Promise<void> => {
  const tableName = process.env.TABLE_NAME;

  await client.send(
    new UpdateItemCommand({
      TableName: tableName,
      Key: {
        PK: { S: 'StoreItem' },
        SK: { S: itemId },
      },
      UpdateExpression: 'SET stock = stock - :quantity',
      ExpressionAttributeValues: {
        ':quantity': { N: quantity.toString() },
      },
    }),
  );
};


Enter fullscreen mode Exit fullscreen mode

在这个 lambda 表达式中,输入是一样的。为了在不知道商品当前价值的情况下更新其库存,我使用UpdateExpression带有SET关键字的 an 。我还使用 anExpressionAttributeValuesquantity值传递给表达式。

最后,让我们创建createOrder函数:



import { DynamoDBClient, PutItemCommand } from '@aws-sdk/client-dynamodb';
import { v4 as uuid } from 'uuid';

const client = new DynamoDBClient({});

export const handler = async ({ order }: { order: { itemId: string; quantity: number }[] }): Promise<void> => {
  const tableName = process.env.TABLE_NAME;

  await client.send(
    new PutItemCommand({
      TableName: tableName,
      Item: {
        PK: { S: 'Order' },
        SK: { S: uuid() },
        order: {
          L: order.map(({ itemId, quantity }) => ({
            M: { itemId: { S: itemId }, quantity: { N: quantity.toString() } },
          })),
        },
      },
    }),
  );
};


Enter fullscreen mode Exit fullscreen mode

这里,我在数据库中添加了一个新项,其中包含一个随机数SK和一个常量PKorder属性是一个包含多个项的数组(L 代表 List,M 代表 Map),每个项由 和itemId一个组成quantity

代码写完了!注意,我创建的 Lambda 函数都没有超过一个职责。它们都非常简单,而且易于调试。

此外,它们不返回任何值:状态机将处理它们之间的数据流。

创建状态机来协调应用程序

本文精彩的部分来了!我们将配置一个新的状态机,它将处理 3 个 Lambda 函数isItemInStockupdateItemStock和之间的数据流createOrder

首先,让我们确定一个数据结构:状态机的输入将是以下类型:{ order: { itemId: string, quantity: number }[] }

注意它如何与我为 Lambda 函数定义的输入类型交互:

  • isItemInStockupdateItemStock期望输入的单个项目
  • createOrder期望完整的项目数组。

我们先从创建任务开始。任务是状态机的基石。我将创建 3 个任务,每个 Lambda 函数一个:



import { JsonPath } from 'aws-cdk-lib/aws-stepfunctions';

// ...

// Create a Map task, iterating over the items of the input
const isItemInStockMappedTask = new cdk.aws_stepfunctions.Map(this, 'isItemInStockMappedTask', {
  itemsPath: '$.order',
  resultPath: JsonPath.DISCARD,
  parameters: {
    'item.$': '$$.Map.Item.Value',
  },
}).iterator(
  new cdk.aws_stepfunctions_tasks.LambdaInvoke(this, 'isItemInStockTask', {
    lambdaFunction: isItemInStock,
  }),
);

// Create a Map task, iterating over the items of the input
const updateItemStockMappedTask = new cdk.aws_stepfunctions.Map(this, 'updateItemStockMappedTask', {
  itemsPath: '$.order',
  resultPath: JsonPath.DISCARD,
  parameters: {
    'item.$': '$$.Map.Item.Value',
  },
}).iterator(
  new cdk.aws_stepfunctions_tasks.LambdaInvoke(this, 'updateItemStockTask', {
    lambdaFunction: updateItemStock,
  }),
);

// Create simple task, calling the createOrder Lambda function
const createOrderTask = new cdk.aws_stepfunctions_tasks.LambdaInvoke(this, 'createOrderTask', {
  lambdaFunction: createOrder,
});


Enter fullscreen mode Exit fullscreen mode

上面的代码片段有很多需要讨论的地方:

  • Map 任务用于遍历项目数组,并为每个项目执行一项任务。

    • itemsPath用于指定在输入中查找项目数组的位置。根据输入,项目数组位于。这是一种特殊的 AWS 语法,您可以在此处$.order了解更多信息
    • parameters用于指定任务的输入。在这里,我想将迭代过的项目传递给 Lambda 函数。我使用特殊语法$$来访问迭代的当前项目。更多详细信息,请参阅 AWS 文档。
    • resultPath用于指定任务结果的存储位置。这里我不想存储任务结果,所以就舍弃了它。
  • LambdaInvoke 任务用于调用 Lambda 函数。我们只需根据之前预置的内容指定要调用的 Lambda 函数即可。

状态机的三个构建块已经准备好了,让我们通过创建一个新的状态机来协调它们:



const parallelState = new cdk.aws_stepfunctions.Parallel(this, 'parallelState', {});

parallelState.branch(updateItemStockMappedTask, createOrderTask);

const definition = isItemInStockMappedTask.next(parallelState);

const myFirstStateMachine = new cdk.aws_stepfunctions.StateMachine(this, 'myFirstStateMachine', {
  definition,
});


Enter fullscreen mode Exit fullscreen mode

查看定义,我们发现首先执行isItemInStockMappedTask,然后执行parallelStateparallelState由两个分支组成:

  • updateItemStockMappedTask
  • createOrderTask

很酷的是,并行状态只有在isItemInStockMappedTask成功时才会执行(即 Map 中的每个任务都成功)。如果失败,状态机就会停止。

使用 API 网关触发状态机

锦上添花的是:我们将通过 API 调用来触发状态机。为了演示无函数编程,我们将不使用 Lambda 函数来触发状态机,而是直接使用 API 网关。

为了做到这一点,我要做两件事:

  • 创建一个 IAM 角色,允许 API 网关触发状态机
  • 在 API 网关上创建一个新的 POST 路由,并将其链接到状态机。该路由将“承担” IAM 角色,以便能够触发状态机。

创建 IAM 角色来执行状态机

IAM 角色的创建非常简单:



const invokeStateMachineRole = new cdk.aws_iam.Role(this, 'invokeStateMachineRole', {
  assumedBy: new cdk.aws_iam.ServicePrincipal('apigateway.amazonaws.com'),
});

invokeStateMachineRole.addToPolicy(
  new cdk.aws_iam.PolicyStatement({
    actions: ['states:StartExecution'],
    resources: [myFirstStateMachine.stateMachineArn],
  }),
);


Enter fullscreen mode Exit fullscreen mode

首先,通过设置“主体”,我定义该角色将由 API 承担。然后,我添加一条允许该 API 启动状态机执行的策略,并指定状态机 ARN,以将策略范围缩小到仅此状态机。

使用 POST 路由触发状态机

这部分有点复杂:



const createOrderResource = myFirstApi.root.addResource('create-order');

createOrderResource.addMethod(
  'POST',
  new cdk.aws_apigateway.Integration({
    type: cdk.aws_apigateway.IntegrationType.AWS,
    integrationHttpMethod: 'POST',
    uri: `arn:aws:apigateway:${cdk.Aws.REGION}:states:action/StartExecution`,
    options: {
      credentialsRole: invokeStateMachineRole,
      requestTemplates: {
        'application/json': `{
        "input": "{\\"order\\": $util.escapeJavaScript($input.json('$'))}",
        "stateMachineArn": "${myFirstStateMachine.stateMachineArn}"
      }`,
      },
      integrationResponses: [
        {
          statusCode: '200',
          responseTemplates: {
            'application/json': `{
            "statusCode": 200,
            "body": { "message": "OK!" }"
          }`,
          },
        },
      ],
    },
  }),
  {
    methodResponses: [
      {
        statusCode: '200',
      },
    ],
  },
);


Enter fullscreen mode Exit fullscreen mode

首先,我向 API 添加一个新的 POST 路由。然后,我不再像我们习惯的那样将其插入 Lambda 函数,而是定义了一个自定义集成。此集成允许我直接调用状态机。

值得注意的一点是:

  • uri是一种特殊的 AWS 语法,允许调用特定的 AWS 服务。在这里,我调用的是StartExecution该服务的操作states,即 Step Functions 服务。
  • options

    • 我指定要承担的 IAM 角色(我之前创建的角色)
    • 我指定了状态机的输入。在字段中input,我使用与之前相同的 AWS 语法,将 POST 路由的主体转换为状态机所需的数据结构。在字段中stateMachineArn,我指定了要调用的状态机的 ARN。
    • 我指定了 API 的响应。在这里,我只是返回一个带有消息的 200 状态代码。
  • 最后,在方法的响应中,我指定了 API 的响应。这里,我直接返回 200 状态码。此响应必须与 中定义的响应之一匹配,integrationResponses否则 API 将返回 500 状态码。

这部分内容相当高级,我找到了一篇非常棒的文章,对这个集成进行了更详细的介绍。强烈推荐你阅读!如有任何疑问,请随时联系我,我很乐意为您提供帮助!

代码写完了!接下来就是部署了:



npm run cdk deploy

Enter fullscreen mode Exit fullscreen mode




测试应用程序

首先,让我们从创建新产品开始,这是一个示例请求:

创建产品

然后,让我们创建一个新的订单:

创建订单

查看状态机的日志,我们可以看到它成功了,并且我们发现了与文章开头计划相同的结构:

状态机日志

最后,我们来检查一下DB的内容:

数据库内容

我们可以看到订单已经创建完成,并且各个产品的库存也已经更新。

如果我们尝试重复相同的请求,状态机最终将失败,因为产品不再有库存:

状态机失败

结论

本教程是对 Step Functions 的简单介绍。这项服务还有更多功能,尤其是使用 SES 或 SNS 等消息服务向用户发送通知。敬请期待,我将在后续文章中探讨这些主题!

我计划每两个月更新一次这一系列文章。我已经介绍了如何创建简单的 Lambda 函数和 REST API,以及如何与 DynamoDB 数据库和 S3 存储桶进行交互。您可以在我的代码库中关注这些进展!我将介绍一些新的主题,例如创建事件驱动的应用程序、类型安全等等。如果您有任何建议,请随时联系我!

如果您能回复并分享这篇文章给您的朋友和同事,我将不胜感激。这将极大地帮助我扩大读者群。另外,别忘了订阅,以便及时收到下一篇文章的更新!

如果您想与我保持联系,请访问我的Twitter 账号。我经常发布或转发有关 AWS 和无服务器的有趣内容,欢迎关注我!

在 Twitter 上关注我🚀

鏂囩珷鏉ユ簮锛�https://dev.to/slsbytheodo/learn-serverless-on-aws-step-by-step-step-functions-4m7c