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

Azure Digital Twins子孪生更新后父孪生数据延迟问题求助

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>;; 多了一个分号,需修正。

解决步骤

  1. 等待数据一致性(推荐)
    在查询子孪生数据前,加入短暂延迟(500ms-1s),给ADT足够时间完成数据持久化。可在获取smartPlugAggregatorId后添加:

    await Task.Delay(500); // 等待ADT完成数据同步
    
  2. 直接使用事件中的最新数据
    Event Grid消息中已包含子孪生的最新数据,无需再查询ADT获取,直接从事件数据中提取最新值,结合其他子孪生数据计算总和,避免查询旧数据:

    • 从eventGridEvent.Data中提取当前子孪生的ActiveEnergyWh、ActivePowerW等最新值
    • 结合父孪生当前汇总值或其他子孪生数据完成计算
  3. 修正语法错误
    将List<SmartPlug> SmartPlugList = new List<SmartPlug>;;改为List<SmartPlug> SmartPlugList = new List<SmartPlug>();

  4. 优化查询逻辑
    若ADT支持,可在查询语句中添加WITH (NO_CACHE)参数,强制获取最新数据,避免查询缓存导致的旧数据问题。

内容的提问来源于stack exchange,提问作者Ricardo de Castro

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 09:33:15