diff --git a/src/main/java/com/celnet/datadump/job/DataDumpNewJob.java b/src/main/java/com/celnet/datadump/job/DataDumpNewJob.java index ff38dd9..d0f07eb 100644 --- a/src/main/java/com/celnet/datadump/job/DataDumpNewJob.java +++ b/src/main/java/com/celnet/datadump/job/DataDumpNewJob.java @@ -191,7 +191,7 @@ public class DataDumpNewJob { return new ReturnT<>(500, "参数解析失败!"); } - return commonService.updateLinkType(param); + return dataImportNewService.updateLinkTypeJob(param); } /** diff --git a/src/main/java/com/celnet/datadump/service/DataImportNewService.java b/src/main/java/com/celnet/datadump/service/DataImportNewService.java index bbf75fd..1946b17 100644 --- a/src/main/java/com/celnet/datadump/service/DataImportNewService.java +++ b/src/main/java/com/celnet/datadump/service/DataImportNewService.java @@ -39,5 +39,11 @@ public interface DataImportNewService { ReturnT uploadFileNew(SalesforceParam param, DataObject dataObject) throws Exception; - + /** + * 更新关联类型 + * @param param 参数 + * @return ReturnT + * @throws Exception exception + */ + ReturnT updateLinkTypeJob(SalesforceParam param) throws Exception; } diff --git a/src/main/java/com/celnet/datadump/service/impl/DataImportBatchServiceImpl.java b/src/main/java/com/celnet/datadump/service/impl/DataImportBatchServiceImpl.java index d0689a0..3d8c113 100644 --- a/src/main/java/com/celnet/datadump/service/impl/DataImportBatchServiceImpl.java +++ b/src/main/java/com/celnet/datadump/service/impl/DataImportBatchServiceImpl.java @@ -517,7 +517,7 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { } /** - * 【单彪】组装Update执行参数 + * 【单表】组装Update执行参数 */ public ReturnT updateSfDataBatch(SalesforceParam param, List> futures) throws Exception { List apis; diff --git a/src/main/java/com/celnet/datadump/service/impl/DataImportNewServiceImpl.java b/src/main/java/com/celnet/datadump/service/impl/DataImportNewServiceImpl.java index 0fb1f06..04f618f 100644 --- a/src/main/java/com/celnet/datadump/service/impl/DataImportNewServiceImpl.java +++ b/src/main/java/com/celnet/datadump/service/impl/DataImportNewServiceImpl.java @@ -1,5 +1,6 @@ package com.celnet.datadump.service.impl; +import cn.hutool.core.lang.UUID; import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.JSONArray; import com.alibaba.fastjson2.JSONObject; @@ -16,12 +17,13 @@ import com.celnet.datadump.mapper.CustomMapper; import com.celnet.datadump.param.DataDumpParam; import com.celnet.datadump.param.SalesforceParam; import com.celnet.datadump.service.*; -import com.celnet.datadump.util.DataUtil; -import com.celnet.datadump.util.EmailUtil; -import com.celnet.datadump.util.HttpUtil; -import com.celnet.datadump.util.OssUtil; +import com.celnet.datadump.util.*; import com.google.common.collect.Lists; import com.google.common.collect.Maps; +import com.sforce.async.BatchInfo; +import com.sforce.async.BulkConnection; +import com.sforce.async.JobInfo; +import com.sforce.async.OperationEnum; import com.sforce.soap.partner.*; import com.sforce.soap.partner.sobject.SObject; import com.xxl.job.core.biz.model.ReturnT; @@ -1665,4 +1667,240 @@ public class DataImportNewServiceImpl implements DataImportNewService { } } + + /** + * updateLinkType入口 + */ + @Override + public ReturnT updateLinkTypeJob(SalesforceParam param) throws Exception { + List> futures = Lists.newArrayList(); + try { + if (param.getType() == 1){ + return updateLinkType(param, futures); + }else { + return updateLinkTypeBatch(param,futures); + } + } catch (InterruptedException e) { + throw e; + } catch (Throwable throwable) { + salesforceExecutor.remove(futures.toArray(new Future[]{})); + log.error("updateLinkTypeJob error", throwable); + throw throwable; + } + } + + /** + * 【batch逻辑】组装updateLinkType执行参数 + */ + public ReturnT updateLinkTypeBatch(SalesforceParam param, List> futures) throws Exception { + List apis; + if (StringUtils.isBlank(param.getApi())) { + apis = dataObjectService.list().stream().map(DataObject::getName).collect(Collectors.toList()); + } else { + apis = DataUtil.toIdList(param.getApi()); + } + String beginDateStr = null; + String endDateStr = null; + if (param.getBeginCreateDate() != null && param.getEndCreateDate() != null){ + Date beginDate = param.getBeginCreateDate(); + Date endDate = param.getEndCreateDate(); + beginDateStr = DateUtil.format(beginDate, "yyyy-MM-dd HH:mm:ss"); + endDateStr = DateUtil.format(endDate, "yyyy-MM-dd HH:mm:ss"); + } + QueryWrapper wrapper = new QueryWrapper<>(); + wrapper.in("api", apis); + Map resultMap = dataObjectService.list().stream() + .collect(Collectors.toMap( + DataObject::getKeyPrefix, + DataObject::getName, + (existing, replacement) -> replacement)); + + for (String api : apis) { + try { + List salesforceParams = null; + QueryWrapper dbQw = new QueryWrapper<>(); + dbQw.eq("name", api); + if (StringUtils.isNotEmpty(beginDateStr) && StringUtils.isNotEmpty(endDateStr)) { + dbQw.eq("sync_start_date", beginDateStr); // 等于开始时间 + dbQw.eq("sync_end_date", endDateStr); // 等于结束时间 + } + List list = dataBatchService.list(dbQw); + AtomicInteger batch = new AtomicInteger(1); + if (CollectionUtils.isNotEmpty(list)) { + salesforceParams = list.stream().map(t -> { + SalesforceParam salesforceParam = param.clone(); + salesforceParam.setApi(t.getName()); + salesforceParam.setBeginCreateDate(t.getSyncStartDate()); + salesforceParam.setEndCreateDate(t.getSyncEndDate()); + salesforceParam.setBatch(batch.getAndIncrement()); + return salesforceParam; + }).collect(Collectors.toList()); + } + + // 手动任务优先执行 + for (SalesforceParam salesforceParam : salesforceParams) { + Future future = salesforceExecutor.execute(() -> { + try { + updateLinkBatch(salesforceParam, resultMap); + } catch (Throwable throwable) { + log.error("salesforceExecutor error", throwable); + throw new RuntimeException(throwable); + } + }, salesforceParam.getBatch(), 1); + futures.add(future); + } + // 等待当前所有线程执行完成 + salesforceExecutor.waitForFutures(futures.toArray(new Future[]{})); + } catch (Exception e) { + throw e; + } + } + return null; + } + + /** + * 执行数据updateLinkType + */ + private void updateLinkBatch(SalesforceParam param, Map resultMap) throws Exception { + + QueryWrapper dbQw = new QueryWrapper<>(); + dbQw.eq("api", param.getApi()); + List linkConfigs = linkConfigService.list(dbQw); + + String beginDateStr = null; + String endDateStr = null; + Date beginDate = param.getBeginCreateDate(); + Date endDate = param.getEndCreateDate(); + if (param.getBeginCreateDate() != null && param.getEndCreateDate() != null){ + beginDateStr = DateUtil.format(beginDate, "yyyy-MM-dd HH:mm:ss"); + endDateStr = DateUtil.format(endDate, "yyyy-MM-dd HH:mm:ss"); + } + // 表内数据总量 + Integer count = customMapper.countBySQL(param.getApi(), "where CreatedDate >= '" + beginDateStr + "' and CreatedDate < '" + endDateStr + "'"); + + log.info("表api:{} 存在" +count+ "条数据!", param.getApi()); + + if (count >0 ) { + int page = count % 2000 == 0 ? count / 2000 : (count / 2000) + 1; + + for (int i = 0; i < page; i++) { + + log.info("表api:{},数据量:{},执行数据更新!", param.getApi(), 2000* (i+1)); + + List> mapList = customMapper.list("*", param.getApi(), "1=1 order by Id limit " + i * 2000 + ",2000"); + + List> updateMapList = new ArrayList<>(); + for (int j = 1; j <= mapList.size(); j++) { + Map map = mapList.get(j - 1); + for (LinkConfig config : linkConfigs) { + if (map.get(config.getField()) != null){ + String type = resultMap.get(map.get(config.getField()).toString().substring(0, 3)); + Map paramMap = Maps.newHashMap(); + paramMap.put("key", config.getLinkField()); + paramMap.put("value", type); + updateMapList.add(paramMap); + } + } + customMapper.updateById(param.getApi(), updateMapList, String.valueOf(mapList.get(j - 1).get("Id") != null?mapList.get(j - 1).get("Id") : mapList.get(j - 1).get("id"))); + + } + } + } + + } + + /** + * 组装updateLinkType执行参数 + */ + public ReturnT updateLinkType(SalesforceParam param, List> futures) throws Exception { + List apis; + if (StringUtils.isBlank(param.getApi())) { + apis = dataObjectService.list().stream().map(DataObject::getName).collect(Collectors.toList()); + } else { + apis = DataUtil.toIdList(param.getApi()); + } + + QueryWrapper wrapper = new QueryWrapper<>(); + wrapper.in("api", apis).eq("is_link",1); + Map resultMap = dataObjectService.list().stream() + .collect(Collectors.toMap( + DataObject::getKeyPrefix, + DataObject::getName, + (existing, replacement) -> replacement)); + + for (String api : apis) { + Future future = salesforceExecutor.execute(() -> { + try { + updateLink(api,param, resultMap); + } catch (Throwable throwable) { + log.error("salesforceExecutor error", throwable); + throw new RuntimeException(throwable); + } + }, 0, 1); + futures.add(future); + } + // 等待当前所有线程执行完成 + salesforceExecutor.waitForFutures(futures.toArray(new Future[]{})); + + return ReturnT.SUCCESS; + } + + + /** + * 执行数据updateLinkType + */ + private void updateLink(String api,SalesforceParam param, Map resultMap) throws Exception { + + QueryWrapper dbQw = new QueryWrapper<>(); + dbQw.eq("api", param.getApi()); + List linkConfigs = linkConfigService.list(dbQw); + + String beginDateStr = null; + Date beginDate = param.getBeginModifyDate(); + if (param.getBeginModifyDate() != null ){ + beginDateStr = DateUtil.format(beginDate, "yyyy-MM-dd HH:mm:ss"); + } + String sql1 = null; + String sql2 = null; + if (beginDateStr != null){ + sql1 = "where SystemModstamp >= '" + beginDateStr; + sql2 = "SystemModstamp >= "+ beginDateStr +" order by Id limit " ; + }else { + sql1 = null; + sql2 = "1=1 order by Id limit " ; + } + + // 表内数据总量 + Integer count = customMapper.countBySQL(api, sql1); + + log.info("表api:{} 存在" +count+ "条数据!", param.getApi()); + + if (count >0 ) { + int page = count % 2000 == 0 ? count / 2000 : (count / 2000) + 1; + + for (int i = 0; i < page; i++) { + + log.info("表api:{},数据量:{},执行数据更新!", api, 2000* (i+1)); + + List> mapList = customMapper.list("*", api, sql2 + i * 2000 + ",2000" ); + + List> updateMapList = new ArrayList<>(); + for (int j = 1; j <= mapList.size(); j++) { + Map map = mapList.get(j - 1); + for (LinkConfig config : linkConfigs) { + if (map.get(config.getField()) != null){ + String type = resultMap.get(map.get(config.getField()).toString().substring(0, 3)); + Map paramMap = Maps.newHashMap(); + paramMap.put("key", config.getLinkField()); + paramMap.put("value", type); + updateMapList.add(paramMap); + } + customMapper.updateById(api, updateMapList, String.valueOf(mapList.get(j - 1).get("Id") != null?mapList.get(j - 1).get("Id") : mapList.get(j - 1).get("id"))); + } + } + } + } + + } + }