From 09c58e109e0c49615c1744080a9cef0668c4ac29 Mon Sep 17 00:00:00 2001 From: Kris <2893855659@qq.com> Date: Wed, 19 Nov 2025 10:31:54 +0800 Subject: [PATCH] =?UTF-8?q?=E3=80=90feat=E3=80=91=20bulk=E8=B0=83=E6=95=B4?= =?UTF-8?q?=E4=BC=98=E5=8C=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit (cherry picked from commit c12fbf263a274ab56160647290ff40d53e472e3e) --- .../service/impl/CommonBatchServiceImpl.java | 83 +++++++++ .../impl/DataImportBatchServiceImpl.java | 169 +++++++++++------- .../com/celnet/datadump/util/BulkUtil.java | 1 + .../datadump/util/CsvConverterUtil.java | 83 ++++++++- 4 files changed, 263 insertions(+), 73 deletions(-) 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 01d8916..5d30ebc 100644 --- a/src/main/java/com/celnet/datadump/service/impl/CommonBatchServiceImpl.java +++ b/src/main/java/com/celnet/datadump/service/impl/CommonBatchServiceImpl.java @@ -132,6 +132,8 @@ public class CommonBatchServiceImpl implements CommonBatchService { if (result != null) { return result; } + }else { + autoBatchDump(param, futures); } return ReturnT.SUCCESS; } catch (Throwable throwable) { @@ -949,6 +951,87 @@ public class CommonBatchServiceImpl implements CommonBatchService { return ReturnT.SUCCESS; } + private void autoBatchDump(SalesforceParam param, List> futures) throws InterruptedException { + QueryWrapper qw = new QueryWrapper<>(); + qw.eq("data_work", 1) + .eq("data_lock", 0) + .orderByAsc("data_index") + .last(" limit 10"); + while (true) { + List dataObjects = dataObjectService.list(qw); + if (CollectionUtils.isEmpty(dataObjects)) { + break; + } + + // 根据参数获取sql + for (DataObject update : dataObjects) { + String api = update.getName(); + TimeUnit.MILLISECONDS.sleep(1); + try { + commonService.checkApi(api, true); + if (!dataFieldService.hasCreatedDate(api)) { + DataDumpSpecialParam dataDumpSpecialParam = new DataDumpSpecialParam(); + dataDumpSpecialParam.setApi(api); + Future future = getData(dataDumpSpecialParam, salesforceConnect.createBulkConnect()); + // 等待当前所有线程执行完成 + salesforceExecutor.waitForFutures(future); + continue; + } + param.setApi(api); + List salesforceParams = null; + update.setName(api); + update.setDataLock(1); + dataObjectService.updateById(update); + QueryWrapper dbQw = new QueryWrapper<>(); + dbQw.eq("name", api) + .isNull("first_sync_date"); + if (param.getBeginDate() != null && param.getEndDate() != null) { + dbQw.eq("sync_start_date", param.getBeginDate()); + dbQw.eq("sync_end_date", param.getEndDate()); + } + List list = dataBatchService.list(dbQw); + AtomicInteger batch = new AtomicInteger(1); + if (CollectionUtils.isNotEmpty(list)) { + salesforceParams = list.stream().map(t -> { + SalesforceParam salesforceParam = param.clone(); + salesforceParam.setApi(t.getName()); + salesforceParam.setBeginCreateDate(t.getSyncStartDate()); + salesforceParam.setEndCreateDate(t.getSyncEndDate()); + salesforceParam.setBatch(batch.getAndIncrement()); + return salesforceParam; + }).collect(Collectors.toList()); + } else { + salesforceParams = DataUtil.splitTask(param); + } + // 手动任务优先执行 + for (SalesforceParam salesforceParam : salesforceParams) { + Future future = salesforceExecutor.execute(() -> { + try { + dumpBatchData(salesforceParam); + } catch (Throwable throwable) { + log.error("salesforceExecutor error", throwable); + throw new RuntimeException(throwable); + } + }, salesforceParam.getBatch(), 1); + futures.add(future); + } + + // 等待当前所有线程执行完成 + salesforceExecutor.waitForFutures(futures.toArray(new Future[]{})); + + update.setDataWork(0); + } catch (Throwable e) { + log.error("manualDump error", e); + throw new RuntimeException(e); + } finally { + update.setName(api); + update.setDataLock(0); + dataObjectService.updateById(update); + } + } + } + } + private BulkConnection dumpBatchData(SalesforceParam param) throws Throwable { String api = param.getApi(); BulkConnection bulkConnect = null; diff --git a/src/main/java/com/celnet/datadump/service/impl/DataImportBatchServiceImpl.java b/src/main/java/com/celnet/datadump/service/impl/DataImportBatchServiceImpl.java index d5a95e9..3440ca9 100644 --- a/src/main/java/com/celnet/datadump/service/impl/DataImportBatchServiceImpl.java +++ b/src/main/java/com/celnet/datadump/service/impl/DataImportBatchServiceImpl.java @@ -286,6 +286,11 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { String api = param.getApi(); QueryWrapper dbQw = new QueryWrapper<>(); dbQw.eq("api", api); + dbQw.and(wrapper -> wrapper.eq("is_createable", 1) + .eq("is_nillable", 0) + .eq("is_defaulted_on_create", 0) + .or().eq("field", "CreatedDate") + .or().eq("field", "CreatedById")); List list = dataFieldService.list(dbQw); String beginDateStr = null; @@ -346,10 +351,12 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { String message = null; JSONObject account = new JSONObject(); for (DataField dataField : list) { - if ("OwnerId".equals(dataField.getField()) || "Owner_Type".equals(dataField.getField()) - || "Id".equals(dataField.getField())){ + if ("OwnerId".equals(dataField.getField()) || "Owner_Type".equals(dataField.getField())){ continue; } + if ("Id".equals(dataField.getField())){ + account.put("Id", data.get(j - 1).get(dataField.getField())); + } if ("CreatedDate".equals(dataField.getField()) && dataField.getIsCreateable()){ // 转换为UTC时间并格式化 LocalDateTime localDateTime = LocalDateTime.parse(String.valueOf(data.get(j - 1).get("CreatedDate")), inputFormatter); @@ -369,6 +376,7 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { } continue; } + if (dataField.getIsCreateable() !=null && dataField.getIsCreateable() && !dataField.getIsNillable() && !dataField.getIsDefaultedOnCreate()) { if ("reference".equals(dataField.getSfType()) && data.get(j - 1).get(dataField.getField()) != null){ //引用类型 @@ -419,7 +427,6 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { account.put("old_sfdc_id__c", data.get(j-1).get("Id")); } - ids[j-1] = data.get(j-1).get("Id").toString(); insertList.add(account); if (i*8000+j >= count){ break; @@ -435,7 +442,7 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { salesforceInsertJob = BulkUtil.createJob(bulkConnection, api, OperationEnum.insert); - BatchInfo batchInfo = CsvConverterUtil.writeToCsvNewOne(insertList, UUID.randomUUID().toString(), list, true, bulkConnection, salesforceInsertJob); + BatchInfo batchInfo = CsvConverterUtil.writeToCsvOne(insertList, UUID.randomUUID().toString(), list, true, bulkConnection, salesforceInsertJob,ids); BulkUtil.awaitOneCompletion(bulkConnection, salesforceInsertJob, batchInfo); @@ -536,37 +543,6 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { while (retryCount <= maxRetries) { try { rdr = new CSVReader(connection.getBatchResultStream(job.getId(), b.getId())); - List resultHeader = rdr.nextRecord(); - int resultCols = resultHeader.size(); - - List row; - while ((row = rdr.nextRecord()) != null) { - Map resultInfo = new HashMap(); - for (int i = 0; i < resultCols; i++) { - resultInfo.put(resultHeader.get(i), row.get(i)); - } - boolean insertStatus = Boolean.valueOf(resultInfo.get("Success")); - boolean created = Boolean.valueOf(resultInfo.get("Created")); - String id = resultInfo.get("Id"); - String error = resultInfo.get("Error"); - if (insertStatus && created) { - List> maps = new ArrayList<>(); - Map m = new HashMap<>(); - m.put("key", "new_id"); - m.put("value", id); - maps.add(m); - customMapper.updateById(api, maps, ids[index]); - } else{ - List> maps = new ArrayList<>(); - Map linkMap1 = new HashMap<>(); - linkMap1.put("key", "error_message"); - linkMap1.put("value", JSON.toJSONString(error)); - maps.add(linkMap1); - customMapper.updateById(api, maps, ids[index]); - log.info("Id:{},saveResults: {}",ids[index], error); - } - index ++; - } break; // 成功执行,跳出重试循环 }catch (Exception e) { retryCount++; @@ -579,6 +555,38 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { } } } + + List resultHeader = rdr.nextRecord(); + int resultCols = resultHeader.size(); + + List row; + while ((row = rdr.nextRecord()) != null) { + Map resultInfo = new HashMap(); + for (int i = 0; i < resultCols; i++) { + resultInfo.put(resultHeader.get(i), row.get(i)); + } + boolean insertStatus = Boolean.valueOf(resultInfo.get("Success")); + boolean created = Boolean.valueOf(resultInfo.get("Created")); + String id = resultInfo.get("Id"); + String error = resultInfo.get("Error"); + if (insertStatus && created) { + List> maps = new ArrayList<>(); + Map m = new HashMap<>(); + m.put("key", "new_id"); + m.put("value", id); + maps.add(m); + customMapper.updateById(api, maps, ids[index]); + } else{ + List> maps = new ArrayList<>(); + Map linkMap1 = new HashMap<>(); + linkMap1.put("key", "error_message"); + linkMap1.put("value", JSON.toJSONString(error)); + maps.add(linkMap1); + customMapper.updateById(api, maps, ids[index]); + log.info("Id:{},saveResults: {}",ids[index], error); + } + index ++; + } return index; } @@ -818,6 +826,8 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { String api = param.getApi(); QueryWrapper dbQw = new QueryWrapper<>(); dbQw.eq("api", api); + dbQw.and(wrapper -> wrapper.eq("is_updateable", 1) + .or().eq("field","Id")); List list = dataFieldService.list(dbQw); String beginDateStr = null; @@ -934,13 +944,17 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { if (map.get(field) != null && StringUtils.isNotBlank(dataField.getSfType())) { account.put(field, DataUtil.localBulkDataToSfData(dataField.getSfType(), String.valueOf(map.get(field)))); }else { - account.put(field, value); + if (map.get(field) == null){ + account.put(field, ""); + }else { + account.put(field, value); + } } } } if (dataObject.getIsEditable()){ - account.put("old_owner_id__c", map.get("OwnerId")); + account.put("old_owner_id__c", map.get("OwnerId") == null?"":map.get("OwnerId")); account.put("old_sfdc_id__c", map.get("Id")); } updateList.add(account); @@ -955,11 +969,12 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { salesforceUpdateJob = BulkUtil.createJob(bulkConnection, api, OperationEnum.update); - BatchInfo batchInfo = CsvConverterUtil.writeToCsvNewOne(updateList, UUID.randomUUID().toString(), list, true, bulkConnection, salesforceUpdateJob); + BatchInfo batchInfo = CsvConverterUtil.writeToCsvNewOne(updateList, UUID.randomUUID().toString(), list, dataObject.getIsEditable(), bulkConnection, salesforceUpdateJob); BulkUtil.awaitOneCompletion(bulkConnection, salesforceUpdateJob, batchInfo); dataLog1.setEndTime(new Date()); + dataLogService.save(dataLog1); DataLog dataLog2 = new DataLog(api, "开始时间:"+beginDateStr+";结束时间:"+endDateStr, new Date(), null, "数据更新(BULK),更新DB第" + i + "批数据", "UpdateDB"); @@ -1145,8 +1160,8 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { update.setDataLock(1); dataObjectService.updateById(update); - if (dataFieldService.hasCreatedDate(api)){ - insertSingleShareData(api,bulkConnection); + if (!dataFieldService.hasCreatedDate(api)){ + insertSingleShareData(api,bulkConnection,update); continue; } @@ -1223,8 +1238,8 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { update.setDataLock(1); dataObjectService.updateById(update); - if (api.endsWith("Share") || "GroupMember".equals(api)){ - insertSingleShareData(api,bulkConnection); + if (dataFieldService.hasCreatedDate(api)){ + insertSingleShareData(api,bulkConnection,dataObject); continue; } @@ -1281,6 +1296,8 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { String api = param.getApi(); QueryWrapper dbQw = new QueryWrapper<>(); dbQw.eq("api", api); + dbQw.and(wrapper -> wrapper.eq("is_createable", 1) + .or().eq("field","Id")); List list = dataFieldService.list(dbQw); TimeUnit.MILLISECONDS.sleep(1); @@ -1293,16 +1310,8 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { endDateStr = DateUtil.format(endDate, "yyyy-MM-dd HH:mm:ss"); } - String sql = ""; - String sql2 = ""; - - if (api.endsWith("Share")){ - sql = "where RowCause NOT IN ('Owner','Rule','Territory','Team','ImplicitChild','ImplicitParent','TerritoryRule','ImplicitCallCenter','PortalRole','Portal','ImplicitPerson','ImplicitGrant') and new_id is null "; - sql2 = "RowCause NOT IN ('Owner','Rule','Territory','Team','ImplicitChild','ImplicitParent','TerritoryRule','ImplicitCallCenter','PortalRole','Portal','ImplicitPerson','ImplicitGrant') and new_id is null limit 8000"; - }else { - sql = "where new_id is null "; - sql2 = "new_id is null limit 8000"; - } + String sql = "where new_id is null and CreatedDate >= '" + beginDateStr + "' and CreatedDate < '" + endDateStr + "'"; + String sql2 = "new_id is null and CreatedDate >= '" + beginDateStr + "' and CreatedDate < '" + endDateStr + "' order by Id asc limit 8000"; //表内数据总量 Integer count = customMapper.countBySQL(api, sql); @@ -1354,9 +1363,12 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { String message = null; for (DataField dataField : list) { - if ("Owner_Type".equals(dataField.getField()) || "Id".equals(dataField.getField())){ + if ("Owner_Type".equals(dataField.getField())){ continue; } + if ("Id".equals(dataField.getField())){ + account.put("Id", data.get(j - 1).get(dataField.getField())); + } if ("CreatedDate".equals(dataField.getField()) && dataField.getIsCreateable()){ // 转换为UTC时间并格式化 LocalDateTime localDateTime = LocalDateTime.parse(String.valueOf(data.get(j - 1).get("CreatedDate")), inputFormatter); @@ -1405,7 +1417,11 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { 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()) ); + if (data.get(j - 1).get(dataField.getField()) == null){ + account.put(dataField.getField(), ""); + }else { + account.put(dataField.getField(), data.get(j - 1).get(dataField.getField()) ); + } } } } @@ -1420,10 +1436,9 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { continue; } if (dataObject.getIsEditable()) { - account.put("old_owner_id__c", data.get(j - 1).get("OwnerId")); + account.put("old_owner_id__c", data.get(j - 1).get("OwnerId") == null? "":data.get(j - 1).get("OwnerId")); account.put("old_sfdc_id__c", data.get(j - 1).get("Id")); } - ids[j-1] = data.get(j-1).get("Id").toString(); insertList.add(account); if (i*8000+j >= count){ break; @@ -1440,7 +1455,7 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { salesforceInsertJob = BulkUtil.createJob(bulkConnection, api, OperationEnum.insert); - BatchInfo batchInfo = CsvConverterUtil.writeToCsvNewOne(insertList, UUID.randomUUID().toString(), list, dataObject.getIsEditable(), bulkConnection, salesforceInsertJob); + BatchInfo batchInfo = CsvConverterUtil.writeToCsvOne(insertList, UUID.randomUUID().toString(), list, dataObject.getIsEditable(), bulkConnection, salesforceInsertJob,ids); BulkUtil.awaitOneCompletion(bulkConnection, salesforceInsertJob, batchInfo); @@ -1453,6 +1468,7 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { sfNum = sfNum + checkInsertResultsOne(bulkConnection, salesforceInsertJob, batchInfo, api, ids); dataLog2.setEndTime(new Date()); + dataLogService.save(dataLog2); } catch (Exception e) { @@ -1482,7 +1498,7 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { /** * 执行一次性Insert Share数据 */ - private void insertSingleShareData(String api, BulkConnection bulkConnection) throws Exception { + private void insertSingleShareData(String api, BulkConnection bulkConnection, DataObject dataObject) throws Exception { QueryWrapper dbQw = new QueryWrapper<>(); dbQw.eq("api", api); @@ -1492,8 +1508,19 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { String beginDateStr = null; String endDateStr = null; + String sql = ""; + String sql2 = ""; + + if (api.endsWith("Share")){ + sql = "where RowCause NOT IN ('Owner','Rule','Territory','Team','ImplicitChild','ImplicitParent','TerritoryRule','ImplicitCallCenter','PortalRole','Portal','ImplicitPerson','ImplicitGrant') and new_id is null "; + sql2 = "RowCause NOT IN ('Owner','Rule','Territory','Team','ImplicitChild','ImplicitParent','TerritoryRule','ImplicitCallCenter','PortalRole','Portal','ImplicitPerson','ImplicitGrant') and new_id is null limit 8000"; + }else { + sql = "where new_id is null "; + sql2 = "new_id is null limit 8000"; + } + //表内数据总量 - Integer count = customMapper.countBySQL(api, "where new_id is null and RowCause NOT IN ('Owner','Rule','Territory','Team','ImplicitChild','ImplicitParent','TerritoryRule','ImplicitCallCenter','PortalRole','Portal','ImplicitPerson','ImplicitGrant')"); + Integer count = customMapper.countBySQL(api, sql); if (count == 0) { return; @@ -1522,7 +1549,7 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { DataLog dataLog = new DataLog(api, "开始时间:"+beginDateStr+";结束时间:"+endDateStr, new Date(), null, "一次性插入(BULK),查询并组装" + i + "批数据", "QueryBD"); - List data = customMapper.listJsonObject("*", api, "new_id is null and RowCause NOT IN ('Owner','Rule','Territory','Team','ImplicitChild','ImplicitParent','TerritoryRule','ImplicitCallCenter','PortalRole','Portal','ImplicitPerson','ImplicitGrant') order by Id asc limit " + i * 10000 + ",10000"); + List data = customMapper.listJsonObject("*", api, sql2); int size = data.size(); log.error("总Insert数据 count:{};当前批次:{},-开始时间:{};-结束时间:{};-api:{};", count,i, beginDateStr, endDateStr, api); @@ -1537,9 +1564,12 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { String message = null; for (DataField dataField : list) { - if ("Owner_Type".equals(dataField.getField()) || "Id".equals(dataField.getField())){ + if ("Owner_Type".equals(dataField.getField())){ continue; } + if ("Id".equals(dataField.getField())){ + account.put("Id", data.get(j - 1).get(dataField.getField())); + } if (dataField.getIsCreateable() !=null && dataField.getIsCreateable()) { if ("reference".equals(dataField.getSfType()) && data.get(j - 1).get(dataField.getField()) != null){ //引用类型 @@ -1574,7 +1604,11 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { 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()) ); + if (data.get(j - 1).get(dataField.getField()) == null){ + account.put(dataField.getField(), ""); + }else { + account.put(dataField.getField(), data.get(j - 1).get(dataField.getField()) ); + } } } } @@ -1590,7 +1624,6 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { continue; } - ids[j-1] = data.get(j-1).get("Id").toString(); insertList.add(account); if (i*8000+j == count){ break; @@ -1608,9 +1641,9 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { salesforceInsertJob = BulkUtil.createJob(bulkConnection, api, OperationEnum.insert); - List batchInfos = CsvConverterUtil.writeToCsvNew(insertList, UUID.randomUUID().toString(),list,true,bulkConnection,salesforceInsertJob); + BatchInfo batchInfo = CsvConverterUtil.writeToCsvOne(insertList, UUID.randomUUID().toString(), list,dataObject.getIsEditable(), bulkConnection, salesforceInsertJob,ids); - BulkUtil.awaitCompletion(bulkConnection, salesforceInsertJob, batchInfos); + BulkUtil.awaitOneCompletion(bulkConnection, salesforceInsertJob, batchInfo); dataLog1.setEndTime(new Date()); @@ -1618,16 +1651,16 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { DataLog dataLog2 = new DataLog(api, "开始时间:"+beginDateStr+";结束时间:"+endDateStr, new Date(), null, "一次性插入(BULK),更新DB第" + i + "批数据", "UpdateDB"); - checkInsertResults(bulkConnection, salesforceInsertJob, batchInfos, api, ids); + checkInsertResultsOne(bulkConnection, salesforceInsertJob, batchInfo, api, ids); dataLog2.setEndTime(new Date()); dataLogService.save(dataLog2); + } catch (Exception e) { - log.error("insertSingleShareData error api:{}", api, e); + log.error("createdNewIdBatch error api:{}", api, e); throw e; }finally { - BulkUtil.closeJob(bulkConnection, salesforceInsertJob.getId()); } } diff --git a/src/main/java/com/celnet/datadump/util/BulkUtil.java b/src/main/java/com/celnet/datadump/util/BulkUtil.java index 6919c77..bc2b545 100644 --- a/src/main/java/com/celnet/datadump/util/BulkUtil.java +++ b/src/main/java/com/celnet/datadump/util/BulkUtil.java @@ -165,6 +165,7 @@ public class BulkUtil { job.setObject(sobjectType); job.setOperation(operation); job.setContentType(ContentType.CSV); + job.setConcurrencyMode(ConcurrencyMode.Serial); job = connection.createJob(job); return job; } diff --git a/src/main/java/com/celnet/datadump/util/CsvConverterUtil.java b/src/main/java/com/celnet/datadump/util/CsvConverterUtil.java index 2d06d2a..85d222f 100644 --- a/src/main/java/com/celnet/datadump/util/CsvConverterUtil.java +++ b/src/main/java/com/celnet/datadump/util/CsvConverterUtil.java @@ -192,13 +192,11 @@ public class CsvConverterUtil { if (isEditable){ header = Stream.concat( - fields.stream().filter(dataField ->(dataField.getIsCreateable() != null && dataField.getIsCreateable())) - .map(DataField::getField), - Stream.of("Id","old_owner_id__c", "old_sfdc_id__c") + fields.stream().map(DataField::getField), + Stream.of("old_owner_id__c", "old_sfdc_id__c") ).toArray(String[]::new); }else { - header = fields.stream().filter(dataField ->(dataField.getIsCreateable() != null && dataField.getIsCreateable())) - .map(DataField::getField).toArray(String[]::new); + header = fields.stream().map(DataField::getField).toArray(String[]::new); } BatchInfo batch = null; @@ -240,6 +238,81 @@ public class CsvConverterUtil { return batch; } + public static BatchInfo writeToCsvOne(List jsonList, + String fileName, + List fields, + Boolean isEditable, + BulkConnection connection, + JobInfo jobInfo,String[] ids) { + + // 1. 创建目标目录(不存在则创建) + File targetDir = FileUtil.mkdir("data-dump/dataFile"); + + // 2. 构建完整文件路径 + String fullPath = targetDir.getAbsolutePath() + File.separator + fileName + ".csv"; + + CsvWriter csvWriter = CsvUtil.getWriter(fullPath, CharsetUtil.CHARSET_UTF_8, false); + + String[] header = null; + + File file = null; + + if (isEditable){ + header = Stream.concat( + fields.stream().filter(dataField ->(dataField.getIsCreateable() != null && dataField.getIsCreateable())) + .map(DataField::getField), + Stream.of("old_owner_id__c", "old_sfdc_id__c") + ).toArray(String[]::new); + }else { + header = fields.stream().filter(dataField ->(dataField.getIsCreateable() != null && dataField.getIsCreateable())) + .map(DataField::getField).toArray(String[]::new); + } + + BatchInfo batch = null; + + try { + + // 2. 写入表头(必须使用 String[]) + csvWriter.writeHeaderLine(header); + + int index = 0; + + // 遍历数据列表 + for (JSONObject jsonObject : jsonList) { + // 按表头顺序获取值 + String[] row = new String[header.length]; + for (int i = 0; i < header.length; i++) { + // 将每个值转换为字符串(如果为null则转为空字符串) + Object value = jsonObject.get(header[i]); + row[i] = String.valueOf(value); + } + ids[index] = jsonObject.get("Id").toString(); + + csvWriter.writeLine(row); + + index ++; + } + + // 关闭writer(在try-with-resources中可省略,但这里我们显式关闭) + csvWriter.close(); + + file = new File(fullPath); + + InputStream inputStream = Files.newInputStream(file.toPath()); + + batch = connection.createBatchFromStream(jobInfo, inputStream); + + } catch (IOException e) { + log.error("CSV文件操作失败: {}", fullPath, e); + file.deleteOnExit(); + } catch (AsyncApiException e) { + log.error("Bulk API上传失败: {}", e.getMessage(), e); + file.deleteOnExit(); + } + + return batch; + } + public static String exportToCsv(List> data, String fileName) throws IOException { // 1. 创建目标目录(不存在则创建) File targetDir = FileUtil.mkdir("data-dump/dataFile");