From ff7fff9f90e3bc9e27754608821df3706728870c Mon Sep 17 00:00:00 2001 From: Kris <2893855659@qq.com> Date: Fri, 31 Oct 2025 10:48:18 +0800 Subject: [PATCH] =?UTF-8?q?=E3=80=90feat=E3=80=91=20=E7=9B=AE=E6=A0=87org?= =?UTF-8?q?=E5=A2=9E=E5=8A=A0bulk=20v1=E6=9F=A5=E8=AF=A2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit (cherry picked from commit d0c94cd1911988a26e16ebd557dae52c628130ea) --- .../celnet/datadump/job/DataDumpNewJob.java | 24 ++++++ .../datadump/service/CommonBatchService.java | 2 + .../service/impl/CommonBatchServiceImpl.java | 76 +++++++++++++++++++ 3 files changed, 102 insertions(+) diff --git a/src/main/java/com/celnet/datadump/job/DataDumpNewJob.java b/src/main/java/com/celnet/datadump/job/DataDumpNewJob.java index 0f2054f..a47f1d5 100644 --- a/src/main/java/com/celnet/datadump/job/DataDumpNewJob.java +++ b/src/main/java/com/celnet/datadump/job/DataDumpNewJob.java @@ -162,6 +162,30 @@ public class DataDumpNewJob { return commonBatchService.dumpBigObjectSourceV1(param); } + /** + * 目标系统BulkV1查询 + * + * @param paramStr 参数json + * @return result + */ + @XxlJob("targetSystemBulkV1QueryJob") + public ReturnT 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) * diff --git a/src/main/java/com/celnet/datadump/service/CommonBatchService.java b/src/main/java/com/celnet/datadump/service/CommonBatchService.java index 34775ee..41997f5 100644 --- a/src/main/java/com/celnet/datadump/service/CommonBatchService.java +++ b/src/main/java/com/celnet/datadump/service/CommonBatchService.java @@ -14,5 +14,7 @@ public interface CommonBatchService { ReturnT importBigObject(SalesforceParam param) throws Exception; ReturnT dumpBigObjectSourceV1(SalesforceParam param) throws Exception; + + ReturnT dumpBigObjectTargetV1(SalesforceParam param) throws Exception; } \ No newline at end of file diff --git a/src/main/java/com/celnet/datadump/service/impl/CommonBatchServiceImpl.java b/src/main/java/com/celnet/datadump/service/impl/CommonBatchServiceImpl.java index 0936fa7..008e463 100644 --- a/src/main/java/com/celnet/datadump/service/impl/CommonBatchServiceImpl.java +++ b/src/main/java/com/celnet/datadump/service/impl/CommonBatchServiceImpl.java @@ -197,6 +197,21 @@ public class CommonBatchServiceImpl implements CommonBatchService { return ReturnT.SUCCESS; } + @Override + public ReturnT dumpBigObjectTargetV1(SalesforceParam param) throws Exception { + try { + // 手动任务 + ReturnT result = manualBigObjectTargetV1Dump(param); + if (result != null) { + return result; + } + } catch (Throwable throwable) { + log.error("dump error", throwable); + throw throwable; + } + return ReturnT.SUCCESS; + } + private ReturnT manualBigObjectDump(SalesforceParam param) throws Exception { String api = param.getApi(); @@ -402,6 +417,67 @@ public class CommonBatchServiceImpl implements CommonBatchService { return ReturnT.SUCCESS; } + /** + * 目标系统BulkV1查询实现 + * + * @param param 参数 + * @return 执行结果 + * @throws Exception 异常 + */ + private ReturnT 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 manualBigObjectDumpV2(SalesforceParam param) throws Exception { String api = param.getApi();