如何修改C# Azure Function将Event Hub接收的消息POST到REST API
解决方案
1. 函数代码更新
首先你需要引入HTTP请求相关命名空间,使用静态HttpClient实例发起POST请求(避免反复实例化导致的套接字资源耗尽问题),修改后的完整代码如下:
using System; using System.Collections.Generic; using System.Linq; using System.Text; using System.Threading.Tasks; using Microsoft.Azure.EventHubs; using Microsoft.Azure.WebJobs; using Microsoft.Extensions.Logging; // 新增HTTP相关引用 using System.Net.Http; using System.Net.Http.Headers; namespace PWO.Function.NS { public static class EventHubTrigger { // 声明静态HttpClient实例,全局复用 private static readonly HttpClient _httpClient = new HttpClient(); [FunctionName("EventHubTrigger")] public static async Task Run([EventHubTrigger("my-events", Connection = "Eventhub_recieverpolicy")] EventData[] events, ILogger log) { var exceptions = new List<Exception>(); // 从环境配置读取目标REST API地址,避免硬编码 string targetApiUrl = Environment.GetEnvironmentVariable("TargetRestApiEndpoint"); foreach (EventData eventData in events) { try { string messageBody = Encoding.UTF8.GetString(eventData.Body.Array, eventData.Body.Offset, eventData.Body.Count); log.LogInformation($"C# Event Hub trigger function processed a message: {messageBody}"); // 构造POST请求内容 var content = new StringContent(messageBody, Encoding.UTF8, "application/json"); // 如果你的API需要认证,比如请求头带密钥,可以在这里加 // content.Headers.Add("X-API-Key", Environment.GetEnvironmentVariable("TargetApiKey")); // 发起POST请求 HttpResponseMessage response = await _httpClient.PostAsync(targetApiUrl, content); // 确保请求成功,失败会抛出异常进入catch逻辑 response.EnsureSuccessStatusCode(); log.LogInformation($"消息推送成功,API返回状态码:{(int)response.StatusCode}"); } catch (Exception e) { log.LogError(e, $"消息推送失败,消息内容:{messageBody}"); exceptions.Add(e); } } if (exceptions.Count > 1) throw new AggregateException(exceptions); if (exceptions.Count == 1) throw exceptions.Single(); } } }
2. 配置文件修改说明
function.json
你使用的是C#进程内模型,通过特性声明触发规则,function.json会在编译时自动生成,不需要手动修改。如果需要验证配置正确性,编译后生成的function.json中触发规则需要和你代码中声明的Event Hub名称、连接字符串配置名保持一致即可。
host.json
推荐按如下配置调整,可根据你的实际业务需求修改参数:
{ "version": "2.0", "logging": { "applicationInsights": { "samplingSettings": { "isEnabled": true, "excludedTypes": "Request" } } }, "extensions": { "eventHubs": { // 每处理多少批消息提交一次检查点,1表示每批都提交避免重复消费 "batchCheckpointFrequency": 1, "eventProcessorOptions": { // 单次拉取的最大消息数量 "maxBatchSize": 32, // 预拉取消息数量,提升消费性能 "prefetchCount": 64 } } }, // 函数最大执行超时时间,消费计划下最长可设为00:10:00 "functionTimeout": "00:05:00" }
3. 环境配置补充
你需要在函数应用的配置中添加以下参数:
- 本地调试时,在
local.settings.json的Values节点下添加:{ "Values": { "AzureWebJobsStorage": "你的存储账户连接字符串", "FUNCTIONS_WORKER_RUNTIME": "dotnet", "Eventhub_recieverpolicy": "你的Event Hub接收策略连接字符串", "TargetRestApiEndpoint": "你要推送的REST API完整地址" // 如果API需要密钥,加下面这行 // "TargetApiKey": "你的API访问密钥" } } - 部署到Azure后,在函数应用的「配置」-「应用程序设置」中添加和上面相同的键值对即可。
内容的提问来源于stack exchange,提问作者Stark
相关产品推荐
相关产品推荐

