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

如何优化加速带数据库调用的异步数据同步方法

EF工单同步流程性能优化咨询

大家好,提前感谢各位的帮助,我是相关技术的初学者,若问题存在表述不当还请谅解。
我基于Entity Framework编写了一个数据同步方法:从本地存储工单信息的数据库表(表内至少16.5万行数据)读取数据,再调用外部接口完成数据同步——数据不存在时发起POST请求新建,数据已存在且有变更时发起PATCH请求更新。
我想了解是否有可行方案(例如多线程、并行处理等)优化提速整个同步流程,是否需要在for循环逻辑中应用并行处理?
当前方法的实现代码如下:

public async Task<List<dynamic>> SyncOrdersTaskAsync(int PageSize)
{
    int PageIndex = 0;

    if (PageSize <= 0) PageSize = 100;

    const string phrase = "The fields order, task_code must make a unique set";

    var sorting = new SortingCriteria {
        Properties = new string[] { "WkOpenDate ASC" } };

    List<dynamic> listTest = new List<dynamic>();

    using (var uow = this.Factory.BeginUnitOfWork())
    {
        var repo = uow.GetRepository<IWorkOrderRepository>();

        var count = await repo.CountAllAsync();

        count = 150;

        for (PageIndex = 0; PageIndex <= count / PageSize; PageIndex++)
        {
            var paging = new PagingCriteria
            {
                PageIndex = PageIndex,
                PageSize = PageSize
            };

            var rows = await repo.GetByCriteriaAsync(
                "new {WkID, CompanyID, JobNo, JobTaskNo ,WkNumber, WkYear," +
                "WkYard,WkCustomerID,CuName,WkDivisionID,DvName,BusinessUnit," +
                "BusinessUnitManagerID,BusinessUnitManager,WkWorkTypeID,WtName," +
                "WkActivityID,WkActivityDescription,NoteDescrLavoro,WkWOManagerID," +
                "ProjectManager,IDMaster,ProjectCoordinator,WkOpenDate," +
                "WkDataChiusa,Prov,CodiceSito,CodiceOffice,CodiceLavorazione," +
                "CodiceNodo,DescrizioneNodo,WkPrevisionalStartDate,WkRealStartDate," +
                "WkPrevisionalEndDate,WkRealEndDate,NumeroOrdine," +
                "WkPrevisionalLabourAmount,TotaleCosti,SumOvertimeHours," +
                "SumTravelHours,SumNormalHours,WkProgressPercentage,Stato,CUP,CIG," +
                "TotaleManodopera,TotalePrestazioni,TotaleNoli,TotaleMateriali," +
                "SumAuxiliaryHours,TipoCommessa,TotaleOrdine, WkPreventivoData," +
                "WkConsuntivoData,TotaleFatturato,AggregateTotaleFatturato," +
                "AggregateTotalePrestazioni,Contract,CustomerOrderNumber," +
                "XmeWBECode,LastUpdateDate,PreGestWkID,CommercialNotes,Mandant," +
                "GammaProjectName,WkInventoryDate,WkCloseFlag,WkNote," +
                "TotalRegisteredLabour,TotalRegisteredPerformances," +
                "TotalRegisteredLeasings,TotalRegisteredMaterials,FlagFinalBalance," +
                "FinalBalance,OrderDate,TotalOrderDivision,SearchDescription," +
                "TotaleBefToBeApproved,TotaleBefToBeApprovedLeasings," +
                "TotaleLabourToBeApproved,AggregateLevel, AggregateTotalLabour," +
                "AggregateTotalLeasings,AggregateTotalMaterials," +
                "AggregateTotalRegisteredLabour," +
                "AggregateTotalRegisteredPerformances," +
                "AggregateTotalRegisteredLeasings," +
                "AggregateTotalRegisteredMaterials," +
                "AggregateTotalCost,AggregateSumNormalHours," +
                "AggregateSumAuxiliaryHours,AggregateSumRainHours," +
                "AggregateSumTravelHours,AggregateSumOvertimeHours," +
                "AggregateWkPrevisionalLabourAmount,AggregateFinalBalance," +
                "AggregateTotalOrder,AggregateTotalOrderDivision," +
                "AggregateTotalBefToBeApproved," +
                "AggregateTotalBefToBeApprovedLeasings," +
                "AggregateTotalLabourToBeApproved,TotalProduction," +
                "AggregateTotalProduction,JobTaskDescription}", paging, sorting);

            String url = appSettings.Value.UrlV1 + "order_tasks/";

            using (var httpClient = new HttpClient())
            {
                httpClient.DefaultRequestHeaders.Add("Authorization", "Token " +
                    await this.GetApiKey(true));
                if (rows.Count() > 0)
                {
                    foreach (var row in rows)
                    {
                        var testWork = (Model.WorkOrderCompleteInfo)Mapper
                            .MapWkOrdersCompleteInfo(row);
                        var orderIdDiv = await this.GetOrderForSyncing(httpClient,
                            testWork.JobNo);
                        var jsonTest = new JObject();
                        jsonTest["task_code"] = testWork.JobTaskNo;
                        jsonTest["description"] = testWork.JobTaskDescription;
                        jsonTest["order"] = orderIdDivitel.Id;
                        jsonTest["order_date"] = testWork.OrderDate.HasValue
                            ? testWork.OrderDate.Value.ToString("yyyy-MM-dd")
                            : string.IsNullOrEmpty(testWork.OrderDate.ToString())
                                ? "1970-01-01"
                                : testWork.OrderDate.ToString().Substring(0, 10);
                        jsonTest["progress"] = testWork.WkProgressPercentage;

                        var content = new StringContent(jsonTest.ToString(),
                            Encoding.UTF8, "application/json");
                        var result = await httpClient.PostAsync(url, content);
                        if (result.Content != null)
                        {
                            var responseContent = await result.Content
                                .ReadAsStringAsync();
                            bool alreadyExists = phrase.All(responseContent.Contains);

                            if (alreadyExists)
                            {
                                var taskCase = await GetTaskForSyncing(httpClient,
                                    testWork.JobTaskNo, orderIdDiv.Id.ToString());
                                var idCase = taskCase.Id;
                                String urlPatch = appSettings.Value.UrlV1 +
                                    "order_tasks/" + idCase + "/";
                                bool isSame = taskCase.Equals(testWork
                                    .toSolOrderTask());
                                if (!isSame)
                                {
                                    var resultPatch = await httpClient.PatchAsync(
                                        urlPatch, content);
                                    if (resultPatch != null)
                                    {
                                        var responsePatchContent = await resultPatch
                                            .Content.ReadAsStringAsync();
                                        var jsonPatchContent = JsonConvert
                                            .DeserializeObject<dynamic>(
                                            responsePatchContent);
                                        listTest.Add(jsonPatchContent);
                                    }
                                }
                                else
                                {
                                    listTest.Add(taskCase.JobTaskNo_ +
                                        " is already updated!");
                                }
                            }
                            else
                            {
                                var jsonContent = JsonConvert
                                    .DeserializeObject<dynamic>(responseContent);
                                listTest.Add(jsonContent);
                            }
                        }
                    }
                }
            }
        }

        return listTest;
    }
}

再次提前感谢各位的解答,希望我的问题表述清晰。


回答

先改你现有代码里的几个基础问题,这部分改完性能就能提3倍以上,再考虑并行优化:

  • 不要在循环里每次new HttpClient,会触发套接字耗尽,每次请求都要重新建立TCP连接、做TLS握手,开销极大。全局单例复用HttpClient,或者用IHttpClientFactory创建实例。另外鉴权Token提前缓存,不要每个分页循环都调用一次GetApiKey拿新Token。
  • 删掉硬编码的count = 150;,另外注意代码里的变量笔误:查询拿到的是orderIdDiv,后面赋值用的是orderIdDivitel.Id,运行会直接抛空引用错误。
  • 把“先发POST靠接口返回唯一键冲突判断数据存在”的逻辑换掉。你现在每一条已存在的数据都要多跑一次失败POST、一次单条查询请求,16.5万条数据平白多了几十万次无效HTTP请求。正确做法是同步启动前先把接口侧已有的工单任务按唯一键(order、task_code)拉到本地存成字典,处理时直接查字典判断是新增还是更新,不需要每条都撞错再回查。
  • 分页不要用传统的Skip/Take页码分页,16万条数据越往后翻页越慢,改成基于排序字段WkOpenDate+主键WkID的游标分页,每次查完记录最后一条的这两个字段值,下一页直接查大于该值的数据,全程分页查询耗时稳定在毫秒级。

再谈并行处理的正确方案,不要直接套Parallel.ForEach或者无脑Task.WhenAll:

  • 你的场景是IO密集型场景,90%以上的时间都在等数据库和接口的网络响应,并行收益极高,但一定要做并发限流,不然一下把几万请求打出去,外部接口会直接触发流控返回429,甚至封你的IP。用SemaphoreSlim控制总并发数,初始值从5开始测,根据对方接口的QPS上限慢慢往上调,一般外部开放接口设置10-20并发就足够。
  • 用生产者消费者模式解耦读库和发请求的逻辑:一个独立任务负责按游标分页从EF读数据,映射完直接丢到固定容量(比如1000条)的内存队列里,不需要等当前页数据处理完再读下一页,避免数据库连接空等。
  • 启动N个(和并发数一致)消费任务从队列取数据处理,每个任务独立完成接口比对、POST/PATCH请求逻辑,处理结果存在线程安全的ConcurrentBag<dynamic>里,不要用普通List,多线程写入会出数据错乱。
  • EF的DbContext不是线程安全的,不要把DbContext的生命周期拉长到并行处理环节,读出来的数据映射完就可以释放对应的EF跟踪状态,不要在并行任务里操作DbContext。
  • 加临时异常重试逻辑,遇到网络波动、接口5xx错误、429流控时,用指数退避策略重试2-3次,不要因为单条请求失败中断整个16万条的同步任务。另外加断点记录,每处理完100条就把最后处理成功的游标位置存下来,程序崩溃重启后不需要从头开始同步。

注意:不要给数据库读操作开并行,单实例EF上下文不支持并行查询,而且单线程游标分页的读取速度完全能跟上20并发的接口处理速度,并行读只会额外增加数据库压力,没有收益。


内容的提问来源于stack exchange,提问作者Frank DG

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 17:45:38