Azure Digital Twins子孪生更新后父孪生数据延迟问题求助
问题描述
无C#开发经验,参照Azure官方教程及示例代码实现Azure Digital Twins孪生间更新逻辑:子孪生(如智能插座)更新时,父聚合孪生汇总子孪生数据。但出现数据延迟问题:第一次子孪生更新后父孪生无数据,第二次更新时父孪生显示第一次的汇总结果,始终慢一个迭代。
相关代码
// Default URL for triggering event grid function in the local environment. // http://localhost:7071/runtime/webhooks/EventGrid?functionName={functionname} using IoTHubTrigger = Microsoft.Azure.WebJobs.EventHubTriggerAttribute; using Azure; using Azure.Core.Pipeline; using Azure.DigitalTwins.Core; using Azure.Identity; using Microsoft.Azure.EventGrid.Models; using Microsoft.Azure.WebJobs; using Microsoft.Azure.WebJobs.Extensions.EventGrid; using Microsoft.Extensions.Logging; using Newtonsoft.Json; using Newtonsoft.Json.Linq; using System; using System.Net.Http; using System.Linq; using System.Reflection.Metadata.Ecma335; using System.Threading.Tasks; using System.Collections.Generic; using TwinUpdatesSample.Dto; namespace TwinUpdatesSample { public class ProcessDTRoutedData { private static HttpClient _httpClient = new HttpClient(); private static string _adtServiceUrl = Environment.GetEnvironmentVariable("ADT_SERVICE_URL"); /// <summary> /// 此函数的作用是根据楼层内的房间数据计算平均温度和湿度值。 /// /// 1) 获取房间的入站关系,从而得到楼层孪生ID /// 2) 获取楼层内所有房间的列表,并获取每个房间的湿度和温度属性 /// 3) 计算所有房间的平均温度和湿度 /// 4) 更新楼层孪生的温度和湿度属性 /// </summary> /// <param name="eventGridEvent"></param> /// <param name="log"></param> /// <returns></returns> [FunctionName("ProcessDTRoutedData")] public async Task Run([EventGridTrigger] EventGridEvent eventGridEvent, ILogger log) { log.LogInformation("ProcessDTRoutedData (Start)..."); DigitalTwinsClient client; DefaultAzureCredential credentials; // 如果未设置Azure Digital Twins服务URL,记录错误并退出方法 if (_adtServiceUrl == null) { log.LogError("未设置应用程序配置项\"ADT_SERVICE_URL\""); return; } try { // 认证Azure Digital Twins credentials = new DefaultAzureCredential(); client = new DigitalTwinsClient(new Uri(_adtServiceUrl), credentials, new DigitalTwinsClientOptions { Transport = new HttpClientTransport(_httpClient) }); } catch (Exception ex) { log.LogError($"异常信息: {ex.Message}"); client = null; credentials = null; return; } if (client != null) { if (eventGridEvent != null && eventGridEvent.Data != null) { JObject message = (JObject)JsonConvert.DeserializeObject(eventGridEvent.Data.ToString()); log.LogInformation("检测到Helix Building数字孪生的状态变更"); //log.LogInformation($"消息内容: {message}"); string twinId = eventGridEvent.Subject.ToString(); log.LogInformation($"孪生ID: {twinId}"); string modelId = message["data"]["modelId"].ToString(); log.LogInformation($"模型ID: {modelId}"); string smartPlugAggregatorId = null; if (modelId.Contains("dtmi:digitaltwins:rec_3_3:core:logicalDevice:SmartPlug;1")) { log.LogInformation($"正在从{twinId}的新状态更新ProjectSmartPlug状态"); // 智能插座孪生应关联到智能插座聚合孪生;获取该智能插座对应的聚合孪生ID AsyncPageable<IncomingRelationship> smartPlugAggregatorList = client.GetIncomingRelationshipsAsync(twinId); // 获取源ID(父孪生ID) await foreach (IncomingRelationship smartPlugAggregator in smartPlugAggregatorList) if (smartPlugAggregator.RelationshipName == "observes") { smartPlugAggregatorId = smartPlugAggregator.SourceId; } log.LogInformation($"本次迭代中,{smartPlugAggregatorId}观测到{twinId}的状态变更"); // 如果父孪生ID为空,说明出现异常 if (string.IsNullOrEmpty(smartPlugAggregatorId)) { log.LogError($"调用GetIncomingRelationships({twinId})未获取到observes关系的SourceID,此情况不应发生。"); return; } AsyncPageable<BasicDigitalTwin> queryResponse = client.QueryAsync<BasicDigitalTwin>($"SELECT smartPlug FROM digitaltwins smartPlugAggregator JOIN smartPlug RELATED smartPlugAggregator.observes WHERE smartPlugAggregator.$dtId = '{smartPlugAggregatorId}'"); List<SmartPlug> SmartPlugList = new List<SmartPlug>(); // 遍历所有智能插座孪生,构建列表 await foreach(BasicDigitalTwin twin in queryResponse) { JObject smartPlugPayload = (JObject)JsonConvert.DeserializeObject(twin.Contents["smartPlug"].ToString()); log.LogInformation($"智能插座{twin.Id}的 payload: {smartPlugPayload}"); SmartPlugList.Add(new SmartPlug() { id = twin.Id, ActiveEnergyWh = Convert.ToDouble(smartPlugPayload["ActiveEnergyWh"]), ActivePowerW = Convert.ToDouble(smartPlugPayload["ActivePowerW"]), ReActiveEnergyVARh = Convert.ToDouble(smartPlugPayload["ReActiveEnergyVARh"]), ReActivePowerVAR = Convert.ToDouble(smartPlugPayload["ReActivePowerVAR"]), }); } // 如果没有智能插座数据,说明出现异常并退出 if (SmartPlugList.Count < 1) { log.LogError($"聚合孪生({smartPlugAggregatorId})对应的智能插座列表为空,此情况不应发生。"); return; } // 计算所有智能插座数据的总和 double sumActiveEnergyWh = SmartPlugList.Sum(x => x.ActiveEnergyWh); log.LogInformation($"ActiveEnergyWh总和 : {sumActiveEnergyWh.ToString()}"); double sumActivePowerW = SmartPlugList.Sum(x => x.ActivePowerW); log.LogInformation($"ActivePowerW总和 : {sumActivePowerW.ToString()}"); double sumReActiveEnergyVARh = SmartPlugList.Sum(x => x.ReActiveEnergyVARh); log.LogInformation($"ReActiveEnergyVARh总和 : {sumReActiveEnergyVARh.ToString()}"); double sumReActivePowerVAR = SmartPlugList.Sum(x => x.ReActivePowerVAR); log.LogInformation($"ReActivePowerVAR总和 : {sumReActivePowerVAR.ToString()}"); var updateTwinData = new JsonPatchDocument(); updateTwinData.AppendReplace("/ActiveEnergyWh", Math.Round(sumActiveEnergyWh, 2)); updateTwinData.AppendReplace("/ActivePowerW", Math.Round(sumActivePowerW, 2)); try { log.LogInformation(updateTwinData.ToString()); await client.UpdateDigitalTwinAsync(smartPlugAggregatorId, updateTwinData); log.LogInformation("ProcessDTRoutedData (Done)..."); log.LogInformation(" "); } catch (Exception ex) { log.LogError($"更新失败: {ex.Message}"); } return; } } } } } }
问题原因与解决方法
核心原因
出现延迟的关键原因是ADT事件触发与孪生数据更新之间存在一致性延迟:当子孪生的更新事件触发函数时,ADT的底层存储可能还未完成子孪生数据的持久化,此时执行查询获取子孪生数据,拿到的是更新前的旧数据,导致计算结果滞后一个迭代。
另外,代码中存在一处语法小错误:List<SmartPlug> SmartPlugList = new List<SmartPlug>;; 多了一个分号,需修正。
解决步骤
等待数据一致性(推荐)
在查询子孪生数据前,加入短暂延迟(500ms-1s),给ADT足够时间完成数据持久化。可在获取smartPlugAggregatorId后添加:await Task.Delay(500); // 等待ADT完成数据同步直接使用事件中的最新数据
Event Grid消息中已包含子孪生的最新数据,无需再查询ADT获取,直接从事件数据中提取最新值,结合其他子孪生数据计算总和,避免查询旧数据:- 从
eventGridEvent.Data中提取当前子孪生的ActiveEnergyWh、ActivePowerW等最新值 - 结合父孪生当前汇总值或其他子孪生数据完成计算
- 从
修正语法错误
将List<SmartPlug> SmartPlugList = new List<SmartPlug>;;改为List<SmartPlug> SmartPlugList = new List<SmartPlug>();优化查询逻辑
若ADT支持,可在查询语句中添加WITH (NO_CACHE)参数,强制获取最新数据,避免查询缓存导致的旧数据问题。
内容的提问来源于stack exchange,提问作者Ricardo de Castro

