【feat】 bigobject数据质量校验

(cherry picked from commit 42c003cb2a)
This commit is contained in:
Kris 2025-10-31 10:31:23 +08:00
parent 40265c3c75
commit 902e618459
4 changed files with 427 additions and 0 deletions

View File

@ -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);
}
}

View File

@ -59,6 +59,12 @@ public class DataVerifyParam {
@ApiModelProperty(value = "忽略开始时间", hidden = true)
private Date ignoreBeginDate;
/**
* 校验字段 多个英文逗号分割
*/
@ApiModelProperty(value = "校验字段 多个英文逗号分割")
private String fields;
}

View File

@ -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;
}

View File

@ -1805,6 +1805,225 @@ 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;
@ -2163,6 +2382,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 取出需要发送的数据列
*