【feat】 调整Bulk api方法

This commit is contained in:
Kris 2025-08-06 15:13:30 +08:00
parent f7e2183989
commit 78e24b5fee

View File

@ -286,9 +286,6 @@ public class DataImportBatchServiceImpl implements DataImportBatchService {
QueryWrapper<DataField> dbQw = new QueryWrapper<>();
dbQw.eq("api", api);
List<DataField> list = dataFieldService.list(dbQw);
QueryWrapper<DataObject> obQw = new QueryWrapper<>();
obQw.eq("name", api);
DataObject dataObject = dataObjectService.getOne(obQw);
String beginDateStr = null;
String endDateStr = null;
@ -300,20 +297,24 @@ public class DataImportBatchServiceImpl implements DataImportBatchService {
}
//表内数据总量
Integer count = customMapper.countBySQL(api, "where new_id is null and CreatedDate >= '" + beginDateStr + "' and CreatedDate < '" + endDateStr + "'");
log.error("总Insert数据 count:{}-开始时间:{}-结束时间:{}-api:{}", count, beginDateStr, endDateStr, api);
if (count == 0) {
return;
}
//批量插入10000一次
int page = count%10000 == 0 ? count/10000 : (count/10000) + 1;
//批量插入5000一次
int page = count%5000 == 0 ? count/5000 : (count/5000) + 1;
log.error("总Insert数据 count:{};总批次:{}-开始时间:{}-结束时间:{}-api:{}", count,page, beginDateStr, endDateStr, api);
//总插入数
int sfNum = 0;
for (int i = 0; i < page; i++) {
List<JSONObject> data = customMapper.listJsonObject("*", api, "new_id is null and CreatedDate >= '" + beginDateStr + "' and CreatedDate < '" + endDateStr + "' limit 10000");
List<JSONObject> data = customMapper.listJsonObject("*", api, "new_id is null and CreatedDate >= '" + beginDateStr + "' and CreatedDate < '" + endDateStr + "' limit 5000");
int size = data.size();
log.info("执行api:{}, 执行page:{}, 执行size:{}", api, i+1, size);
log.error("总Insert数据 count:{};当前批次:{}-开始时间:{}-结束时间:{}-api:{}", count,i, beginDateStr, endDateStr, api);
List<JSONObject> insertList = new ArrayList<>();
//查询当前对象多态字段映射
@ -410,7 +411,7 @@ public class DataImportBatchServiceImpl implements DataImportBatchService {
ids[j-1] = data.get(j-1).get("Id").toString();
insertList.add(account);
if (i*10000+j == count){
if (i*5000+j == count){
break;
}
}
@ -418,16 +419,10 @@ public class DataImportBatchServiceImpl implements DataImportBatchService {
JobInfo salesforceQueryJob = null;
String fullPath = null;
try {
//写入csv文件
fullPath = CsvConverterUtil.writeToCsv(insertList, UUID.randomUUID().toString(),list,dataObject.getIsEditable());
salesforceInsertJob = BulkUtil.createJob(bulkConnection, api, OperationEnum.insert);
List<BatchInfo> batchInfos = BulkUtil.createBatchesFromCSVFile(bulkConnection, salesforceInsertJob, fullPath);
List<BatchInfo> batchInfos = CsvConverterUtil.writeToCsvNew(insertList, UUID.randomUUID().toString(),list,true,bulkConnection,salesforceInsertJob);
BulkUtil.awaitCompletion(bulkConnection, salesforceInsertJob, batchInfos);
@ -462,8 +457,6 @@ public class DataImportBatchServiceImpl implements DataImportBatchService {
log.error("createdNewIdBatch error api:{}", api, e);
throw e;
}finally {
new File(fullPath).delete();
BulkUtil.closeJob(bulkConnection, salesforceInsertJob.getId());
BulkUtil.closeJob(bulkConnection, salesforceQueryJob.getId());
@ -767,6 +760,7 @@ public class DataImportBatchServiceImpl implements DataImportBatchService {
}
// 表内数据总量
Integer count = customMapper.countBySQL(api, "where new_id is not null and CreatedDate >= '" + beginDateStr + "' and CreatedDate < '" + endDateStr + "'");
if(count == 0){
return;
}
@ -787,19 +781,20 @@ public class DataImportBatchServiceImpl implements DataImportBatchService {
// 总更新数
int sfNum = 0;
// 批量更新10000一次
int page = count%10000 == 0 ? count/10000 : (count/10000) + 1;
// 批量更新5000一次
int page = count%5000 == 0 ? count/5000 : (count/5000) + 1;
log.error("总Update数据 count:{};总批次:{}-开始时间:{}-结束时间:{}-api:{}", count,page, beginDateStr, endDateStr, api);
for (int i = 0; i < page; i++) {
List<Map<String, Object>> mapList = customMapper.list("*", api, "new_id is not null and CreatedDate >= '" + beginDateStr + "' and CreatedDate < '" + endDateStr + "' order by Id asc limit " + i * 10000 + ",10000");
log.error("总Update数据 count:{};当前批次:{}-开始时间:{}-结束时间:{}-api:{}", count,i, beginDateStr, endDateStr, api);
List<JSONObject> updateList = new ArrayList<>();
for (Map<String, Object> map : mapList) {
JSONObject account = new JSONObject();
//给对象赋值
for (DataField dataField : list) {
String field = dataField.getField();
@ -855,25 +850,24 @@ public class DataImportBatchServiceImpl implements DataImportBatchService {
updateList.add(account);
}
JobInfo salesforceUpdateJob = null;
try {
//写入csv文件
String fullPath = CsvConverterUtil.writeToCsv(updateList, UUID.randomUUID().toString(),list,dataObject.getIsEditable());
JobInfo salesforceInsertJob = BulkUtil.createJob(bulkConnection, api, OperationEnum.update);
salesforceUpdateJob = BulkUtil.createJob(bulkConnection, api, OperationEnum.update);
List<BatchInfo> batchInfos = BulkUtil.createBatchesFromCSVFile(bulkConnection, salesforceInsertJob, fullPath);
List<BatchInfo> batchInfos = CsvConverterUtil.writeToCsvNew(updateList, UUID.randomUUID().toString(),list,true,bulkConnection,salesforceUpdateJob);
BulkUtil.awaitCompletion(bulkConnection, salesforceInsertJob, batchInfos);
BulkUtil.awaitCompletion(bulkConnection, salesforceUpdateJob, batchInfos);
sfNum = sfNum + checkUpdateResults(bulkConnection, salesforceInsertJob, batchInfos,api);
BulkUtil.closeJob(bulkConnection, salesforceInsertJob.getId());
new File(fullPath).delete();
sfNum = sfNum + checkUpdateResults(bulkConnection, salesforceUpdateJob, batchInfos,api);
} catch (Throwable e) {
log.info(e.getMessage());
throw e;
}finally {
BulkUtil.closeJob(bulkConnection, salesforceUpdateJob.getId());
}
}
@ -986,6 +980,11 @@ public class DataImportBatchServiceImpl implements DataImportBatchService {
update.setDataLock(1);
dataObjectService.updateById(update);
if (api.endsWith("Share")){
insertSingleShareData(api,bulkConnection);
continue;
}
QueryWrapper<DataBatch> dbQw = new QueryWrapper<>();
dbQw.eq("name", api);
List<DataBatch> list = dataBatchService.list(dbQw);
@ -1130,19 +1129,20 @@ public class DataImportBatchServiceImpl implements DataImportBatchService {
return;
}
//批量插入10000一次
int page = count%10000 == 0 ? count/10000 : (count/10000) + 1;
int page = count%5000 == 0 ? count/5000 : (count/5000) + 1;
log.error("总Insert数据 count:{}批次:{}-开始时间:{}-结束时间:{}-api:{}", count,page, beginDateStr, endDateStr, api);
log.error("总Insert数据 count:{}批次:{}-开始时间:{}-结束时间:{}-api:{}", count,page, beginDateStr, endDateStr, api);
//总插入数
int sfNum = 0;
for (int i = 0; i < page; i++) {
List<JSONObject> data = customMapper.listJsonObject("*", api, "new_id is null and CreatedDate >= '" + beginDateStr + "' and CreatedDate < '" + endDateStr + "' limit 10000");
List<JSONObject> data = customMapper.listJsonObject("*", api, "new_id is null and CreatedDate >= '" + beginDateStr + "' and CreatedDate < '" + endDateStr + "' limit 5000");
int size = data.size();
log.info("执行api:{}, 执行page:{}, 执行size:{}", api, i+1, size);
log.error("总Insert数据 count:{};当前批次:{}-开始时间:{}-结束时间:{}-api:{}", count,i, beginDateStr, endDateStr, api);
List<JSONObject> insertList = new ArrayList<>();
//查询当前对象多态字段映射
@ -1221,7 +1221,6 @@ public class DataImportBatchServiceImpl implements DataImportBatchService {
}
}
}
}
account.put("old_owner_id__c", data.get(j - 1).get("OwnerId"));
@ -1229,7 +1228,7 @@ public class DataImportBatchServiceImpl implements DataImportBatchService {
ids[j-1] = data.get(j-1).get("Id").toString();
insertList.add(account);
if (i*1000+j == count){
if (i*5000+j == count){
break;
}
}
@ -1238,16 +1237,11 @@ public class DataImportBatchServiceImpl implements DataImportBatchService {
JobInfo salesforceQueryJob = null;
String fullPath = null;
try {
//写入csv文件
fullPath = CsvConverterUtil.writeToCsv(insertList, UUID.randomUUID().toString(),list,true);
salesforceInsertJob = BulkUtil.createJob(bulkConnection, api, OperationEnum.insert);
List<BatchInfo> batchInfos = BulkUtil.createBatchesFromCSVFile(bulkConnection, salesforceInsertJob, fullPath);
List<BatchInfo> batchInfos = CsvConverterUtil.writeToCsvNew(insertList, UUID.randomUUID().toString(),list,true,bulkConnection,salesforceInsertJob);
BulkUtil.awaitCompletion(bulkConnection, salesforceInsertJob, batchInfos);
@ -1282,8 +1276,6 @@ public class DataImportBatchServiceImpl implements DataImportBatchService {
log.error("insertSingleData error api:{}", api, e);
throw e;
}finally {
new File(fullPath).delete();
BulkUtil.closeJob(bulkConnection, salesforceInsertJob.getId());
BulkUtil.closeJob(bulkConnection, salesforceQueryJob.getId());
@ -1306,5 +1298,131 @@ public class DataImportBatchServiceImpl implements DataImportBatchService {
}
/**
* 执行一次性Insert Share数据
*/
private void insertSingleShareData(String api, BulkConnection bulkConnection) throws Exception {
QueryWrapper<DataField> dbQw = new QueryWrapper<>();
dbQw.eq("api", api);
List<DataField> list = dataFieldService.list(dbQw);
TimeUnit.MILLISECONDS.sleep(1);
String beginDateStr = null;
String endDateStr = null;
//表内数据总量
Integer count = customMapper.countBySQL(api, "where new_id is null and RowCause = 'Manual'");
if (count == 0) {
return;
}
//批量插入10000一次
int page = count%5000 == 0 ? count/5000 : (count/5000) + 1;
log.error("总Insert数据 count:{};总批次:{}-开始时间:{}-结束时间:{}-api:{}", count,page, beginDateStr, endDateStr, api);
for (int i = 0; i < page; i++) {
List<JSONObject> data = customMapper.listJsonObject("*", api, "new_id is null and RowCause = 'Manual' order by Id asc limit " + i * 10000 + ",10000");
int size = data.size();
log.error("总Insert数据 count:{};当前批次:{}-开始时间:{}-结束时间:{}-api:{}", count,i, beginDateStr, endDateStr, api);
List<JSONObject> insertList = new ArrayList<>();
//查询当前对象多态字段映射
QueryWrapper<LinkConfig> queryWrapper = new QueryWrapper<>();
queryWrapper.eq("api",api).eq("is_create",1).eq("is_link",true);
List<LinkConfig> configs = linkConfigService.list(queryWrapper);
Map<String, String> fieldMap = new HashMap<>();
if (!configs.isEmpty()) {
fieldMap = configs.stream()
.collect(Collectors.toMap(
LinkConfig::getField, // Key提取器
LinkConfig::getLinkField, // Value提取器
(oldVal, newVal) -> newVal // 解决重复Key冲突保留新值
));
}
//更新对象的new_id
String[] ids = new String[size];
for (int j = 1; j <= size; j++) {
JSONObject account = new JSONObject();
for (DataField dataField : list) {
if ("Owner_Type".equals(dataField.getField()) || "Id".equals(dataField.getField())){
continue;
}
if (dataField.getIsCreateable() !=null && dataField.getIsCreateable()) {
if ("reference".equals(dataField.getSfType()) && data.get(j - 1).get(dataField.getField()) != null){
//引用类型
String reference_to = dataField.getReferenceTo();
//引用类型字段
String linkfield = fieldMap.get(dataField.getField());
if (StringUtils.isNotBlank(linkfield)){
reference_to = data.get(j-1).get(linkfield)!=null?data.get(j-1).get(linkfield).toString():null;
}
if (reference_to == null){
continue;
}
if (reference_to.contains(",User") || reference_to.contains("User,")) {
reference_to = "User";
}
Map<String, Object> m = customMapper.getById("new_id", reference_to, data.get(j - 1).get(dataField.getField()).toString());
if (m != null && !m.isEmpty()) {
account.put(dataField.getField(), m.get("new_id"));
}else {
String message = "对象类型:" + api + "的数据:"+ data.get(j - 1).get("Id") +"的引用对象:" + reference_to + "的数据:"+ data.get(j - 1).get(dataField.getField()) +"不存在!";
EmailUtil.send("DataDump ERROR", message);
log.info(message);
return;
}
}else {
if (data.get(j - 1).get(dataField.getField()) != null && StringUtils.isNotBlank(dataField.getSfType())) {
account.put(dataField.getField(), DataUtil.localBulkDataToSfData(dataField.getSfType(), data.get(j - 1).get(dataField.getField()).toString()));
}else {
account.put(dataField.getField(), data.get(j - 1).get(dataField.getField()) );
}
}
}
}
ids[j-1] = data.get(j-1).get("Id").toString();
insertList.add(account);
if (i*5000+j == count){
break;
}
}
JobInfo salesforceInsertJob = null;
try {
salesforceInsertJob = BulkUtil.createJob(bulkConnection, api, OperationEnum.insert);
List<BatchInfo> batchInfos = CsvConverterUtil.writeToCsvNew(insertList, UUID.randomUUID().toString(),list,true,bulkConnection,salesforceInsertJob);
BulkUtil.awaitCompletion(bulkConnection, salesforceInsertJob, batchInfos);
checkInsertResults(bulkConnection, salesforceInsertJob, batchInfos, api, ids);
} catch (Exception e) {
log.error("insertSingleShareData error api:{}", api, e);
throw e;
}finally {
BulkUtil.closeJob(bulkConnection, salesforceInsertJob.getId());
}
}
}
}