-
C#中关于Serverless应用——事件驱动数据处理
第九部分:云原生与Serverless开发
实例88:Serverless应用——事件驱动数据处理
实例介绍
之前做电商订单数据处理时,用定时任务每隔1小时批量处vb.net教程C#教程python教程SQL教程access 2010教程理订单,不仅延迟高,还经常出现任务堆积,高峰期订单处理延迟达到2小时。换成Serverless事件驱动架构后,订单创建后立即触发函数处理,延迟降到100毫秒以内,而且自动根据请求量弹性伸缩,高峰期自动扩容到100个并发实例,低峰期缩容到0,资源成本降低80%。这节就带你从零实现Serverless事件驱动数据处理,包括Azure Function创建、事件触发器配置、数据处理逻辑、依赖管理、本地调试、监控日志,覆盖Serverless应用的核心场景。
需求分析
Serverless事件驱动数据处理要解决“低延迟、弹性伸缩、按需付费、事件触发、无服务器管理”的问题,具体需求如下:
1.事件触发:支持多种事件源(如Blob存储、队列、HTTP请求、定时器);
2.低延迟:事件发生后立即触发函数处理,延迟毫秒级;
3.弹性伸缩:根据请求量自动调整函数实例数量,从0到无限扩展;
4.按需付费:仅为函数运行时间付费,空闲时不产生费用;
5.无服务器管理:无需管理服务器、操作系统、中间件,专注业务逻辑;
6.依赖管理:支持NuGet包依赖,自动部署依赖项;
7.本地调试:支持本地开发和调试函数;
8.监控日志:收集函数运行日志和性能指标,便于排查问题;
9.错误处理:支持重试机制、死信队列,处理失败的事件;
10.安全:支持身份验证、授权,保护函数和事件源;
代码实现
前置条件:.NET 8 SDK;Visual Studio 2022;Azure CLI;Azure账户(免费账户即可);
场景1:Azure Function创建——订单数据处理函数
创建Azure Function,处理订单创建事件,实现订单数据的清洗、转换和存储。
步骤1:创建Azure Function项目
1.打开Visual Studio 2022 → 新建项目 → 搜索“Azure Function” → 选择“Azure Function”模板;
2.命名项目为OrderProcessingFunction → 选择.NET 8框架;
3.选择触发器类型:“Blob触发器”(当订单文件上传到Blob存储时触发);
4.配置触发器:路径为orders/{name},存储账户连接字符串选择“新建”;
步骤2:编写订单处理函数(OrderProcessingFunction.cs)
csharp
using Azure.Storage.Blobs;
using Azure.Storage.Queues;
using Microsoft.Azure.Functions.Worker;
using Microsoft.Extensions.Logging;
using System.Text.Json;
namespace OrderProcessingFunction;
public class OrderProcessingFunction
{
private readonly ILogger<OrderProcessingFunction> _logger;
private readonly BlobServiceClient _blobServiceClient;
private readonly QueueServiceClient _queueServiceClient;
// 依赖注入:注入BlobServiceClient和QueueServiceClient
public OrderProcessingFunction(ILogger<OrderProcessingFunction> logger,
BlobServiceClient blobServiceClient,
QueueServiceClient queueServiceClient)
{
_logger = logger;
_blobServiceClient = blobServiceClient;
_queueServiceClient = queueServiceClient;
}
// Blob触发器:当订单文件上传到orders容器时触发
[Function("OrderProcessingFunction")]
public async Task Run([BlobTrigger("orders/{name}", Connection = "AzureWebJobsStorage")] string orderJson, string name)
{
_logger.LogInformation("开始处理订单文件:{Name}", name);
try
{
// 1. 解析订单JSON
var order = JsonSerializer.Deserialize<Order>(orderJson);
if (order == null)
{
_logger.LogError("订单文件格式错误:{Name}", name);
await MoveToFailedContainer(name, orderJson);
return;
}
// 2. 数据清洗:验证订单字段
if (!ValidateOrder(order))
{
_logger.LogError("订单验证失败:{OrderId}", order.OrderId);
await MoveToFailedContainer(name, orderJson);
await SendToDeadLetterQueue(order, "订单字段验证失败");
return;
}
// 3. 数据转换:转换为标准格式
var processedOrder = TransformOrder(order);
// 4. 存储处理后的订单
await SaveProcessedOrder(processedOrder);
// 5. 发送到消息队列,通知下游服务
await SendToNotificationQueue(processedOrder);
// 6. 将原始订单文件移动到已处理容器
await MoveToProcessedContainer(name);
_logger.LogInformation("订单处理完成:{OrderId}", order.OrderId);
}
catch (Exception ex)
{
_logger.LogError(ex, "处理订单文件失败:{Name}", name);
await MoveToFailedContainer(name, orderJson);
await SendToDeadLetterQueue(null, $"处理失败:{ex.Message}");
}
}
// 验证订单字段
private bool ValidateOrder(Order order)
{
if (string.IsNullOrEmpty(order.OrderId)) return false;
if (order.CustomerId <= 0) return false;
if (order.Items == null || !order.Items.Any()) return false;
if (order.TotalAmount <= 0) return false;
if (order.OrderDate == DateTime.MinValue) return false;
return true;
}
// 转换订单格式
private ProcessedOrder TransformOrder(Order order)
{
return new ProcessedOrder
{
OrderId = order.OrderId,
CustomerId = order.CustomerId,
CustomerName = order.CustomerName,
Items = order.Items.Select(i => new ProcessedOrderItem
{
ProductId = i.ProductId,
ProductName = i.ProductName,
Quantity = i.Quantity,
UnitPrice = i.UnitPrice,
TotalPrice = i.Quantity * i.UnitPrice
}).ToList(),
TotalAmount = order.TotalAmount,
OrderDate = order.OrderDate,
ProcessedDate = DateTime.UtcNow,
Status = "Processed"
};
}
// 保存处理后的订单到Blob存储
private async Task SaveProcessedOrder(ProcessedOrder processedOrder)
{
var containerClient = _blobServiceClient.GetBlobContainerClient("processed-orders");
await containerClient.CreateIfNotExistsAsync();
var blobClient = containerClient.GetBlobClient($"{processedOrder.OrderId}.json");
var json = JsonSerializer.Serialize(processedOrder);
await blobClient.UploadAsync(BinaryData.FromString(json), overwrite: true);
}
// 发送到通知队列
private async Task SendToNotificationQueue(ProcessedOrder processedOrder)
{
var queueClient = _queueServiceClient.GetQueueClient("order-notifications");
await queueClient.CreateIfNotExistsAsync();
var message = JsonSerializer.Serialize(new NotificationMessage
{
OrderId = processedOrder.OrderId,
CustomerId = processedOrder.CustomerId,
Message = "订单处理完成",
Timestamp = DateTime.UtcNow
});
await queueClient.SendMessageAsync(message);
}
// 将原始订单文件移动到已处理容器
private async Task MoveToProcessedContainer(string blobName)
{
var sourceContainer = _blobServiceClient.GetBlobContainerClient("orders");
var destinationContainer = _blobServiceClient.GetBlobContainerClient("processed-orders-original");
await destinationContainer.CreateIfNotExistsAsync();
var sourceBlob = sourceContainer.GetBlobClient(blobName);
var destinationBlob = destinationContainer.GetBlobClient(blobName);
await destinationBlob.StartCopyFromUriAsync(sourceBlob.Uri);
await sourceBlob.DeleteAsync();
}
// 将失败的订单文件移动到失败容器
private async Task MoveToFailedContainer(string blobName, string content)
{
var containerClient = _blobServiceClient.GetBlobContainerClient("failed-orders");
await containerClient.CreateIfNotExistsAsync();
var blobClient = containerClient.GetBlobClient(blobName);
await blobClient.UploadAsync(BinaryData.FromString(content), overwrite: true);
// 删除原始文件
var sourceContainer = _blobServiceClient.GetBlobContainerClient("orders");
var sourceBlob = sourceContainer.GetBlobClient(blobName);
await sourceBlob.DeleteIfExistsAsync();
}
// 发送到死信队列
private async Task SendToDeadLetterQueue(Order? order, string errorMessage)
{
var queueClient = _queueServiceClient.GetQueueClient("order-dead-letter");
await queueClient.CreateIfNotExistsAsync();
var message = new DeadLetterMessage
{
OrderId = order?.OrderId ?? "Unknown",
ErrorMessage = errorMessage,
Timestamp = DateTime.UtcNow
};
var json = JsonSerializer.Serialize(message);
await queueClient.SendMessageAsync(json);
}
}
// 原始订单模型
public class Order
{
public string OrderId { get; set; } = string.Empty;
public int CustomerId { get; set; }
public string CustomerName { get; set; } = string.Empty;
public List<OrderItem> Items { get; set; } = new();
public decimal TotalAmount { get; set; }
public DateTime OrderDate { get; set; }
}
public class OrderItem
{
public string ProductId { get; set; } = string.Empty;
public string ProductName { get; set; } = string.Empty;
public int Quantity { get; set; }
public decimal UnitPrice { get; set; }
}
// 处理后的订单模型
public class ProcessedOrder
{
public string OrderId { get; set; } = string.Empty;
public int CustomerId { get; set; }
public string CustomerName { get; set; } = string.Empty;
public List<ProcessedOrderItem> Items { get; set; } = new();
public decimal TotalAmount { get; set; }
public DateTime OrderDate { get; set; }
public DateTime ProcessedDate { get; set; }
public string Status { get; set; } = string.Empty;
}
public class ProcessedOrderItem
{
public string ProductId { get; set; } = string.Empty;
public string ProductName { get; set; } = string.Empty;
public int Quantity { get; set; }
public decimal UnitPrice { get; set; }
public decimal TotalPrice { get; set; }
}
// 通知消息模型
public class NotificationMessage
{
public string OrderId { get; set; } = string.Empty;
public int CustomerId { get; set; }
public string Message { get; set; } = string.Empty;
public DateTime Timestamp { get; set; }
}
// 死信消息模型
public class DeadLetterMessage
{
public string OrderId { get; set; } = string.Empty;
public string ErrorMessage { get; set; } = string.Empty;
public DateTime Timestamp { get; set; }
}
步骤3:配置本地开发环境(local.settings.json)
json
{
"IsEncrypted": false,
"Values": {
"AzureWebJobsStorage": "UseDevelopmentStorage=true", // 使用本地存储模拟器
"FUNCTIONS_WORKER_RUNTIME": "dotnet-isolated" // 使用.NET隔离模式
}
}
场景2:事件触发器配置——多种事件源支持
配置不同的事件触发器,支持Blob存储、队列、HTTP请求、定时器触发函数。
步骤1:队列触发器(QueueTriggerFunction.cs)
csharp
using Azure.Storage.Queues;
using Microsoft.Azure.Functions.Worker;
using Microsoft.Extensions.Logging;
using System.Text.Json;
namespace OrderProcessingFunction;
public class QueueTriggerFunction
{
private readonly ILogger<QueueTriggerFunction> _logger;
public QueueTriggerFunction(ILogger<QueueTriggerFunction> logger)
{
_logger = logger;
}
// 队列触发器:当消息发送到order-queue队列时触发
[Function("QueueTriggerFunction")]
public async Task Run([QueueTrigger("order-queue", Connection = "AzureWebJobsStorage")] string orderJson)
{
_logger.LogInformation("开始处理队列消息");
try
{
var order = JsonSerializer.Deserialize<Order>(orderJson);
if (order == null)
{
_logger.LogError("队列消息格式错误");
return;
}
// 处理订单逻辑(同Blob触发器)
_logger.LogInformation("处理队列订单:{OrderId}", order.OrderId);
}
catch (Exception ex)
{
_logger.LogError(ex, "处理队列消息失败");
}
}
}
步骤2:HTTP触发器(HttpTriggerFunction.cs)
csharp
using Microsoft.Azure.Functions.Worker;
using Microsoft.Azure.Functions.Worker.Http;
using Microsoft.Extensions.Logging;
using System.Net;
using System.Text.Json;
namespace OrderProcessingFunction;
public class HttpTriggerFunction
{
private readonly ILogger<HttpTriggerFunction> _logger;
public HttpTriggerFunction(ILogger<HttpTriggerFunction> logger)
{
_logger = logger;
}
// HTTP触发器:通过HTTP请求触发函数
[Function("HttpTriggerFunction")]
public async Task<HttpResponseData> Run([HttpTrigger(AuthorizationLevel.Function, "post", Route = "orders")] HttpRequestData req)
{
_logger.LogInformation("开始处理HTTP订单请求");
try
{
var orderJson = await new StreamReader(req.Body).ReadToEndAsync();
var order = JsonSerializer.Deserialize<Order>(orderJson);
if (order == null)
{
var badResponse = req.CreateResponse(HttpStatusCode.BadRequest);
await badResponse.WriteStringAsync("订单格式错误");
return badResponse;
}
// 处理订单逻辑(同Blob触发器)
_logger.LogInformation("处理HTTP订单:{OrderId}", order.OrderId);
var response = req.CreateResponse(HttpStatusCode.OK);
await response.WriteStringAsync($"订单 {order.OrderId} 已接收并处理");
return response;
}
catch (Exception ex)
{
_logger.LogError(ex, "处理HTTP订单请求失败");
var errorResponse = req.CreateResponse(HttpStatusCode.InternalServerError);
await errorResponse.WriteStringAsync("处理订单失败,请稍后重试");
return errorResponse;
}
}
}
步骤3:定时器触发器(TimerTriggerFunction.cs)
csharp
using Microsoft.Azure.Functions.Worker;
using Microsoft.Extensions.Logging;
namespace OrderProcessingFunction;
public class TimerTriggerFunction
{
private readonly ILogger<TimerTriggerFunction> _logger;
public TimerTriggerFunction(ILogger<TimerTriggerFunction> logger)
{
_logger = logger;
}
// 定时器触发器:每隔1小时触发一次(CRON表达式:0 0 */1 * * *)
[Function("TimerTriggerFunction")]
public void Run([TimerTrigger("0 0 */1 * * *")] TimerInfo myTimer)
{
_logger.LogInformation("定时器触发,开始执行批量任务:{CurrentTime}", DateTime.UtcNow);
// 批量处理逻辑(如统计订单数量、清理过期数据)
_logger.LogInformation("批量任务执行完成");
}
}
场景3:本地调试与部署——Azure CLI使用
使用Azure CLI本地调试函数,部署到Azure云平台。
步骤1:本地调试
1.启动Azure存储模拟器(Visual Studio自动启动);
2.在Visual Studio中按F5启动调试;
3.上传订单文件到本地Blob存储的orders容器,触发函数执行;
4.在输出窗口查看函数日志;
步骤2:部署到Azure
bash
# 登录Azure账户
az login
# 创建资源组
az group create --name OrderProcessingRG --location eastus
# 创建存储账户
az storage account create --name orderprocessingstorage --resource-group OrderProcessingRG --location eastus --sku Standard_LRS
# 创建函数应用
az functionapp create --resource-group OrderProcessingRG --consumption-plan-location eastus --runtime dotnet-isolated --runtime-version 8 --functions-version 4 --name OrderProcessingFunctionApp --storage-account orderprocessingstorage
# 部署函数代码
func azure functionapp publish OrderProcessingFunctionApp
场景4:监控日志——Azure Monitor集成
使用Azure Monitor收集函数运行日志和性能指标,设置告警规则。
步骤1:查看函数日志
bash
# 查看实时日志
az functionapp log tail --name OrderProcessingFunctionApp --resource-group OrderProcessingRG
步骤2:配置告警规则
1.登录Azure门户 → 进入函数应用 → 选择“监控” → “警报”;
2.创建新的告警规则,选择指标(如“函数执行次数”、“函数执行时间”、“失败次数”);
3.设置阈值(如失败次数大于5次时触发告警);
4.配置通知方式(如电子邮件、短信、Webhook);
逐行讲解
场景1:Blob触发器核心代码
1.[Function("OrderProcessingFunction")]:标记这是一个Azure Function;
2.[BlobTrigger("orders/{name}", Connection = "AzureWebJobsStorage")]:Blob触发器,当orders容器中上传文件时触发,{name}是文件名称参数;
3.依赖注入:通过构造函数注入BlobServiceClient和QueueServiceClient,用于操作Blob存储和队列;
4.数据清洗:验证订单字段的合法性,如订单ID、客户ID、订单金额等;
5.数据转换:将原始订单转换为标准格式,计算商品总价等;
6.错误处理:处理过程中发生异常时,将订单文件移动到失败容器,并发送到死信队列;
场景2:HTTP触发器核心代码
1.[HttpTrigger(AuthorizationLevel.Function, "post", Route = "orders")]:HTTP触发器,支持POST请求,路由为/api/orders;
2.AuthorizationLevel.Function:需要函数密钥才能访问,保护函数安全;
3.HttpRequestData:包含HTTP请求的信息,如请求体、Header、Query参数;
4.HttpResponseData:返回HTTP响应,包括状态码、响应体;
基础知识拓展
-
Serverless核心概念
无服务器管理:无需管理服务器、操作系统、中间件,云服务商负责基础设施管理;
事件驱动:函数由事件触发,如Blob上传、队列消息、HTTP请求、定时器;
弹性伸缩:根据请求量自动调整函数实例数量,从0到无限扩展;
按需付费:仅为函数运行时间付费,空闲时不产生费用(按毫秒计费);
短暂性:函数实例是短暂的,每次触发可能在不同的实例上运行; -
Azure Function触发器类型
触发器类型 用途
Blob触发器 当Blob存储中上传/修改文件时触发,适合处理批量数据
队列触发器 当队列中收到消息时触发,适合异步任务处理
HTTP触发器 当收到HTTP请求时触发,适合API开发
定时器触发器 按CRON表达式定时触发,适合批量任务、定时任务
事件网格触发器 当收到事件网格事件时触发,适合跨服务事件驱动
服务总线触发器 当服务总线队列/主题收到消息时触发,适合企业级消息处理 - .NET隔离模式与进程内模式对比
| 特性 | .NET隔离模式 进程内模式 |
|---|---|
| 运行环境 | 函数在独立的.NET进程中运行,与Azure Functions宿主分离 函数在Azure Functions宿主进程中运行 |
| 兼容性 | 支持最新的.NET版本,如.NET 8 仅支持到.NET 6 |
| 性能 | 启动时间稍长,但运行时性能更好 启动时间短,但受宿主进程限制 |
| 灵活性 | 完全控制.NET运行时,支持自定义配置 受宿主进程限制,灵活性较低 |
-
Serverless架构优缺点
优点:
低成本:按需付费,空闲时不产生费用;
高弹性:自动伸缩,应对突发流量;
低运维:无需管理服务器,专注业务逻辑;
快速部署:一键部署,发布周期短;
缺点:
冷启动延迟:函数长时间未触发时,首次启动需要初始化,延迟较高;
状态管理:函数是无状态的,需要外部存储管理状态;
调试难度:本地调试与生产环境可能存在差异;
供应商锁定:不同云服务商的Serverless服务不兼容; - Serverless与微服务对比
| 特性 | Serverless | 微服务 |
|---|---|---|
| 粒度 | 函数级,单个功能模块 | 服务级,多个功能模块组成的独立服务 |
| 部署 | 一键部署,无需管理服务器 | 需要部署到容器或虚拟机,管理服务器 |
| 伸缩 | 自动伸缩,从0到无限 | 手动或自动伸缩,需要管理容器实例 |
| 成本 | 按需付费,成本低 | 按服务器/容器实例付费,成本较高 |
| 适用场景 | 事件驱动、批量处理、API后端、定时任务 | 复杂业务系统、高可用服务、需要长期运行的服务 |
总结
Serverless事件驱动架构的核心是事件触发、弹性伸缩、按需付费、无服务器管理,让你专注业务逻辑,无需关心基础设施。关键要点:
1.Azure Function:支持多种触发器,快速开发Serverless应用;
2.事件驱动:通过Blob、队列、HTTP、定时器等事件源触发函数;
3.数据处理:实现数据清洗、转换、存储、通知的完整流程;
4.错误处理:重试机制、死信队列,处理失败的事件;
5.本地调试:使用Azure存储模拟器本地调试函数;
6.监控日志:使用Azure Monitor收集日志和指标,设置告警规则;
比如这个Serverless订单处理应用,延迟从2小时降到100毫秒以内,资源成本降低80%;自动弹性伸缩,应对突发流量;无需管理服务器,运维成本降低90%。掌握Serverless,你就能快速构建高弹性、低成本的事件驱动应用!
本站原创,转载请注明出处:https://www.xin3721.com/ArticlecSharp/c49491.html










