Compare commits
4 Commits
eea0ef28a6
...
e72501f17d
| Author | SHA1 | Date | |
|---|---|---|---|
| e72501f17d | |||
| ff7fff9f90 | |||
| 902e618459 | |||
| 40265c3c75 |
@ -138,6 +138,54 @@ public class DataDumpNewJob {
|
||||
return commonBatchService.dumpBigObject(param);
|
||||
}
|
||||
|
||||
/**
|
||||
* 源系统BulkV1查询
|
||||
*
|
||||
* @param paramStr 参数json
|
||||
* @return result
|
||||
*/
|
||||
@XxlJob("sourceSystemBulkV1QueryJob")
|
||||
public ReturnT<String> sourceSystemBulkV1QueryJob(String paramStr) throws Exception {
|
||||
log.info("sourceSystemBulkV1QueryJob execute start ..................");
|
||||
SalesforceParam param = new SalesforceParam();
|
||||
try {
|
||||
if (StringUtils.isNotBlank(paramStr)) {
|
||||
param = JSON.parseObject(paramStr, SalesforceParam.class);
|
||||
}
|
||||
} catch (Throwable throwable) {
|
||||
return new ReturnT<>(500, "参数解析失败!");
|
||||
}
|
||||
// 参数转换
|
||||
param.setBeginCreateDate(param.getBeginDate());
|
||||
param.setEndCreateDate(param.getEndDate());
|
||||
|
||||
return commonBatchService.dumpBigObjectSourceV1(param);
|
||||
}
|
||||
|
||||
/**
|
||||
* 目标系统BulkV1查询
|
||||
*
|
||||
* @param paramStr 参数json
|
||||
* @return result
|
||||
*/
|
||||
@XxlJob("targetSystemBulkV1QueryJob")
|
||||
public ReturnT<String> targetSystemBulkV1QueryJob(String paramStr) throws Exception {
|
||||
log.info("targetSystemBulkV1QueryJob execute start ..................");
|
||||
SalesforceParam param = new SalesforceParam();
|
||||
try {
|
||||
if (StringUtils.isNotBlank(paramStr)) {
|
||||
param = JSON.parseObject(paramStr, SalesforceParam.class);
|
||||
}
|
||||
} catch (Throwable throwable) {
|
||||
return new ReturnT<>(500, "参数解析失败!");
|
||||
}
|
||||
// 参数转换
|
||||
param.setBeginCreateDate(param.getBeginDate());
|
||||
param.setEndCreateDate(param.getEndDate());
|
||||
|
||||
return commonBatchService.dumpBigObjectTargetV1(param);
|
||||
}
|
||||
|
||||
/**
|
||||
* 存量任务(BigObject_V2)
|
||||
*
|
||||
|
||||
@ -123,4 +123,27 @@ public class DataVerifyJob {
|
||||
return dataVerifyService.verifyIncrementalQuality(param);
|
||||
}
|
||||
|
||||
/**
|
||||
* BigObject数据质量校验任务
|
||||
* 专门用于校验BigObject类型的数据质量
|
||||
*
|
||||
* @param paramStr 参数json
|
||||
* @return result
|
||||
*/
|
||||
@XxlJob("dataVerifyBigObjectQualityJob")
|
||||
public ReturnT<String> dataVerifyBigObjectQualityJob(String paramStr) throws Exception {
|
||||
log.info("dataVerifyBigObjectQualityJob execute start ..................");
|
||||
DataVerifyParam param = new DataVerifyParam();
|
||||
try {
|
||||
if (StringUtils.isNotBlank(paramStr)) {
|
||||
param = JSON.parseObject(paramStr, DataVerifyParam.class);
|
||||
}
|
||||
} catch (Throwable throwable) {
|
||||
return new ReturnT<>(500, "参数解析失败!");
|
||||
}
|
||||
|
||||
// 调用BigObject数据质量校验服务方法
|
||||
return dataVerifyService.verifyBigObjectQuality(param);
|
||||
}
|
||||
|
||||
}
|
||||
@ -59,6 +59,12 @@ public class DataVerifyParam {
|
||||
@ApiModelProperty(value = "忽略开始时间", hidden = true)
|
||||
private Date ignoreBeginDate;
|
||||
|
||||
/**
|
||||
* 校验字段 多个英文逗号分割
|
||||
*/
|
||||
@ApiModelProperty(value = "校验字段 多个英文逗号分割")
|
||||
private String fields;
|
||||
|
||||
|
||||
|
||||
}
|
||||
|
||||
@ -11,7 +11,10 @@ public interface CommonBatchService {
|
||||
|
||||
ReturnT<String> dumpBigObjectV2(SalesforceParam param) throws Exception;
|
||||
|
||||
|
||||
ReturnT<String> importBigObject(SalesforceParam param) throws Exception;
|
||||
|
||||
ReturnT<String> dumpBigObjectSourceV1(SalesforceParam param) throws Exception;
|
||||
|
||||
ReturnT<String> dumpBigObjectTargetV1(SalesforceParam param) throws Exception;
|
||||
|
||||
}
|
||||
@ -27,5 +27,12 @@ public interface DataVerifyService {
|
||||
* @return ReturnT
|
||||
*/
|
||||
ReturnT<String> verifyIncrementalQuality(DataVerifyParam param) throws Exception;
|
||||
|
||||
/**
|
||||
* BigObject数据质量校验
|
||||
* @param param 参数
|
||||
* @return ReturnT
|
||||
*/
|
||||
ReturnT<String> verifyBigObjectQuality(DataVerifyParam param) throws Exception;
|
||||
|
||||
}
|
||||
@ -182,6 +182,36 @@ public class CommonBatchServiceImpl implements CommonBatchService {
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public ReturnT<String> dumpBigObjectSourceV1(SalesforceParam param) throws Exception {
|
||||
try {
|
||||
// 手动任务
|
||||
ReturnT<String> result = manualBigObjectSourceV1Dump(param);
|
||||
if (result != null) {
|
||||
return result;
|
||||
}
|
||||
} catch (Throwable throwable) {
|
||||
log.error("dump error", throwable);
|
||||
throw throwable;
|
||||
}
|
||||
return ReturnT.SUCCESS;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ReturnT<String> dumpBigObjectTargetV1(SalesforceParam param) throws Exception {
|
||||
try {
|
||||
// 手动任务
|
||||
ReturnT<String> result = manualBigObjectTargetV1Dump(param);
|
||||
if (result != null) {
|
||||
return result;
|
||||
}
|
||||
} catch (Throwable throwable) {
|
||||
log.error("dump error", throwable);
|
||||
throw throwable;
|
||||
}
|
||||
return ReturnT.SUCCESS;
|
||||
}
|
||||
|
||||
private ReturnT<String> manualBigObjectDump(SalesforceParam param) throws Exception {
|
||||
|
||||
String api = param.getApi();
|
||||
@ -199,7 +229,7 @@ public class CommonBatchServiceImpl implements CommonBatchService {
|
||||
.filter(field -> !"SystemModstamp".equals(field) && !"LastModifiedDate".equals(field))
|
||||
.collect(Collectors.joining(", "));
|
||||
|
||||
String sql = "select Id__c,OppId__c from " + param.getApi();
|
||||
String sql = "select " + fieldStr + " from " + param.getApi() ;
|
||||
log.info("构建查询SQL: {}", sql);
|
||||
|
||||
job = BulkUtil.createJob(bulkConnect, api, OperationEnum.queryAll);
|
||||
@ -326,6 +356,128 @@ public class CommonBatchServiceImpl implements CommonBatchService {
|
||||
return ReturnT.SUCCESS;
|
||||
}
|
||||
|
||||
/**
|
||||
* 源系统BulkV1查询实现
|
||||
*
|
||||
* @param param 参数
|
||||
* @return 执行结果
|
||||
* @throws Exception 异常
|
||||
*/
|
||||
private ReturnT<String> manualBigObjectSourceV1Dump(SalesforceParam param) throws Exception {
|
||||
|
||||
String api = param.getApi();
|
||||
|
||||
String sql = param.getSql();
|
||||
|
||||
Integer type = param.getType();
|
||||
|
||||
BulkConnection bulkConnect = salesforceConnect.createBulkConnect();
|
||||
|
||||
JobInfo job = null;
|
||||
|
||||
try {
|
||||
|
||||
log.info("开始源系统BulkV1查询BigObject数据, API: {}, SQL: {}", api,sql);
|
||||
|
||||
// 创建Bulk作业
|
||||
job = BulkUtil.createJob(bulkConnect, api, OperationEnum.queryAll);
|
||||
log.info("创建Bulk V1查询作业成功, 作业ID: {}", job.getId());
|
||||
|
||||
// 提交查询批次
|
||||
BatchInfo batchInfo = bulkConnect.createBatchFromStream(job, new ByteArrayInputStream(sql.getBytes(StandardCharsets.UTF_8)));
|
||||
log.info("创建查询批次成功, 批次ID: {}", batchInfo.getId());
|
||||
|
||||
if (type == 0){
|
||||
return ReturnT.SUCCESS;
|
||||
}
|
||||
|
||||
// 等待作业完成
|
||||
int completionOne = BulkUtil.awaitCompletionOne(bulkConnect, job, batchInfo);
|
||||
log.info("查询批次处理完成, 处理记录数: {}", completionOne);
|
||||
|
||||
// 如果没有数据,直接结束
|
||||
if (completionOne == 0){
|
||||
log.info("无更多数据, 结束处理");
|
||||
BulkUtil.closeJob(bulkConnect, job.getId());
|
||||
return ReturnT.SUCCESS;
|
||||
}
|
||||
|
||||
// 获取查询结果列表
|
||||
QueryResultList queryResultList = bulkConnect.getQueryResultList(job.getId(), batchInfo.getId());
|
||||
log.info("获取查询结果列表, 结果数量: {}", queryResultList.getResult().length);
|
||||
|
||||
BulkUtil.closeJob(bulkConnect, job.getId());
|
||||
log.info("关闭作业成功, 作业ID: {}", job.getId());
|
||||
|
||||
}catch (Exception e) {
|
||||
log.error("源系统BulkV1查询BigObject数据失败, API: {}", api, e);
|
||||
throw e;
|
||||
}
|
||||
log.info("源系统BulkV1查询BigObject数据任务完成, API: {}", api);
|
||||
return ReturnT.SUCCESS;
|
||||
}
|
||||
|
||||
/**
|
||||
* 目标系统BulkV1查询实现
|
||||
*
|
||||
* @param param 参数
|
||||
* @return 执行结果
|
||||
* @throws Exception 异常
|
||||
*/
|
||||
private ReturnT<String> manualBigObjectTargetV1Dump(SalesforceParam param) throws Exception {
|
||||
|
||||
String api = param.getApi();
|
||||
|
||||
String sql = param.getSql();
|
||||
|
||||
Integer type = param.getType();
|
||||
|
||||
BulkConnection bulkConnect = salesforceTargetConnect.createBulkConnect();
|
||||
|
||||
JobInfo job = null;
|
||||
|
||||
try {
|
||||
|
||||
log.info("开始目标系统BulkV1查询BigObject数据, API: {}, SQL: {}", api,sql);
|
||||
|
||||
// 创建Bulk作业
|
||||
job = BulkUtil.createJob(bulkConnect, api, OperationEnum.queryAll);
|
||||
log.info("创建Bulk V1查询作业成功, 作业ID: {}", job.getId());
|
||||
|
||||
// 提交查询批次
|
||||
BatchInfo batchInfo = bulkConnect.createBatchFromStream(job, new ByteArrayInputStream(sql.getBytes(StandardCharsets.UTF_8)));
|
||||
log.info("创建查询批次成功, 批次ID: {}", batchInfo.getId());
|
||||
|
||||
if (type == 0){
|
||||
return ReturnT.SUCCESS;
|
||||
}
|
||||
|
||||
// 等待作业完成
|
||||
int completionOne = BulkUtil.awaitCompletionOne(bulkConnect, job, batchInfo);
|
||||
log.info("查询批次处理完成, 处理记录数: {}", completionOne);
|
||||
|
||||
// 如果没有数据,直接结束
|
||||
if (completionOne == 0){
|
||||
log.info("无更多数据, 结束处理");
|
||||
BulkUtil.closeJob(bulkConnect, job.getId());
|
||||
return ReturnT.SUCCESS;
|
||||
}
|
||||
|
||||
// 获取查询结果列表
|
||||
QueryResultList queryResultList = bulkConnect.getQueryResultList(job.getId(), batchInfo.getId());
|
||||
log.info("获取查询结果列表, 结果数量: {}", queryResultList.getResult().length);
|
||||
|
||||
BulkUtil.closeJob(bulkConnect, job.getId());
|
||||
log.info("关闭作业成功, 作业ID: {}", job.getId());
|
||||
|
||||
}catch (Exception e) {
|
||||
log.error("目标系统BulkV1查询BigObject数据失败, API: {}", api, e);
|
||||
throw e;
|
||||
}
|
||||
log.info("目标系统BulkV1查询BigObject数据任务完成, API: {}", api);
|
||||
return ReturnT.SUCCESS;
|
||||
}
|
||||
|
||||
private ReturnT<String> manualBigObjectDumpV2(SalesforceParam param) throws Exception {
|
||||
|
||||
String api = param.getApi();
|
||||
|
||||
@ -1265,7 +1265,7 @@ public class DataImportNewServiceImpl implements DataImportNewService {
|
||||
//批量插入200一次
|
||||
int page = count % 200 == 0 ? count / 200 : (count / 200) + 1;
|
||||
for (int i = 0; i < page; i++) {
|
||||
List<Map<String, Object>> linkList = customMapper.list("Id,LinkedEntityId,ContentDocumentId,LinkedEntity_Type,ShareType,Visibility", api, "ShareType = 'V' and new_id = '0' order by Id limit " + i * 200 + ",200");
|
||||
List<Map<String, Object>> linkList = customMapper.list("Id,LinkedEntityId,ContentDocumentId,LinkedEntity_Type,ShareType,Visibility", api, "ShareType = 'V' and new_id = '0' order by Id limit 200");
|
||||
SObject[] accounts = new SObject[linkList.size()];
|
||||
String[] ids = new String[linkList.size()];
|
||||
int index = 0;
|
||||
@ -1289,13 +1289,27 @@ public class DataImportNewServiceImpl implements DataImportNewService {
|
||||
if (dMap != null){
|
||||
account.setField("ContentDocumentId", dMap.get("new_id").toString());
|
||||
}else {
|
||||
log.info("ContentDocumentLink Id: {},对应的ContentDocumentId: {} 数据不存在! " , id, contentDocumentId);
|
||||
String errorMessage = String.format("ContentDocumentLink Id: %s,对应的ContentDocumentId: %s 数据不存在!", id, contentDocumentId);
|
||||
log.info(errorMessage);
|
||||
List<Map<String, Object>> maps = new ArrayList<>();
|
||||
Map<String, Object> linkMap1 = new HashMap<>();
|
||||
linkMap1.put("key", "error_message");
|
||||
linkMap1.put("value", errorMessage);
|
||||
maps.add(linkMap1);
|
||||
customMapper.updateById(api, maps, id);
|
||||
continue;
|
||||
}
|
||||
if (lMap != null){
|
||||
account.setField("LinkedEntityId", lMap.get("new_id").toString());
|
||||
}else {
|
||||
log.info("ContentDocumentLink Id: {},对应的LinkedEntityId: {} 数据不存在! " , id, linkedEntityId);
|
||||
String errorMessage = String.format("ContentDocumentLink Id: %s,对应的LinkedEntityId: %s 数据不存在!", id, linkedEntityId);
|
||||
log.info(errorMessage);
|
||||
List<Map<String, Object>> maps = new ArrayList<>();
|
||||
Map<String, Object> linkMap1 = new HashMap<>();
|
||||
linkMap1.put("key", "error_message");
|
||||
linkMap1.put("value", errorMessage);
|
||||
maps.add(linkMap1);
|
||||
customMapper.updateById(api, maps, id);
|
||||
continue;
|
||||
}
|
||||
account.setField("ShareType", shareType);
|
||||
@ -1310,18 +1324,24 @@ public class DataImportNewServiceImpl implements DataImportNewService {
|
||||
printlnAccountsDetails(accounts,dataFields);
|
||||
}
|
||||
SaveResult[] saveResults = connection.create(accounts);
|
||||
;
|
||||
for (int j = 0; j < saveResults.length; j++) {
|
||||
if (!saveResults[j].getSuccess()) {
|
||||
String format = String.format("数据导入 error, api name: %s, \nparam: %s, \ncause:\n%s", api, com.alibaba.fastjson2.JSON.toJSONString(DataDumpParam.getFilter()), JSON.toJSONString(saveResults[j]));
|
||||
log.error(format);
|
||||
} else {
|
||||
List<Map<String, Object>> dList = new ArrayList<>();
|
||||
Map<String, Object> linkMap = new HashMap<>();
|
||||
linkMap.put("key", "new_id");
|
||||
linkMap.put("value", saveResults[j].getId());
|
||||
dList.add(linkMap);
|
||||
if (saveResults[j].getSuccess()) {
|
||||
List<Map<String, Object>> maps = new ArrayList<>();
|
||||
Map<String, Object> m = new HashMap<>();
|
||||
m.put("key", "new_id");
|
||||
m.put("value", saveResults[j].getId());
|
||||
maps.add(m);
|
||||
customMapper.updateById(api, maps, ids[j]);
|
||||
log.info("ContentDocumentLink Id: {},对应的new_id: {} 更新成功! " , ids[j], saveResults[j].getId());
|
||||
customMapper.updateById("ContentDocumentLink", dList, ids[j]);
|
||||
}else{
|
||||
List<Map<String, Object>> maps = new ArrayList<>();
|
||||
Map<String, Object> linkMap1 = new HashMap<>();
|
||||
linkMap1.put("key", "error_message");
|
||||
linkMap1.put("value", JSON.toJSONString(saveResults[j].getErrors()));
|
||||
maps.add(linkMap1);
|
||||
customMapper.updateById(api, maps, ids[j]);
|
||||
log.error("Id:{},saveResults: {}",ids[j], JSON.toJSONString(saveResults[j]));
|
||||
}
|
||||
}
|
||||
} catch (Exception e) {
|
||||
@ -1521,7 +1541,7 @@ public class DataImportNewServiceImpl implements DataImportNewService {
|
||||
break;
|
||||
}else {
|
||||
failCount++;
|
||||
log.error("文件下载失败,第"+ failCount +"次!, id: "+ id + ",返回体信息:" + response.message());
|
||||
log.info("文件下载失败,第"+ failCount +"次!, id: "+ id + ",返回体信息:" + response.message());
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@ -1805,11 +1805,231 @@ private void sendIncrementalQualityCheckEmail(DataVerifyParam param, String file
|
||||
return isError;
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理BigObject数据
|
||||
*/
|
||||
private boolean processBigObjectIncrementalData(String objectApi, List<DataField> fields,
|
||||
ExcelWriter excelWriter, ExcelWriter objectWriter, WriteSheet summarySheet, WriteSheet detailSheet,
|
||||
PartnerConnection sourceConn, PartnerConnection targetConn,
|
||||
Map<String, Map<String, Field>> fieldMap,Set< String> objectHashSet,String checkFields) {
|
||||
|
||||
|
||||
String[] checkFieldArray = checkFields != null ? checkFields.split(",") : new String[0];
|
||||
List<String> checkFieldList = Arrays.asList(checkFieldArray);
|
||||
// 分页处理记录
|
||||
String maxId = null;
|
||||
int totalRecordNum = 0;
|
||||
int targetRecordNum = 0;
|
||||
int sourceRecordNum = 0;
|
||||
int sorceOnlyNum = 0;
|
||||
int targetOnlyNum = 0;
|
||||
boolean isError = false;
|
||||
|
||||
Map<String, Field> sourceFieldMap = fieldMap.get("source");
|
||||
Map<String, Field> targetFieldMap = fieldMap.get("target");
|
||||
|
||||
List<DataField> fieldList = dataFieldService.list(new QueryWrapper<DataField>().eq("api", objectApi));
|
||||
String fieldStr = fieldList.stream()
|
||||
.map(DataField::getField)
|
||||
.filter(field -> !"SystemModstamp".equals(field) && !"LastModifiedDate".equals(field))
|
||||
.collect(Collectors.joining(", "));
|
||||
|
||||
try {
|
||||
DescribeSObjectResult dsr = sourceConn.describeSObject(objectApi);
|
||||
Field[] dsrFields = dsr.getFields();
|
||||
|
||||
String sql = "SELECT "+ fieldStr +" FROM " + objectApi + " LIMIT 100";
|
||||
|
||||
log.info("源ORG校验查询 sql: " + sql);
|
||||
|
||||
QueryResult queryResult = sourceConn.queryAll(sql);
|
||||
if (ObjectUtils.isEmpty(queryResult) || ObjectUtils.isEmpty(queryResult.getRecords())) {
|
||||
return false;
|
||||
}
|
||||
sourceRecordNum = queryResult.getSize();
|
||||
|
||||
SObject[] records = queryResult.getRecords();
|
||||
|
||||
JSONArray objects = DataUtil.toJsonArray(records, dsrFields);
|
||||
|
||||
for (int i = 0; i < objects.size(); i++) {
|
||||
|
||||
String targetSql = null;
|
||||
JSONObject sourceRecord = objects.getJSONObject(i);
|
||||
|
||||
StringBuilder whereClause = new StringBuilder("SELECT "+ fieldStr +" FROM " + objectApi + " WHERE ");
|
||||
boolean firstCondition = true;
|
||||
|
||||
List<String> ids = new ArrayList<>();
|
||||
|
||||
for (String checkField : checkFieldList) {
|
||||
if (!firstCondition) {
|
||||
whereClause.append(" AND ");
|
||||
}
|
||||
|
||||
Object fieldValue = sourceRecord.get(checkField);
|
||||
whereClause.append(checkField).append(" = '").append(fieldValue).append("'");
|
||||
|
||||
ids.add(String.valueOf(fieldValue));
|
||||
firstCondition = false;
|
||||
}
|
||||
|
||||
String id = String.join(",", ids);
|
||||
|
||||
targetSql = whereClause.toString();
|
||||
|
||||
log.info("目标ORG校验查询 sql: " + targetSql);
|
||||
|
||||
QueryResult targetQeruyResult = targetConn.queryAll(targetSql);
|
||||
|
||||
JSONArray tagetObject = DataUtil.toJsonArray(targetQeruyResult.getRecords(), dsrFields);
|
||||
|
||||
JSONObject targetRecord = tagetObject.getJSONObject(0);
|
||||
|
||||
if (sourceRecord != null && targetRecord == null) {
|
||||
sorceOnlyNum ++;
|
||||
List<String> row = Lists.newArrayList();
|
||||
row.add(" ");
|
||||
row.add(id);
|
||||
row.add(" ");
|
||||
row.add(objectApi);
|
||||
row.add(" ");
|
||||
row.add(" ");
|
||||
row.add(" ");
|
||||
row.add(" ");
|
||||
row.add("源ORG仅有");
|
||||
row.add(" ");
|
||||
row.add(" ");
|
||||
row.add(" ");
|
||||
row.add(" ");
|
||||
row.add(" ");
|
||||
row.add(" ");
|
||||
row.add(" ");
|
||||
// 写入到Excel的detailSheet,而不是添加到某个detailData列表
|
||||
objectWriter.write(Collections.singletonList(row), detailSheet);
|
||||
isError = true;
|
||||
continue;
|
||||
}
|
||||
|
||||
targetRecordNum ++ ;
|
||||
|
||||
// 比对字段值
|
||||
for (DataField field : fields) {
|
||||
String fieldName = field.getField();
|
||||
|
||||
if (Const.EXCLUDE_FIELDS.contains(fieldName))
|
||||
continue; // 跳过字段比较
|
||||
|
||||
Object sourceValue = sourceRecord.get(fieldName);
|
||||
Object targetValue = targetRecord.get(fieldName);
|
||||
|
||||
Field sorceField = sourceFieldMap.get(fieldName);
|
||||
Field targetField = targetFieldMap.get(fieldName);
|
||||
|
||||
if (sourceValue == null && targetValue == null) {
|
||||
continue;
|
||||
}
|
||||
|
||||
if (targetField == null){
|
||||
objectHashSet.add(fieldName);
|
||||
continue;
|
||||
}
|
||||
|
||||
//判断reference_to内是否包含User字符串
|
||||
String referenceTo = field.getReferenceTo();
|
||||
if (referenceTo !=null && (referenceTo.contains(",User") || referenceTo.contains("User,"))) {
|
||||
referenceTo = "User";
|
||||
}
|
||||
|
||||
if (referenceTo != null &&field.getSfType().equals("reference") && !referenceTo.equals("data_picklist")) {
|
||||
|
||||
String[] split = referenceTo.split(",");
|
||||
if (split.length > 1 && !"User".equals(referenceTo)){
|
||||
continue;
|
||||
}
|
||||
if (StringUtils.isBlank(referenceTo)){
|
||||
continue;
|
||||
}
|
||||
|
||||
Map<String, Object> objectMap = customMapper.getById("new_id", referenceTo, String.valueOf(sourceValue));
|
||||
if (objectMap == null || !String.valueOf(objectMap.get("new_id")).equals(String.valueOf(targetValue))) {
|
||||
List<String> row = Lists.newArrayList();
|
||||
row.add(" ");
|
||||
row.add(id);
|
||||
row.add(" ");
|
||||
row.add(objectApi);
|
||||
row.add(field.getField());
|
||||
row.add(field.getName());
|
||||
row.add(String.valueOf(sourceValue));
|
||||
row.add(String.valueOf(targetValue));
|
||||
row.add("字段值不一致");
|
||||
row.add(String.valueOf(sorceField.getType()));
|
||||
row.add(String.valueOf(targetField.getType()));
|
||||
row.add(String.valueOf(sorceField.isCreateable()));
|
||||
row.add(String.valueOf(targetField.isCreateable()));
|
||||
row.add(String.valueOf(sourceRecord.get("CreatedDate")));
|
||||
row.add(String.valueOf(sourceRecord.get("LastModifiedDate")));
|
||||
row.add(String.valueOf(targetRecord.get("LastModifiedDate")));
|
||||
// 写入到Excel的detailSheet,而不是添加到某个detailData列表
|
||||
objectWriter.write(Collections.singletonList(row), detailSheet);
|
||||
isError = true;
|
||||
}
|
||||
} else if (!Objects.equals(sourceValue, targetValue)){
|
||||
List<String> row = Lists.newArrayList();
|
||||
row.add(" ");
|
||||
row.add(id);
|
||||
row.add(" ");
|
||||
row.add(objectApi);
|
||||
row.add(field.getField());
|
||||
row.add(field.getName());
|
||||
row.add(String.valueOf(sourceValue));
|
||||
row.add(String.valueOf(targetValue));
|
||||
row.add("字段值不一致");
|
||||
row.add(String.valueOf(sorceField.getType()));
|
||||
row.add(String.valueOf(targetField.getType()));
|
||||
row.add(String.valueOf(sorceField.isCreateable()));
|
||||
row.add(String.valueOf(targetField.isCreateable()));
|
||||
row.add(String.valueOf(sourceRecord.get("CreatedDate")));
|
||||
row.add(String.valueOf(sourceRecord.get("LastModifiedDate")));
|
||||
row.add(String.valueOf(targetRecord.get("LastModifiedDate")));
|
||||
// 写入到Excel的detailSheet,而不是添加到某个detailData列表
|
||||
objectWriter.write(Collections.singletonList(row), detailSheet);
|
||||
isError = true;
|
||||
}
|
||||
}
|
||||
}
|
||||
TimeUnit.MILLISECONDS.sleep(1);
|
||||
} catch (InterruptedException e) {
|
||||
return isError;
|
||||
} catch (Exception throwable) {
|
||||
log.error("verify error", throwable);
|
||||
}
|
||||
|
||||
|
||||
// 写入批次汇总结果
|
||||
List<String> row = Lists.newArrayList();
|
||||
row.add(" ");
|
||||
row.add(objectApi);
|
||||
row.add("");
|
||||
row.add("");
|
||||
row.add(String.valueOf(sourceRecordNum));
|
||||
row.add(String.valueOf(targetRecordNum));
|
||||
row.add(String.valueOf(totalRecordNum));
|
||||
row.add(String.valueOf(sorceOnlyNum));
|
||||
row.add(String.valueOf(targetOnlyNum));
|
||||
row.add(isError?"不通过":"通过");
|
||||
|
||||
// 写入Excel的summarySheet
|
||||
excelWriter.write(Collections.singletonList(row), summarySheet);
|
||||
|
||||
return isError;
|
||||
}
|
||||
private Map<String, SObject> queryRecords(PartnerConnection connect, String objectApi ,List<String> ids){
|
||||
String idCondition = "Id IN (" + String.join(",", ids) + ")";
|
||||
String idSql = "SELECT FIELDS(ALL) FROM " + objectApi + " WHERE " + idCondition;
|
||||
// 将结果转换为Map便于比对
|
||||
Map<String, SObject> queryMap = Maps.newHashMap();
|
||||
log.info("查询ORG数据,SQL: {}", idSql);
|
||||
try {
|
||||
QueryResult queryResult = connect.queryAll(idSql);
|
||||
if (queryResult.getRecords() != null) {
|
||||
@ -1830,6 +2050,7 @@ private void sendIncrementalQualityCheckEmail(DataVerifyParam param, String file
|
||||
// 查询本地数据库获取ID列表
|
||||
String idSql = " CreatedDate >= '" + startDate + "'" +
|
||||
" AND CreatedDate < '" + endDate + "'" +
|
||||
" AND new_id is not null " +
|
||||
" ORDER BY Id ASC LIMIT " + (page * MAX_PAGE_RECORDS) + "," + MAX_PAGE_RECORDS;
|
||||
|
||||
return customMapper.list("Id,new_id", objectApi, idSql);
|
||||
@ -2163,6 +2384,178 @@ private void sendIncrementalQualityCheckEmail(DataVerifyParam param, String file
|
||||
return qw;
|
||||
}
|
||||
|
||||
/**
|
||||
* BigObject数据质量校验
|
||||
* 专门用于校验BigObject类型的数据质量
|
||||
*
|
||||
* @param param 参数
|
||||
* @return ReturnT
|
||||
*/
|
||||
@Override
|
||||
public ReturnT<String> verifyBigObjectQuality(DataVerifyParam param) throws Exception {
|
||||
log.info(">>> 开始执行BigObject数据质量校验任务 <<<");
|
||||
log.info("输入参数: {}", param != null ? param.toString() : "null");
|
||||
|
||||
try {
|
||||
// 获取需要校验的BigObject数据对象列表
|
||||
log.info("步骤1: 获取需要校验的BigObject数据对象列表");
|
||||
QueryWrapper<DataObject> wrapper = new QueryWrapper<>();
|
||||
if (StringUtils.isNotBlank(param.getApi())) {
|
||||
wrapper.in("name", DataUtil.toIdList(param.getApi()));
|
||||
}else {
|
||||
log.warn("没有需要校验的BigObject数据对象,终止执行");
|
||||
return ReturnT.FAIL;
|
||||
}
|
||||
List<DataObject> dataObjects = dataObjectService.list(wrapper);
|
||||
log.info("成功获取BigObject数据对象列表,共 {} 个对象", dataObjects.size());
|
||||
|
||||
// 如果没有需要校验的对象,直接返回失败
|
||||
if (dataObjects.isEmpty()) {
|
||||
log.warn("没有需要校验的BigObject数据对象,终止执行");
|
||||
return ReturnT.FAIL;
|
||||
}
|
||||
log.info("校验对象列表: {}", dataObjects.stream().map(DataObject::getName).collect(Collectors.joining(", ")));
|
||||
|
||||
// 创建Salesforce连接
|
||||
log.info("步骤2: 创建Salesforce连接");
|
||||
log.debug("正在创建源系统Salesforce连接...");
|
||||
PartnerConnection sourceConnect = salesforceConnect.createConnect();
|
||||
log.debug("源系统Salesforce连接创建成功");
|
||||
|
||||
log.debug("正在创建目标系统Salesforce连接...");
|
||||
PartnerConnection targetConnect = salesforceTargetConnect.createConnect();
|
||||
log.debug("目标系统Salesforce连接创建成功");
|
||||
log.info("Salesforce连接创建完成");
|
||||
|
||||
// 生成文件路径
|
||||
log.info("步骤3: 生成BigObject质量检查报告文件路径");
|
||||
String filePath = this.generateQualityCheckFilePath();
|
||||
log.info("BigObject质量检查报告文件路径: {}", filePath);
|
||||
|
||||
// 执行BigObject数据质量校验并写入Excel文件
|
||||
log.info("步骤4: 执行BigObject数据质量校验并写入Excel文件");
|
||||
long startTime = System.currentTimeMillis();
|
||||
log.info("开始时间: {}", new Date(startTime));
|
||||
|
||||
try {
|
||||
this.performBigObjectQualityCheckAndWriteToExcel(dataObjects, sourceConnect, targetConnect, filePath, param.getSendEmail(),param.getFields());
|
||||
long endTime = System.currentTimeMillis();
|
||||
log.info("BigObject数据质量校验执行完成,耗时: {} ms", endTime - startTime);
|
||||
} catch (Exception e) {
|
||||
long errorTime = System.currentTimeMillis();
|
||||
log.info("BigObject数据质量校验执行失败,耗时: {} ms", errorTime - startTime, e);
|
||||
return new ReturnT<>(ReturnT.FAIL.getCode(), "执行BigObject数据质量校验失败: " + e.getMessage());
|
||||
}
|
||||
|
||||
// 记录日志
|
||||
log.info("BigObject数据质量校验完成,报告文件路径: {}", filePath);
|
||||
|
||||
// 发送邮件(如果需要)
|
||||
log.info("步骤5: 检查并发送BigObject质量检查邮件");
|
||||
if (param.getSendEmail()) {
|
||||
log.debug("需要发送邮件,开始执行邮件发送逻辑");
|
||||
try {
|
||||
this.sendObjectQualityCheckEmail(param.getApi(), filePath);
|
||||
log.info("邮件发送完成");
|
||||
} catch (Exception e) {
|
||||
log.info("邮件发送失败", e);
|
||||
// 邮件发送失败不影响主流程,仅记录错误日志
|
||||
}
|
||||
} else {
|
||||
log.info("无需发送邮件,跳过邮件发送步骤");
|
||||
}
|
||||
|
||||
log.info(">>> BigObject数据质量校验任务执行完成 <<<");
|
||||
return ReturnT.SUCCESS;
|
||||
} catch (Exception e) {
|
||||
log.info("BigObject数据质量校验过程中发生未预期的错误", e);
|
||||
return new ReturnT<>(ReturnT.FAIL.getCode(), "BigObject数据质量校验失败: " + e.getMessage());
|
||||
} finally {
|
||||
log.info("BigObject数据质量校验任务结束");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 执行增量质量检查并将结果写入Excel文件
|
||||
*
|
||||
* @param dataObjects 数据对象列表
|
||||
* @param sourceConn 源系统连接
|
||||
* @param targetConn 目标系统连接
|
||||
* @param filePath 文件路径
|
||||
* @param sendEmail 是否发送邮件
|
||||
* @throws Exception 异常
|
||||
*/
|
||||
private void performBigObjectQualityCheckAndWriteToExcel(List<DataObject> dataObjects,
|
||||
PartnerConnection sourceConn,
|
||||
PartnerConnection targetConn,
|
||||
String filePath, boolean sendEmail,String fields) throws Exception {
|
||||
try (ExcelWriter excelWriter = EasyExcel.write(filePath).build()) {
|
||||
// 1. 创建汇总 Sheet
|
||||
WriteSheet summarySheet = this.createSummarySheet(excelWriter);
|
||||
// 2. 创建字段类型 Sheet
|
||||
WriteSheet fieldTypeSheet = this.createFieldTypeSheet(excelWriter);
|
||||
|
||||
List<Future<?>> futures = Lists.newArrayList();
|
||||
dataObjects.forEach(dataObject -> {
|
||||
Future<?> future = salesforceExecutor.execute(() -> {
|
||||
this.processIncrementalBigObjectDataObject(dataObject, excelWriter, summarySheet, fieldTypeSheet,
|
||||
sourceConn, targetConn, sendEmail, fields);
|
||||
}, 1, 1);
|
||||
futures.add(future);
|
||||
});
|
||||
|
||||
// 等待所有线程执行完毕
|
||||
try {
|
||||
salesforceExecutor.waitForFutures(futures.toArray(new Future<?>[0]));
|
||||
} catch (InterruptedException e) {
|
||||
log.info("线程执行被中断", e);
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理单个数据对象(增量)
|
||||
*/
|
||||
private void processIncrementalBigObjectDataObject(DataObject dataObject, ExcelWriter excelWriter, WriteSheet summarySheet,
|
||||
WriteSheet fieldTypeSheet, PartnerConnection sourceConn,
|
||||
PartnerConnection targetConn, boolean sendEmail,String checkFields) {
|
||||
// 1. 准备对象相关信息
|
||||
String objectApi = dataObject.getName();
|
||||
List<DataField> fields = this.getVerifiableFields(objectApi);
|
||||
if (fields.isEmpty()) {
|
||||
log.info("对象: {} 没有需要校验的字段!!!", objectApi);
|
||||
return;
|
||||
}
|
||||
|
||||
// 2. 处理字段相关信息
|
||||
Map<String, Map<String, Field>> fieldMap = this.handlerFieldInfo(fields, fieldTypeSheet, targetConn, sourceConn, dataObject, excelWriter);
|
||||
|
||||
String qualityCheckFilePath = this.generateObjectQualityCheckFilePath(objectApi);
|
||||
|
||||
HashSet<String> objectHashSet = new HashSet<>();
|
||||
|
||||
boolean hasError = false;
|
||||
try (ExcelWriter objectWriter = EasyExcel.write(qualityCheckFilePath).build()) {
|
||||
WriteSheet detailSheet = this.createDetailSheet(objectWriter, dataObject); // 修复:使用正确的ExcelWriter
|
||||
|
||||
hasError = this.processBigObjectIncrementalData(objectApi, fields, excelWriter, objectWriter, summarySheet, detailSheet,
|
||||
sourceConn, targetConn, fieldMap,objectHashSet,checkFields);
|
||||
} catch (Exception e) {
|
||||
log.info("处理对象 {} 数据时发生异常", objectApi, e);
|
||||
}
|
||||
|
||||
if (!objectHashSet.isEmpty()){
|
||||
String missingFields = "对象 " +objectApi + "在目标ORG中不存在以下字段:" + String.join(", ", objectHashSet);
|
||||
log.info(missingFields);
|
||||
EmailUtil.send("数据质量校验 ERROR", missingFields);
|
||||
}
|
||||
|
||||
if (sendEmail && hasError) {
|
||||
this.sendObjectQualityCheckEmail(objectApi, qualityCheckFilePath);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 统计sf db数据量 并更新batch 取出需要发送的数据列
|
||||
*
|
||||
|
||||
Loading…
Reference in New Issue
Block a user