You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何修改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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.30 23:39:04