如何优化加速带数据库调用的异步数据同步方法
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
相关产品推荐
相关产品推荐

