【feat】 bulk调整优化

(cherry picked from commit c12fbf263a)
This commit is contained in:
Kris 2025-11-19 10:31:54 +08:00
parent 93d56c627d
commit 09c58e109e
4 changed files with 263 additions and 73 deletions

View File

@ -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<Future<?>> futures) throws InterruptedException {
QueryWrapper<DataObject> qw = new QueryWrapper<>();
qw.eq("data_work", 1)
.eq("data_lock", 0)
.orderByAsc("data_index")
.last(" limit 10");
while (true) {
List<DataObject> 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<SalesforceParam> salesforceParams = null;
update.setName(api);
update.setDataLock(1);
dataObjectService.updateById(update);
QueryWrapper<DataBatch> 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<DataBatch> 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;

View File

@ -286,6 +286,11 @@ public class DataImportBatchServiceImpl implements DataImportBatchService {
String api = param.getApi();
QueryWrapper<DataField> 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<DataField> 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<String> resultHeader = rdr.nextRecord();
int resultCols = resultHeader.size();
List<String> row;
while ((row = rdr.nextRecord()) != null) {
Map<String, String> resultInfo = new HashMap<String, String>();
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<Map<String, Object>> maps = new ArrayList<>();
Map<String, Object> m = new HashMap<>();
m.put("key", "new_id");
m.put("value", id);
maps.add(m);
customMapper.updateById(api, maps, ids[index]);
} else{
List<Map<String, Object>> maps = new ArrayList<>();
Map<String, Object> 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<String> resultHeader = rdr.nextRecord();
int resultCols = resultHeader.size();
List<String> row;
while ((row = rdr.nextRecord()) != null) {
Map<String, String> resultInfo = new HashMap<String, String>();
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<Map<String, Object>> maps = new ArrayList<>();
Map<String, Object> m = new HashMap<>();
m.put("key", "new_id");
m.put("value", id);
maps.add(m);
customMapper.updateById(api, maps, ids[index]);
} else{
List<Map<String, Object>> maps = new ArrayList<>();
Map<String, Object> 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<DataField> dbQw = new QueryWrapper<>();
dbQw.eq("api", api);
dbQw.and(wrapper -> wrapper.eq("is_updateable", 1)
.or().eq("field","Id"));
List<DataField> 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<DataField> dbQw = new QueryWrapper<>();
dbQw.eq("api", api);
dbQw.and(wrapper -> wrapper.eq("is_createable", 1)
.or().eq("field","Id"));
List<DataField> 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<DataField> 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<JSONObject> 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<JSONObject> 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<BatchInfo> 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());
}
}

View File

@ -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;
}

View File

@ -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<JSONObject> jsonList,
String fileName,
List<DataField> 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<Map<String, Object>> data, String fileName) throws IOException {
// 1. 创建目标目录不存在则创建
File targetDir = FileUtil.mkdir("data-dump/dataFile");