package com.celnet.datadump.service.impl; import com.alibaba.fastjson2.JSON; import com.alibaba.fastjson2.JSONArray; import com.alibaba.fastjson2.JSONObject; import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper; import com.celnet.datadump.config.SalesforceConnect; import com.celnet.datadump.config.SalesforceExecutor; import com.celnet.datadump.entity.DataBatch; import com.celnet.datadump.entity.DataObject; import com.celnet.datadump.global.Const; import com.celnet.datadump.param.DataDumpSpecialParam; import com.celnet.datadump.param.SalesforceParam; import com.celnet.datadump.mapper.CustomMapper; import com.celnet.datadump.service.*; import com.celnet.datadump.util.DataUtil; import com.celnet.datadump.util.EmailUtil; import com.celnet.datadump.util.SqlUtil; import com.google.common.collect.Lists; import com.google.common.collect.Maps; import com.google.common.collect.Sets; import com.sforce.soap.partner.DescribeSObjectResult; import com.sforce.soap.partner.Field; import com.sforce.soap.partner.PartnerConnection; import com.sforce.soap.partner.QueryResult; import com.sforce.soap.partner.sobject.SObject; import com.xxl.job.core.biz.model.ReturnT; import com.xxl.job.core.log.XxlJobLogger; import lombok.extern.slf4j.Slf4j; import org.apache.commons.collections.CollectionUtils; import org.apache.commons.lang3.ObjectUtils; import org.apache.commons.lang3.StringUtils; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import java.util.Date; import java.util.List; import java.util.Map; import java.util.Set; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; /** * @author Red * @description * @date 2023/01/10 */ @Slf4j @Service public class DataDumpSpecialServiceImpl implements DataDumpSpecialService { @Autowired private CommonService commonService; @Autowired private DataObjectService dataObjectService; @Autowired private DataFieldService dataFieldService; @Autowired private DataBatchService dataBatchService; @Autowired private CustomMapper customMapper; @Autowired private SalesforceConnect salesforceConnect; @Autowired private SalesforceExecutor salesforceExecutor; @Override public ReturnT dataDumpSpecial(DataDumpSpecialParam param) throws Throwable { PartnerConnection connect = salesforceConnect.createConnect(); if (StringUtils.isBlank(param.getApi())) { // dataObject查询 QueryWrapper qw = new QueryWrapper<>(); qw.eq("data_work", 1) .eq("data_lock", 0) .orderByAsc("data_index") .last(" limit 10"); while (true) { TimeUnit.MILLISECONDS.sleep(1); List list = dataObjectService.list(qw); if (CollectionUtils.isEmpty(list)) { break; } List> futures = Lists.newArrayList(); try { for (DataObject dataObject : list) { Future future = salesforceExecutor.execute(() -> { DataObject update = new DataObject(); try { // 存在isDeleted 只查询IsDeleted为false的 if (dataFieldService.hasDeleted(param.getApi())) { param.setIsDeleted(false); } else { // 不存在 过滤 param.setIsDeleted(null); } update.setName(dataObject.getName()); String api = dataObject.getName(); log.info("dump apis: {}", api); XxlJobLogger.log("dump apis: {}", api); update.setDataLock(1); dataObjectService.updateById(update); param.setApi(api); salesforceExecutor.waitForFutures(getData(param, connect)); update.setDataWork(0); } catch (Throwable e) { throw new RuntimeException(e); } finally { if (StringUtils.isNotBlank(update.getName())) { update.setDataLock(0); dataObjectService.updateById(update); } } }, 0, 0); futures.add(future); } // 等待当前所有线程执行完成 salesforceExecutor.waitForFutures(futures.toArray(new Future[]{})); } finally { salesforceExecutor.remove(futures.toArray(new Future[]{})); } } } else { // 指定api Set ids = Sets.newHashSet(); if (param.getApi().contains(",")) { ids.addAll(DataUtil.toIdList(param.getApi())); } else { ids.add(param.getApi()); } List> futures = Lists.newArrayList(); try { for (String api : ids) { futures.add(getData(param, connect)); } // 等待当前所有线程执行完成 salesforceExecutor.waitForFutures(futures.toArray(new Future[]{})); } finally { salesforceExecutor.remove(futures.toArray(new Future[]{})); } } return ReturnT.SUCCESS; } @Override public Future getData(DataDumpSpecialParam param, PartnerConnection connect) { String api = param.getApi(); return salesforceExecutor.execute(() -> { try { // 检测表 不存在就生成 但不生成批次 commonService.checkApi(api, true); log.info("dataDumpSpecial api: {}, field:{}", api, param.getField()); Map map = Maps.newHashMap(); SalesforceParam salesforceParam = new SalesforceParam(); salesforceParam.setApi(api); salesforceParam.setIdField(param.getField()); String maxId = null; Field[] dsrFields; List fields = Lists.newArrayList(); // 获取sf字段 { DescribeSObjectResult dsr = connect.describeSObject(api); dsrFields = dsr.getFields(); for (Field field : dsrFields) { // 不查询文件 if ("base64".equalsIgnoreCase(field.getType().toString())) { continue; } fields.add(field.getName()); } salesforceParam.setSelect(StringUtils.join(fields, ",")); } JSONArray objects = null; int count = 0, failCount = 0; // 遍历 while (true) { try { // 判断是否存在要排除的id salesforceParam.setMaxId(maxId); map.put("param", salesforceParam); String sql = SqlUtil.showSql("com.celnet.datadump.mapper.SalesforceMapper.listOrderById", map); log.info("query sql: {}", sql); XxlJobLogger.log("query sql: {}", sql); QueryResult queryResult = connect.queryAll(sql); if (ObjectUtils.isEmpty(queryResult) || ObjectUtils.isEmpty(queryResult.getRecords())) { break; } SObject[] records = queryResult.getRecords(); objects = DataUtil.toJsonArray(records, dsrFields); maxId = ((JSONObject) objects.get(objects.size() - 1)).getString(Const.ID); // 存储更新 commonService.saveOrUpdate(api, fields, records, objects, true); count += records.length; TimeUnit.MILLISECONDS.sleep(1); log.info("dump success count: {}", count); XxlJobLogger.log("dump success count: {}", count); failCount = 0; objects = null; } catch (InterruptedException e) { throw e; } catch (Throwable throwable) { failCount++; log.error("dataDumpSpecial error api:{}, data:{}", api, JSON.toJSONString(objects), throwable); if (failCount > Const.MAX_FAIL_COUNT) { throwable.addSuppressed(new Exception("dataDumpSpecial error data:" + JSON.toJSONString(objects))); throw throwable; } TimeUnit.MINUTES.sleep(1); } } updateDataBatch(connect, api); } catch (Throwable throwable) { String format = String.format("数据特殊表迁移 error, api name: %s, \nparam: %s, \ncause:\n%s", api, JSON.toJSONString(param), throwable); EmailUtil.send("DataDump ERROR", format); log.error("dataDumpSpecial error ", throwable); throw new RuntimeException(throwable); } }, 0, 0); } /** * 更新batch * * @param connect connect * @param api 表名 */ private void updateDataBatch(PartnerConnection connect, String api) throws Exception { if (Const.BATCH_FILTERS.stream().anyMatch(t -> api.indexOf(t) > 0)) { return; } QueryWrapper qw = new QueryWrapper<>(); qw.eq("name", api); DataBatch dataBatch = dataBatchService.getOne(qw); DataBatch one = new DataBatch(); SalesforceParam salesforceParam = new SalesforceParam(); salesforceParam.setApi(api); // 存在isDeleted 只查询IsDeleted为false的 if (dataFieldService.hasDeleted(api)) { salesforceParam.setIsDeleted(false); } else { // 不存在 过滤 salesforceParam.setIsDeleted(null); } // sf count Integer sfNum = commonService.countSfNum(connect, salesforceParam); // db count Integer dbNum = customMapper.count(salesforceParam); Date now = new Date(); if (dataBatch != null) { one.setId(dataBatch.getId()); if (dataBatch.getFirstSyncDate() == null) { one.setFirstDbNum(dbNum); one.setFirstSfNum(sfNum); one.setFirstSyncDate(now); } } else { one.setFirstDbNum(dbNum); one.setFirstSfNum(sfNum); one.setFirstSyncDate(now); } one.setDbNum(dbNum); one.setSfNum(sfNum); one.setName(api); one.setSyncStatus(dbNum.equals(sfNum) ? 1 : 0); // 没创建时间 开始结束时间为空 dataBatchService.saveOrUpdate(one); } }