From 902e618459378528804ca4b8f866e95d6129fdf1 Mon Sep 17 00:00:00 2001 From: Kris <2893855659@qq.com> Date: Fri, 31 Oct 2025 10:31:23 +0800 Subject: [PATCH] =?UTF-8?q?=E3=80=90feat=E3=80=91=20bigobject=E6=95=B0?= =?UTF-8?q?=E6=8D=AE=E8=B4=A8=E9=87=8F=E6=A0=A1=E9=AA=8C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit (cherry picked from commit 42c003cb2a9de43aa5f1d96167bec2d9f5e66af9) --- .../celnet/datadump/job/DataVerifyJob.java | 23 ++ .../datadump/param/DataVerifyParam.java | 6 + .../datadump/service/DataVerifyService.java | 7 + .../service/impl/DataVerifyServiceImpl.java | 391 ++++++++++++++++++ 4 files changed, 427 insertions(+) diff --git a/src/main/java/com/celnet/datadump/job/DataVerifyJob.java b/src/main/java/com/celnet/datadump/job/DataVerifyJob.java index a5cf530..068301b 100644 --- a/src/main/java/com/celnet/datadump/job/DataVerifyJob.java +++ b/src/main/java/com/celnet/datadump/job/DataVerifyJob.java @@ -123,4 +123,27 @@ public class DataVerifyJob { return dataVerifyService.verifyIncrementalQuality(param); } + /** + * BigObject数据质量校验任务 + * 专门用于校验BigObject类型的数据质量 + * + * @param paramStr 参数json + * @return result + */ + @XxlJob("dataVerifyBigObjectQualityJob") + public ReturnT 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); + } + } \ No newline at end of file diff --git a/src/main/java/com/celnet/datadump/param/DataVerifyParam.java b/src/main/java/com/celnet/datadump/param/DataVerifyParam.java index 42d8594..0ed7e32 100644 --- a/src/main/java/com/celnet/datadump/param/DataVerifyParam.java +++ b/src/main/java/com/celnet/datadump/param/DataVerifyParam.java @@ -59,6 +59,12 @@ public class DataVerifyParam { @ApiModelProperty(value = "忽略开始时间", hidden = true) private Date ignoreBeginDate; + /** + * 校验字段 多个英文逗号分割 + */ + @ApiModelProperty(value = "校验字段 多个英文逗号分割") + private String fields; + } diff --git a/src/main/java/com/celnet/datadump/service/DataVerifyService.java b/src/main/java/com/celnet/datadump/service/DataVerifyService.java index 25a5e69..cc8f48d 100644 --- a/src/main/java/com/celnet/datadump/service/DataVerifyService.java +++ b/src/main/java/com/celnet/datadump/service/DataVerifyService.java @@ -27,5 +27,12 @@ public interface DataVerifyService { * @return ReturnT */ ReturnT verifyIncrementalQuality(DataVerifyParam param) throws Exception; + + /** + * BigObject数据质量校验 + * @param param 参数 + * @return ReturnT + */ + ReturnT verifyBigObjectQuality(DataVerifyParam param) throws Exception; } \ No newline at end of file diff --git a/src/main/java/com/celnet/datadump/service/impl/DataVerifyServiceImpl.java b/src/main/java/com/celnet/datadump/service/impl/DataVerifyServiceImpl.java index 64ab2a7..12651c7 100644 --- a/src/main/java/com/celnet/datadump/service/impl/DataVerifyServiceImpl.java +++ b/src/main/java/com/celnet/datadump/service/impl/DataVerifyServiceImpl.java @@ -1805,6 +1805,225 @@ private void sendIncrementalQualityCheckEmail(DataVerifyParam param, String file return isError; } + /** + * 处理BigObject数据 + */ + private boolean processBigObjectIncrementalData(String objectApi, List fields, + ExcelWriter excelWriter, ExcelWriter objectWriter, WriteSheet summarySheet, WriteSheet detailSheet, + PartnerConnection sourceConn, PartnerConnection targetConn, + Map> fieldMap,Set< String> objectHashSet,String checkFields) { + + + String[] checkFieldArray = checkFields != null ? checkFields.split(",") : new String[0]; + List 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 sourceFieldMap = fieldMap.get("source"); + Map targetFieldMap = fieldMap.get("target"); + + List fieldList = dataFieldService.list(new QueryWrapper().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 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 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 objectMap = customMapper.getById("new_id", referenceTo, String.valueOf(sourceValue)); + if (objectMap == null || !String.valueOf(objectMap.get("new_id")).equals(String.valueOf(targetValue))) { + List 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 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 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 queryRecords(PartnerConnection connect, String objectApi ,List 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 verifyBigObjectQuality(DataVerifyParam param) throws Exception { + log.info(">>> 开始执行BigObject数据质量校验任务 <<<"); + log.info("输入参数: {}", param != null ? param.toString() : "null"); + + try { + // 获取需要校验的BigObject数据对象列表 + log.info("步骤1: 获取需要校验的BigObject数据对象列表"); + QueryWrapper wrapper = new QueryWrapper<>(); + if (StringUtils.isNotBlank(param.getApi())) { + wrapper.in("name", DataUtil.toIdList(param.getApi())); + }else { + log.warn("没有需要校验的BigObject数据对象,终止执行"); + return ReturnT.FAIL; + } + List 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 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> 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 fields = this.getVerifiableFields(objectApi); + if (fields.isEmpty()) { + log.info("对象: {} 没有需要校验的字段!!!", objectApi); + return; + } + + // 2. 处理字段相关信息 + Map> fieldMap = this.handlerFieldInfo(fields, fieldTypeSheet, targetConn, sourceConn, dataObject, excelWriter); + + String qualityCheckFilePath = this.generateObjectQualityCheckFilePath(objectApi); + + HashSet 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 取出需要发送的数据列 *