2025-03-28 17:38:34 +08:00
|
|
|
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<String> dataDumpSpecial(DataDumpSpecialParam param) throws Throwable {
|
|
|
|
|
PartnerConnection connect = salesforceConnect.createConnect();
|
|
|
|
|
|
|
|
|
|
if (StringUtils.isBlank(param.getApi())) {
|
|
|
|
|
// dataObject查询
|
|
|
|
|
QueryWrapper<DataObject> 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<DataObject> list = dataObjectService.list(qw);
|
|
|
|
|
if (CollectionUtils.isEmpty(list)) {
|
|
|
|
|
break;
|
|
|
|
|
}
|
|
|
|
|
List<Future<?>> 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<String> ids = Sets.newHashSet();
|
|
|
|
|
if (param.getApi().contains(",")) {
|
|
|
|
|
ids.addAll(DataUtil.toIdList(param.getApi()));
|
|
|
|
|
} else {
|
|
|
|
|
ids.add(param.getApi());
|
|
|
|
|
}
|
|
|
|
|
List<Future<?>> 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<String, Object> map = Maps.newHashMap();
|
|
|
|
|
SalesforceParam salesforceParam = new SalesforceParam();
|
|
|
|
|
salesforceParam.setApi(api);
|
|
|
|
|
salesforceParam.setIdField(param.getField());
|
|
|
|
|
String maxId = null;
|
|
|
|
|
Field[] dsrFields;
|
|
|
|
|
List<String> 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);
|
2025-08-19 11:28:14 +08:00
|
|
|
EmailUtil.send("DataDump ERROR", format);
|
|
|
|
|
log.error("dataDumpSpecial error ", throwable);
|
|
|
|
|
throw new RuntimeException(throwable);
|
|
|
|
|
}
|
|
|
|
|
}, 0, 0);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Override
|
|
|
|
|
public Future<?> getDataNew(DataDumpSpecialParam param, PartnerConnection connect,DataObject dataObject) {
|
|
|
|
|
String api = param.getApi();
|
|
|
|
|
return salesforceExecutor.execute(() -> {
|
|
|
|
|
try {
|
|
|
|
|
Map<String, Object> map = Maps.newHashMap();
|
|
|
|
|
SalesforceParam salesforceParam = new SalesforceParam();
|
|
|
|
|
salesforceParam.setApi(api);
|
|
|
|
|
salesforceParam.setIdField(param.getField());
|
|
|
|
|
String maxId = null;
|
|
|
|
|
Field[] dsrFields;
|
|
|
|
|
List<String> 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);
|
|
|
|
|
salesforceParam.setBeginModifyDate(dataObject.getLastUpdateDate());
|
|
|
|
|
map.put("param", salesforceParam);
|
|
|
|
|
String sql = SqlUtil.showSql("com.celnet.datadump.mapper.SalesforceMapper.listOrderByIdNew", 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);
|
|
|
|
|
dataObject.setLastUpdateDate(new Date());
|
|
|
|
|
dataObjectService.updateById(dataObject);
|
|
|
|
|
} catch (Throwable throwable) {
|
|
|
|
|
String format = String.format("数据特殊表迁移 error, api name: %s, \nparam: %s, \ncause:\n%s", api, JSON.toJSONString(param), throwable);
|
2025-03-28 17:38:34 +08:00
|
|
|
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<DataBatch> 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);
|
|
|
|
|
}
|
|
|
|
|
}
|