瑞典马工

云服务的致命伤:无法组合

案例缘起

客户需要对一个上游数据源做一个查询应用,供自己下辖的 200 多位职工使用。一个职工每天需要查很多次,时间次数都不定。要做一个查询应用给职工服务。

数据源是上级组织下发的一个 API 接口,不知道是服务性能问题还是故意限制,QPS 平均为 0.5,也就意味着这个 API 后面的服务是单线程,而且单次请求处理时长达到 2 秒。

案例架构

这个案例最难的部分不是前端页面和查询逻辑,用户的请求简单包装一下直接调上游 API 接口,基本上没有什么工作量。

最难的只有一点,就是在上游处理能力极为有限的前提下,如何在多人同时查询时,保证 API 的次序调用,避免超过限制导致失败。

所以在设计架构时,需要将主体服务和上游 API 解耦,以发布者订阅者的形式来处理查询过程。

这就引出了消息队列,它是一种跨进程的通信机制,主要用于传递上下游消息。

Image
案例架构图

在这个案例中,我们引入消息队列,来充当主体服务和查询处理服务 的缓冲区。

1. 主体服务的实现思路

当主体服务收到用户的查询请求后,会将用户的查询信息发送到消息队列中。另外一边,查询服务会按照上游 API 的合理速度,拉取消息队列中的消息,并执行查询和发邮件的过程。

如果有多个查询请求同时进来,则只会积压消息在消息队列中,查询处理服务仍然按照自己的速度消费队列中的消息。(对于用户来说,发送请求到收到邮件的时间间隔会延长)

我可以在主体服务中编写逻辑,通过时间戳 ID 等形式控制前端的查询请求频率,避免一个人短时间发送大量的查询请求独占资源。

主体服务可以使用无服务器形态的Lambda/云函数来承载,能应对短时大并发访问(虽然这里也用不到),降低了维护成本和空闲的运行成本。

2. 查询处理服务的实现思路

查询服务这里就不能盲目云函数了,因为这里有一个致命问题,就是上游 API 的单线程限制,如果消息队列并发触发大量的云函数来处理,则肯定会在上游 API 调用这里发生拥塞,造成的局面肯定是大量的失败。

这种情况必须有一个可以按照固定速率处理任务的服务。其实一开始我还寄希望于消息队列能够完成这个过程,可以控制消息的分发速度和并发数量,但调研了一圈发现满足要求又价格亲民的消息队列几乎没有,于是我打算自己做这个服务。

消息队列我基于价格和易用性考量,确认使用 AWS 的 SQS 服务。由于 SQS 服务 基于轮询的,本身并没有提供事件监听机制,因此我需要在服务中建立长轮询过程。

我的思路是,启动一个EC2服务器或 Fargate引擎(取决于价格因素)运行查询服务,在服务中轮询消息队列,当消息队列有数据时,进行查询处理和发邮件的流程,并标记消息已消费,然后继续轮询消息队列,直到下个消息到来。

这是一个 NodeJS 实现例子:

const AWS = require('aws-sdk');
const sqs = new AWS.SQS({ region: 'ap-southeast-1' }); // 初始化SQS消息队列客户端
const queueUrl = "https://sqs.ap-southeast-1.amazonaws.com/0/queue"; // SQS消息队列URL

async function processMessages() {
console.log(new Date().getTime(),'开始新一轮处理!')
let data = await sqs.receiveMessage({ // 监听接收消息
QueueUrl: queueUrl,
MaxNumberOfMessages: 1, // 每次只能处理一条
WaitTimeSeconds: 20 // 有消息并达到MaxNumberOfMessages数量后立刻返回,无消息最长20秒
}).promise();
if (data.Messages.length != 0) {
let message = data.Messages[0];
console.log(new Date().getTime(),'处理消息: ', message.MessageId);
// 调用上游API服务,大约2秒处理
// ...
// 发送邮件信息
await sqs.deleteMessage({ // 消息已经消费,删除消息
QueueUrl: queueUrl,
ReceiptHandle: message.ReceiptHandle
}).promise();
}
return;
}

async function run() {
while (true) { // 轮询监听
await processMessages();
}
}

run().catch(console.error);

在这个例子中,WaitTimeSeconds 为 20,在 20 秒内如果没有任何消息的话,sqs.receiveMessage 就会返回空消息。

这意味着,抛除有效的消息接收,这个服务每个月最低就要消耗掉 3600秒*24小时*30天/20=129600 次请求,用于轮询监听。按照0.4美元/一百万次的计费标准,每个月轮询SQS服务空耗费用为0.052 美元 ≈ 0.37 元人民币。

4. Terraform 配置

这个案例的接收查询请求的 lambda 函数代码,以及处理查询的服务代码因为没有多大的通用性,就不提供了,只提供在 AWS 中的 tf 配置,给大家展示全部服务的关联组成。

如果你也有类似的需求,可以后台和我联系。

// 定义本地变量
locals {
// 全局变量,用于构建AWS资源的ARN
global = "aws"
// AWS区域
region = "ap-southeast-1"
// 主体服务的Lambda函数的压缩包,案例不提供,有需要可以自己编写。
listenerfunction = "http_listener.zip"
}

// 定义AWS提供者
provider "aws" {
// 使用本地变量中定义的区域
region = local.region
}

// 获取当前AWS账户的信息
data "aws_caller_identity" "current" {}

// 创建SQS队列
resource "aws_sqs_queue" "my_queue" {
// 队列名称
name = "my-queue"
}

// 创建Lambda函数
resource "aws_lambda_function" "http_listener" {
// 函数名称
function_name = "http_listener"
// 执行角色
role = aws_iam_role.lambda_exec.arn
// 处理器
handler = "index.handler"
// 运行时环境
runtime = "nodejs18.x"
// 函数文件名
filename = local.listenerfunction
// 环境变量
environment {
variables = {
// SQS队列的URL
SQS_QUEUE_URL = aws_sqs_queue.my_queue.url
}
}
}

// 创建IAM角色
resource "aws_iam_role" "lambda_exec" {
// 角色名称
name = "lambda_exec_role"

// 角色的信任策略
assume_role_policy = <<EOF
{
"Version": "2012-10-17",
"Statement": [
{
"Action": "sts:AssumeRole",
"Principal": {
"Service": "lambda.amazonaws.com"
},
"Effect": "Allow",
"Sid": ""
}
]
}
EOF
}

// 创建IAM角色策略
resource "aws_iam_role_policy" "lambda_exec_policy" {
// 策略名称
name = "lambda_exec_policy"
// 角色ID
role = aws_iam_role.lambda_exec.id

// 策略内容
policy = <<EOF
{
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Action": "logs:CreateLogGroup",
"Resource": "${format("arn:%s:logs:%s:%s:*", local.global, local.region, data.aws_caller_identity.current.account_id)}"
},
{
"Effect": "Allow",
"Action": [
"logs:CreateLogStream",
"logs:PutLogEvents"
],
"Resource": [
"${format("arn:%s:logs:%s:%s:log-group:/aws/lambda/%s:*", local.global, local.region, data.aws_caller_identity.current.account_id, aws_lambda_function.http_listener.function_name)}"
]
},
{
"Effect": "Allow",
"Action": [
"sqs:SendMessage",
"sqs:ReceiveMessage",
"sqs:DeleteMessage",
"sqs:GetQueueAttributes"
],
"Resource": "${aws_sqs_queue.my_queue.arn}"
}
]
}
EOF
}

// 创建API Gateway
resource "aws_apigatewayv2_api" "my_api" {
// API名称
name = "my-api"
// 协议类型
protocol_type = "HTTP"
}

// 创建Lambda权限
resource "aws_lambda_permission" "apigw" {
// 声明ID
statement_id = "AllowExecutionFromAPIGateway"
// 动作
action = "lambda:InvokeFunction"
// 函数名称
function_name = aws_lambda_function.http_listener.function_name
// 主体
principal = "apigateway.amazonaws.com"

// 来源ARN
source_arn = "${aws_apigatewayv2_api.my_api.execution_arn}/*/*"
}

// 创建API Gateway的默认阶段
resource "aws_apigatewayv2_stage" "default" {
// API ID
api_id = aws_apigatewayv2_api.my_api.id
// 阶段名称
name = "$default"
// 自动部署
auto_deploy = true
}

// 创建API Gateway的集成
resource "aws_apigatewayv2_integration" "my_integration" {
// API ID
api_id = aws_apigatewayv2_api.my_api.id
// 集成类型
integration_type = "AWS_PROXY"

// 连接类型
connection_type = "INTERNET"
// 描述
description = "Lambda integration"
// 集成方法
integration_method = "POST"
// 集成URI
integration_uri = aws_lambda_function.http_listener.invoke_arn
// 载荷格式版本
payload_format_version = "2.0"
}

// 创建API Gateway的路由
resource "aws_apigatewayv2_route" "my_route" {
// API ID
api_id = aws_apigatewayv2_api.my_api.id
// 路由键
route_key = "ANY /call"

// 目标
target = "integrations/${aws_apigatewayv2_integration.my_integration.id}"
}

// 输出API端点
output "api_endpoint" {
// 描述
description = "The endpoint of the API Gateway"
// 值
value = "${aws_apigatewayv2_api.my_api.api_endpoint}/call"
}

// 创建EC2执行角色
resource "aws_iam_role" "ec2_exec" {
// 角色名称
name = "ec2_exec_role"

// 角色的信任策略
assume_role_policy = <<EOF
{
"Version": "2012-10-17",
"Statement": [
{
"Action": "sts:AssumeRole",
"Principal": {
"Service": "ec2.amazonaws.com"
},
"Effect": "Allow",
"Sid": ""
}
]
}
EOF
}

// 创建EC2执行策略
resource "aws_iam_role_policy" "ec2_exec_policy" {
// 策略名称
name = "ec2_exec_policy"
// 角色ID
role = aws_iam_role.ec2_exec.id

// 策略内容
policy = <<EOF
{
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Action": [
"sqs:SendMessage",
"sqs:ReceiveMessage",
"sqs:DeleteMessage",
"sqs:GetQueueAttributes"
],
"Resource": "${aws_sqs_queue.my_queue.arn}"
}
]
}
EOF
}

// 创建EC2实例配置文件
resource "aws_iam_instance_profile" "ec2_exec_profile" {
// 配置文件名称
name = "ec2_exec_profile"
// 角色名称
role = aws_iam_role.ec2_exec.name
}

// 创建EC2实例请求
resource "aws_spot_instance_request" "sqs_processor" {
// 镜像ID
ami = "ami-0fa377108253bf620"
// 实例类型
instance_type = "t2.micro"
// 竞价价格
spot_price = "0.01"
// IAM实例配置文件
iam_instance_profile = aws_iam_instance_profile.ec2_exec_profile.name

// 用户数据,用于在实例启动时运行脚本,这里只是举个例子,服务代码需要处理鉴权、日志等事项。
user_data = <<-EOF
#!/bin/bash
// 安装Node.js
curl -sL https://deb.nodesource.com/setup_18.x | sudo -E bash -
sudo apt-get install -y nodejs
// 从S3下载服务内容
aws s3 cp s3://my-bucket/sqs_processor.zip /home/ubuntu/
// 解压Lambda函数
unzip /home/ubuntu/sqs_processor.zip -d /home/ubuntu/
// 将Lambda函数添加到rc.local,使其在启动时运行
echo "node /home/ubuntu/index.js" >> /etc/rc.local
EOF
}

可复制的 DEMO

由于上面的案例上游API的限制有点太极端了,不太适用我们常见使用消息队列的场景。我自己整理了一个在 AWS 运行的开箱即用的架构配置。如果你感兴趣的话可以从这里开始实践一下,然后按照自己的需求自己修改生产和消费的逻辑。

首先是 AWS 的 tf配置 文件

// 定义一些本地变量
locals {
global = "aws" // 全局变量,表示我们正在使用的云服务提供商
region = "ap-southeast-1" // AWS的区域代码
listenerfunction = "http_listener.zip" // HTTP监听器的Lambda函数的文件名
processorfunction = "sqs_processor.zip" // SQS处理器的Lambda函数的文件名
}

// 定义AWS提供商,并设置区域
provider "aws" {
region = local.region
}

// 获取当前调用者的身份信息
data "aws_caller_identity" "current" {}

// 创建一个SQS队列
resource "aws_sqs_queue" "my_queue" {
name = "my-queue"
}

// 创建一个Lambda函数,用于HTTP监听
resource "aws_lambda_function" "http_listener" {
function_name = "http_listener"
role = aws_iam_role.lambda_exec.arn
handler = "index.handler"
runtime = "nodejs18.x"
filename = local.listenerfunction
environment {
variables = {
SQS_QUEUE_URL = aws_sqs_queue.my_queue.url // 环境变量,指向我们的SQS队列
}
}
}

// 创建一个Lambda函数,用于处理SQS消息
resource "aws_lambda_function" "sqs_processor" {
function_name = "sqs_processor"
role = aws_iam_role.lambda_exec.arn
handler = "index.handler"
runtime = "nodejs18.x"
filename = local.processorfunction
}

// 创建一个IAM角色,供Lambda函数使用
resource "aws_iam_role" "lambda_exec" {
name = "lambda_exec_role"

assume_role_policy = <<EOF
{
"Version": "2012-10-17",
"Statement": [
{
"Action": "sts:AssumeRole",
"Principal": {
"Service": "lambda.amazonaws.com"
},
"Effect": "Allow",
"Sid": ""
}
]
}
EOF
}

// 创建一个IAM策略,允许Lambda函数创建日志组,创建日志流,写入日志事件,以及操作SQS队列
resource "aws_iam_role_policy" "lambda_exec_policy" {
name = "lambda_exec_policy"
role = aws_iam_role.lambda_exec.id

policy = <<EOF
{
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Action": "logs:CreateLogGroup",
"Resource": "${format("arn:%s:logs:%s:%s:*", local.global, local.region, data.aws_caller_identity.current.account_id)}"
},
{
"Effect": "Allow",
"Action": [
"logs:CreateLogStream",
"logs:PutLogEvents"
],
"Resource": [
"${format("arn:%s:logs:%s:%s:log-group:/aws/lambda/%s:*", local.global, local.region, data.aws_caller_identity.current.account_id, aws_lambda_function.http_listener.function_name)}",
"${format("arn:%s:logs:%s:%s:log-group:/aws/lambda/%s:*", local.global, local.region, data.aws_caller_identity.current.account_id, aws_lambda_function.sqs_processor.function_name)}"
]
},
{
"Effect": "Allow",
"Action": [
"sqs:SendMessage",
"sqs:ReceiveMessage",
"sqs:DeleteMessage",
"sqs:GetQueueAttributes"
],
"Resource": "${aws_sqs_queue.my_queue.arn}"
}
]
}
EOF
}

// 创建一个Lambda事件源映射,将SQS队列与SQS处理器Lambda函数关联起来
resource "aws_lambda_event_source_mapping" "sqs_processor_mapping" {
depends_on = [aws_iam_role_policy.lambda_exec_policy]
event_source_arn = aws_sqs_queue.my_queue.arn
function_name = aws_lambda_function.sqs_processor.arn
}

// 创建一个API Gateway
resource "aws_apigatewayv2_api" "my_api" {
name = "my-api"
protocol_type = "HTTP"
}

// 创建一个Lambda权限,允许API Gateway调用我们的HTTP监听器Lambda函数
resource "aws_lambda_permission" "apigw" {
statement_id = "AllowExecutionFromAPIGateway"
action = "lambda:InvokeFunction"
function_name = aws_lambda_function.http_listener.function_name
principal = "apigateway.amazonaws.com"

source_arn = "${aws_apigatewayv2_api.my_api.execution_arn}/*/*"
}

// 创建一个API Gateway的默认阶段
resource "aws_apigatewayv2_stage" "default" {
api_id = aws_apigatewayv2_api.my_api.id
name = "$default"
auto_deploy = true
}

// 创建一个API Gateway的集成,将API Gateway与我们的HTTP监听器Lambda函数关联起来
resource "aws_apigatewayv2_integration" "my_integration" {
api_id = aws_apigatewayv2_api.my_api.id
integration_type = "AWS_PROXY"

connection_type = "INTERNET"
description = "Lambda integration"
integration_method = "POST"
integration_uri = aws_lambda_function.http_listener.invoke_arn
payload_format_version = "2.0"
}

// 创建一个API Gateway的路由,将任何请求都转发到我们的集成
resource "aws_apigatewayv2_route" "my_route" {
api_id = aws_apigatewayv2_api.my_api.id
route_key = "ANY /call"

target = "integrations/${aws_apigatewayv2_integration.my_integration.id}"
}

// 输出API Gateway的端点
output "api_endpoint" {
description = "The endpoint of the API Gateway"
value = "${aws_apigatewayv2_api.my_api.api_endpoint}/call"
}

由于 AWS 分为国际区和中国区,我给的配置是国际区的,如果你要在中国区运行,需要注意以下几项:

  1. locals.global 需要改为 aws-cn,并且初始认证时需要中国区的账号

  2. 中国区 API Gateway 需要在备案通过后才能访问,否则会一直返回权限问题。

另外 tf配置 中,还有两个 lambda代码包,我这里用的 nodejs 环境,两个代码包内容如下:

  • http_listener:这里用的时候需要自己做一下依赖配置。

const AWS = require('aws-sdk');
const sqs = new AWS.SQS();

exports.handler = async (event) => {
const params = {
MessageBody: JSON.stringify(event),
QueueUrl: process.env.SQS_QUEUE_URL
};
try {
const data = await sqs.sendMessage(params).promise();
console.log(`MessageID is ${data.MessageId}`);
return { statusCode: 200, body: '消息发送完毕!' };
} catch (error) {
console.error(error);
return { statusCode: 500, body: '发生错误!'+ error.toString() };
}
};

需要执行命令 npm install aws-sdk

{
"dependencies": {
"aws-sdk": "^2.1539.0"
}
}
  • sqs_processor:只有一个打印消息,SQS 有自动确认机制,因此不需要我们手动确认。

exports.handler = async (event) => {
for (const record of event.Records) {
const message = JSON.parse(record.body);
console.log(`Received message: ${JSON.stringify(message)}`);
// 这里可以处理消息
}
};

上面两个代码产物分别 zip 压缩打包,放置 tf配置 同级目录下 ,接着跑一下 tf配置:

terraform init
terraform apply

Image

然后就可以得到一个 API 网关地址,我们访问后效果如下:

Image

这里就是 http_listener 函数运行的结果,消息已经发送到 SQS 了,接着我们去 AWS 控制台看 sqs_processor 函数的日志结果,可以看到相关的打印消息。

Image

如果我们不需要了,我们可以 terraform destroy 直接销毁所有的服务资源。

Image

IaC(基础设施即代码)是一种快速且易复制的云资源操作方式,相比于在控制台一顿操作截图,这种代码的形式能够让你在自己的账号中轻松复制一模一样的资源配置,开箱效率大为增加。terraform 只是 IaC 的一种,如果你喜欢你可以用任何其他的形态。

写在后面

消息队列的使用在我们日常生活服务中很常见,比如我们常用的电话卡通查,银行卡通查,都使用了消息队列机制来处理查询过程。

Image

本来这篇文章想用国内云厂商的产品举例子,找了半天没有什么合适的,供大家复刻的难度有点大了。

1. 阿里云的云消息队列 rabbitmq

有 serverless 版本,开箱难度很小,但是这个静态用户名密码,我本寻思用资源角色鉴权完成这个部分,但是死活不通过。

Image

而且我在访问控制中也没有看到有相关的操作,为了确认我选的服务是对的,我还特意去 rabbitmq 下面抓了下请求包,前缀是ampq无误。

Image

2. 腾讯云的 CKafka 弹性 Topic

腾讯云选择从 Kafka 拆 Topic 来做按量产品,虽然开箱难度降低了,但也没办法提供角色鉴权,赤裸裸的还是密码,根本没有对云上服务场景做适应性改造。

Image

说能用吧倒是也能用,但是用起来就是比较别扭。本来挺自然的云内资源角色授权,被这么一搞还要外挂一个 凭证管理服务。(可以参考《一个不完美的云原生实践,无奈的选择》)

你说密码就密码吧,又不是第一次遇到。最让人崩溃的是,云函数触发器竟然没办法识别这个 弹性Topic,这就没法玩了,只能用本文中最开始的服务脚本大法了。

Image


我还是比较希望国内云厂商能够推出在体验和价格上更加适合小用户的消息队列产品。这样在做架构时,会有更多的选择,比较容易拿出一个在性能和价格上都很有优势的方案出来。

比较好的消息是,AWS中国区的服务种类和操作体验方面与国际区几乎相同,本文的 DEMO 可以完美的在中国区跑通,只需要备案就可以。

访问网址:https://www.amazonaws.cn

除了阿里云,腾讯云等本土厂商,我后续文章中也会展示 AWS中国区 的产品和推荐架构,以一个小用户的视角来看待不同云服务的能力和作用场景。

我计划不断探索和整理一些有意义且可复制的案例,以帮助更多的小型用户全面理解云服务对他们的价值。既然各大云服务商对此类案例不够重视,那我们就自己来填补这一空白。

关于作者:

ZiraLi,微信生态领域 MVP,就职于相关团队担任技术产品经理、架构师;为多家企业组织提供上云架构和微信生态产品咨询服务。欢迎关注我的公众号:咋用云

公众号文章所有的代码均归纳到 Github 仓库中,地址如下:

https://github.com/zzirali/usecloud