Azure Function+Dapr实现Service Bus消息调度及双重Data封装问题排查
Dapr结合Azure Function实现Service Bus延迟消息调度及双重data封装问题
问题描述
初始需求
想要创建基于C#和Dapr API的Azure Function,该函数在Azure Service Bus收到消息时触发,随后向Azure Service Bus主题发布测试消息,但不希望消息立即发布。调研发现DaprPubSubEvent中没有直接的调度设置选项,初始代码如下:
public static class TestFunc { [FunctionName("TestFunc")] public static async Task<IActionResult> Run( [HttpTrigger(AuthorizationLevel.Anonymous, "get", "post", Route = "TestFunc")] HttpRequest req, [DaprPublish(PubSubName = "testSub", Topic = "testTopic")] IAsyncCollector<DaprPubSubEvent> pubEvent, ILogger log) { log.LogInformation(" HTTP trigger function processed a request."); await pubEvent.AddAsync(new DaprPubSubEvent("Test Message")); return new OkObjectResult("TestResponse"); } }
询问是否有办法在Dapr中设置消息调度。
更新后的问题
修改代码添加ScheduledEnqueueTime元数据后,消息出现双重data字段封装问题。新代码如下:
public static class TestFunc { [FunctionName("TestFunc")] public static async Task<IActionResult> Run( [HttpTrigger(AuthorizationLevel.Anonymous, "get", "post", Route = "TestFunc")] HttpRequest req, [DaprPublish(PubSubName = "testSub", Topic = "testTopic")] IAsyncCollector<DaprPubSubEvent> pubEvent, ILogger log) { log.LogInformation(" HTTP trigger function processed a request."); var d = new Dictionary<string, object>() { { "ScheduledEnqueueTime", DateTimeOffset.Now.AddMinutes(5)} }; DaprBindingMessage dm = new DaprBindingMessage("Test Message", d); await pubEvent.AddAsync(new DaprPubSubEvent(dm)); return new OkObjectResult("TestResponse"); } }
Service Bus队列中的消息格式:
{ "data": { "data": { "Test Message" }, "metadata": { "ScheduledEnqueueTime": "2023-12-07T20:01:38.677943+00:00" } }, "datacontenttype": "application/json; charset=utf-8", "id": "20ad432360", "pubsubname": "pubsub", "source": "ca", "specversion": "1.0", "time": "2023-12-07T19:56:38Z", "topic": "testTopic", "traceid": "004b7d4b-00", "traceparent": "00d4b-00", "tracestate": "", "type": "com.dapr.event.sent" }
需要解决消息被双重data封装的问题,同时实现正确的延迟调度。
解决方案
1. 双重data封装的原因
你错误地将Dapr Bindings组件的DaprBindingMessage作为DaprPubSubEvent的构造参数传入,DaprPubSubEvent会自动将传入的对象序列化为data字段,而DaprBindingMessage本身包含data和metadata属性,最终导致data字段双重嵌套。
2. 正确的延迟调度实现
Dapr的Pub/Sub组件针对Azure Service Bus,支持通过ScheduledEnqueueTime元数据实现延迟调度,无需使用DaprBindingMessage。直接给DaprPubSubEvent添加元数据即可,步骤如下:
- 直接创建
DaprPubSubEvent实例并传入消息内容 - 通过
DaprPubSubEvent.Metadata属性添加ScheduledEnqueueTime,注意格式需为ISO 8601标准字符串(用ToString("o")生成)
修改后的代码:
public static class TestFunc { [FunctionName("TestFunc")] public static async Task<IActionResult> Run( [HttpTrigger(AuthorizationLevel.Anonymous, "get", "post", Route = "TestFunc")] HttpRequest req, [DaprPublish(PubSubName = "testSub", Topic = "testTopic")] IAsyncCollector<DaprPubSubEvent> pubEvent, ILogger log) { log.LogInformation("HTTP trigger function processed a request."); // 创建Pub/Sub事件实例,直接传入消息内容 var delayedEvent = new DaprPubSubEvent("Test Message"); // 添加延迟调度元数据,使用ISO 8601格式 delayedEvent.Metadata.Add( "ScheduledEnqueueTime", DateTimeOffset.Now.AddMinutes(5).ToString("o") ); await pubEvent.AddAsync(delayedEvent); return new OkObjectResult("TestResponse"); } }
3. 验证结果
修改后,Service Bus中的消息格式将恢复正常,data字段直接为你的消息内容,同时ScheduledEnqueueTime元数据会被Azure Service Bus识别,实现消息延迟5分钟后投递。
内容的提问来源于stack exchange,提问作者Piyush
相关产品推荐
相关产品推荐

