parent
902e618459
commit
ff7fff9f90
@ -162,6 +162,30 @@ public class DataDumpNewJob {
|
||||
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)
|
||||
*
|
||||
|
||||
@ -14,5 +14,7 @@ public interface CommonBatchService {
|
||||
ReturnT<String> importBigObject(SalesforceParam param) throws Exception;
|
||||
|
||||
ReturnT<String> dumpBigObjectSourceV1(SalesforceParam param) throws Exception;
|
||||
|
||||
ReturnT<String> dumpBigObjectTargetV1(SalesforceParam param) throws Exception;
|
||||
|
||||
}
|
||||
@ -197,6 +197,21 @@ public class CommonBatchServiceImpl implements CommonBatchService {
|
||||
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();
|
||||
@ -402,6 +417,67 @@ public class CommonBatchServiceImpl implements CommonBatchService {
|
||||
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();
|
||||
|
||||
Loading…
Reference in New Issue
Block a user