diff --git a/src/main/java/com/celnet/datadump/controller/JobController.java b/src/main/java/com/celnet/datadump/controller/JobController.java index c3da061..0037874 100644 --- a/src/main/java/com/celnet/datadump/controller/JobController.java +++ b/src/main/java/com/celnet/datadump/controller/JobController.java @@ -144,7 +144,7 @@ public class JobController { @PostMapping("/createApi") @ApiOperation("创建表结构") public Result createApi(String paramStr) throws Exception { - log.info("dataImportBatchJob execute start .................."); + log.info("createApi execute start .................."); SalesforceParam param = new SalesforceParam(); if (StringUtils.isNotBlank(paramStr)) { param = JSON.parseObject(paramStr, SalesforceParam.class); @@ -161,6 +161,31 @@ public class JobController { } } + /** + * 创建表结构 + * + * @return result + */ + @PostMapping("/createBigObjectApi") + @ApiOperation("创建BigOBject表结构") + public Result createBigObjectApi(String paramStr) throws Exception { + log.info("createBigObjectApi execute start .................."); + SalesforceParam param = new SalesforceParam(); + if (StringUtils.isNotBlank(paramStr)) { + param = JSON.parseObject(paramStr, SalesforceParam.class); + } + param.setType(1); + // 参数转换 + param.setBeginCreateDate(param.getBeginDate()); + param.setEndCreateDate(param.getEndDate()); + ReturnT returnT = commonService.createBigObjectApi(param); + if (returnT.getCode() == ReturnT.SUCCESS_CODE) { + return Result.success(); + } else { + return Result.fail(returnT.getMsg()); + } + } + /** * 存量任务 * @@ -251,7 +276,7 @@ public class JobController { } /** - * bulk批量大数据生成newSFID + * 存量任务(大数据量) * @param paramStr * @author kris * @return @@ -277,6 +302,23 @@ public class JobController { return commonBatchService.dumpBatch(param); } + @PostMapping("/dataDumpBigObjectJob") + @ApiOperation("存量任务(BigObject)") + @LogServiceAnnotation(operateType = OperateTypeConstant.TYPE_INSERT, remark = "存量任务(BigObject)") + public ReturnT dataDumpBigObjectJob(String paramStr) throws Exception { + log.info("dataDumpBigObjectJob execute start .................."); + SalesforceParam param = new SalesforceParam(); + try { + if (StringUtils.isNotBlank(paramStr)) { + param = JSON.parseObject(paramStr, SalesforceParam.class); + } + } catch (Throwable throwable) { + return new ReturnT<>(500, "参数解析失败!"); + } + + return commonBatchService.dumpBigObject(param); + } + /** * bulk批量更新大数据量数据 * @param paramStr 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 fb5e03a..065cc99 100644 --- a/src/main/java/com/celnet/datadump/service/impl/CommonBatchServiceImpl.java +++ b/src/main/java/com/celnet/datadump/service/impl/CommonBatchServiceImpl.java @@ -17,10 +17,7 @@ import com.celnet.datadump.param.DataDumpParam; import com.celnet.datadump.param.DataDumpSpecialParam; import com.celnet.datadump.param.SalesforceParam; import com.celnet.datadump.service.*; -import com.celnet.datadump.util.BulkUtil; -import com.celnet.datadump.util.DataUtil; -import com.celnet.datadump.util.EmailUtil; -import com.celnet.datadump.util.SqlUtil; +import com.celnet.datadump.util.*; import com.google.common.collect.Lists; import com.google.common.collect.Maps; @@ -45,7 +42,12 @@ import org.apache.commons.lang3.time.DateUtils; import org.apache.http.Header; import org.apache.http.client.methods.CloseableHttpResponse; import org.apache.http.client.methods.HttpGet; +import org.apache.http.client.methods.HttpPatch; import org.apache.http.client.methods.HttpPost; +import org.apache.http.client.methods.HttpPut; +import org.apache.http.client.methods.HttpPut; +import org.apache.http.entity.FileEntity; +import org.apache.http.entity.StringEntity; import org.apache.http.impl.client.CloseableHttpClient; import org.apache.http.impl.client.HttpClients; import org.apache.http.util.EntityUtils; @@ -56,6 +58,8 @@ import javax.annotation.PostConstruct; import java.io.*; import java.nio.charset.StandardCharsets; import java.nio.file.Files; +import java.nio.file.Path; +import java.nio.file.Paths; import java.util.*; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; @@ -93,7 +97,10 @@ public class CommonBatchServiceImpl implements CommonBatchService { // 批次处理记录数,默认为10000 private int batchProcessCount = 10000; - + + private static final String BIGOBJECT_FILE_PATH = "bigObject/"; + + @Autowired private SystemConfigService systemConfigService; @@ -176,231 +183,241 @@ public class CommonBatchServiceImpl implements CommonBatchService { String api = param.getApi(); BulkConnection bulkConnect = salesforceConnect.createBulkConnect(); - - String beginDateStr = null; - String endDateStr = null; - if (param.getBeginCreateDate() != null && param.getEndCreateDate() != null){ - Date beginDate = param.getBeginCreateDate(); - Date endDate = param.getEndCreateDate(); - beginDateStr = DateUtil.format(beginDate, "yyyy-MM-dd HH:mm:ss"); - endDateStr = DateUtil.format(endDate, "yyyy-MM-dd HH:mm:ss"); - } + JobInfo job = null; try { - List salesforceParams = null; - QueryWrapper dbQw = new QueryWrapper<>(); - dbQw.eq("name", api); - if (StringUtils.isNotEmpty(beginDateStr) && StringUtils.isNotEmpty(endDateStr)) { - dbQw.eq("sync_start_date", beginDateStr); // 等于开始时间 - dbQw.eq("sync_end_date", endDateStr); // 等于结束时间 - } - 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()); + + log.info("开始导出BigObject数据, API: {}", api); + + // 检测路径是否存在 不存在则创建 + File excel = new File( param.getApi() + "/"); + if (!excel.exists()) { + System.out.println("创建目录: " + param.getApi() + "/"); + boolean mkdir = excel.mkdir(); } - // 手动任务优先执行 - for (SalesforceParam salesforceParam : salesforceParams) { - if (param.getIsSingleThread()){ - manualBigObject(salesforceParam, bulkConnect); - }else { - Future future = salesforceExecutor.execute(() -> { - try { - manualBigObject(salesforceParam, bulkConnect); - } catch (Throwable throwable) { - log.error("salesforceExecutor error", throwable); - throw new RuntimeException(throwable); + List fieldList = dataFieldService.list(new QueryWrapper().eq("api", param.getApi())); + String fieldStr = fieldList.stream() + .map(DataField::getField) + .collect(Collectors.joining(", ")); + + String sql = "select " + fieldStr + " from " + param.getApi() ; + log.info("构建查询SQL: {}", sql); + + job = BulkUtil.createJob(bulkConnect, api, OperationEnum.queryAll); + log.info("创建Bulk作业成功, 作业ID: {}", job.getId()); + + BatchInfo batchInfo = bulkConnect.createBatchFromStream(job, new ByteArrayInputStream(sql.getBytes(StandardCharsets.UTF_8))); + log.info("创建批次成功, 批次ID: {}", batchInfo.getId()); + + 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); + + + + // 使用确定性命名策略,避免使用随机时间戳 + String fileName = param.getApi() + "/" + DateFormatUtils.format(new Date(), "yyyyMMddHHmmss"); + final long maxFileSize = 50 * 1024 * 1024; // 50MB + + int fileIndex = 0; + String currentFilePath = fileName + "_" + fileIndex + ".csv"; + + log.info("开始写入数据到文件 : {}", fileName); + + BufferedWriter bufferedWriter = null; + FileWriter fileWriter = null; + + try { + fileWriter = new FileWriter(currentFilePath, true); + bufferedWriter = new BufferedWriter(fileWriter); + long currentFileSize = 0; + + int resultIndex = 0; + for (final String resultId : queryResultList.getResult()) { + log.info("处理结果ID: {}, 进度: {}/{}", resultId, ++resultIndex, queryResultList.getResult().length); + + InputStream resultStream = bulkConnect.getQueryResultStream(job.getId(), batchInfo.getId(), resultId); + + try (BufferedReader bufferedReader = new BufferedReader(new InputStreamReader(resultStream, StandardCharsets.UTF_8))) { + String line; + int lineCount = 0; + while ((line = bufferedReader.readLine()) != null) { + lineCount++; + byte[] lineBytes = (line + "\n").getBytes(StandardCharsets.UTF_8); + + // 如果添加这行会超出文件大小限制,则创建新文件 + if (currentFileSize + lineBytes.length > maxFileSize && currentFileSize > 0) { + // 关闭当前文件 + if (bufferedWriter != null) { + bufferedWriter.close(); + } + if (fileWriter != null) { + fileWriter.close(); + } + + log.info("当前文件已达到大小限制, 创建新文件. 文件路径: {}, 数据行数: {}", currentFilePath, lineCount); + + // 创建新文件 + fileIndex++; + currentFilePath = fileName + "_" + fileIndex + ".csv"; + fileWriter = new FileWriter(currentFilePath, true); + bufferedWriter = new BufferedWriter(fileWriter); + currentFileSize = 0; + lineCount = 0; + } + + // 写入数据 + bufferedWriter.write(line + "\n"); + currentFileSize += lineBytes.length; } - }, salesforceParam.getBatch(), 1); - futures.add(future); + log.info("结果ID {} 处理完成", resultId); + } + } + log.info("所有数据写入完成, 总文件数: {}", fileIndex + 1); + } finally { + if (bufferedWriter != null) { + bufferedWriter.close(); + } + if (fileWriter != null) { + fileWriter.close(); } } - // 等待当前所有线程执行完成 - salesforceExecutor.waitForFutures(futures.toArray(new Future[]{})); - }catch (InterruptedException interruptedException){ - return ReturnT.FAIL; + BulkUtil.closeJob(bulkConnect, job.getId()); + log.info("关闭作业成功, 作业ID: {}", job.getId()); + }catch (Exception e) { + log.error("导出BigObject数据失败, API: {}", api, e); throw e; } + log.info("BigObject数据导出任务完成, API: {}", api); return ReturnT.SUCCESS; } - private void manualBigObject(SalesforceParam param, BulkConnection bulkConnect) throws Exception { - String sql = "select * from " + param.getApi() + " where CreatedDate >= '" + param.getBeginCreateDate() + "' and CreatedDate <= '" + param.getEndCreateDate() + "'"; - - log.info("BigObject query SQL: {}", sql); - XxlJobLogger.log("BigObject query SQL: {}", sql); - - // 获取Salesforce实例URL - String restEndpoint = bulkConnect.getConfig().getRestEndpoint(); - String url = restEndpoint.substring(0, restEndpoint.indexOf("services/")); - - // 创建Bulk API 2.0查询作业 - String jobId = BulkUtil.createBulkV2Job(url, sql, bulkConnect); - log.info("Created Bulk API 2.0 job with ID: {}", jobId); - - // 等待作业完成 - BulkUtil.waitForBulkV2JobCompletion(url, jobId, bulkConnect); - - // 获取查询结果 - String csvResult = BulkUtil.getBulkV2QueryJobResults(url, jobId, bulkConnect); - - // 保存到本地文件 - String fileName = param.getApi() + "_" + System.currentTimeMillis() + ".csv"; - String filePath = "bigobject/" + fileName; - - // 确保目录存在 - File directory = new File("bigobject"); - if (!directory.exists()) { - directory.mkdirs(); - } - - // 写入文件 - try (FileWriter fileWriter = new FileWriter(filePath); - BufferedWriter bufferedWriter = new BufferedWriter(fileWriter)) { - bufferedWriter.write(csvResult); - } - - log.info("BigObject data saved to file: {}", filePath); - XxlJobLogger.log("BigObject data saved to file: {}", filePath); - } - - private ReturnT manualBigObjectImport(SalesforceParam param, List> futures) throws Exception { +private ReturnT manualBigObjectImport(SalesforceParam param, List> futures) throws Exception { String api = param.getApi(); // 创建目标ORG连接 BulkConnection targetBulkConnect = salesforceTargetConnect.createBulkConnect(); - - String beginDateStr = null; - String endDateStr = null; - if (param.getBeginCreateDate() != null && param.getEndCreateDate() != null){ - Date beginDate = param.getBeginCreateDate(); - Date endDate = param.getEndCreateDate(); - beginDateStr = DateUtil.format(beginDate, "yyyy-MM-dd HH:mm:ss"); - endDateStr = DateUtil.format(endDate, "yyyy-MM-dd HH:mm:ss"); - } - - try { - List salesforceParams = null; - QueryWrapper dbQw = new QueryWrapper<>(); - dbQw.eq("name", api); - if (StringUtils.isNotEmpty(beginDateStr) && StringUtils.isNotEmpty(endDateStr)) { - dbQw.eq("sync_start_date", beginDateStr); // 等于开始时间 - dbQw.eq("sync_end_date", endDateStr); // 等于结束时间 - } - 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()); - } - - // 手动任务优先执行 - for (SalesforceParam salesforceParam : salesforceParams) { - if (param.getIsSingleThread()){ - manualImportBigObject(salesforceParam, targetBulkConnect); - }else { - Future future = salesforceExecutor.execute(() -> { - try { - manualImportBigObject(salesforceParam, targetBulkConnect); - } catch (Throwable throwable) { - log.error("salesforceExecutor error", throwable); - throw new RuntimeException(throwable); - } - }, salesforceParam.getBatch(), 1); - futures.add(future); - } - } - // 等待当前所有线程执行完成 - salesforceExecutor.waitForFutures(futures.toArray(new Future[]{})); - - }catch (InterruptedException interruptedException){ - return ReturnT.FAIL; - }catch (Exception e) { - throw e; - } - return ReturnT.SUCCESS; - } - - private void manualImportBigObject(SalesforceParam param, BulkConnection targetBulkConnect) throws Exception { - String api = param.getApi(); - - // 构造CSV文件路径 - 使用与导出时相同的文件名格式(但需要知道确切的时间戳) - // 这里应该查找最新的相关文件,而不是随机生成一个文件名 - File directory = new File("bigobject"); + + File directory = new File(api); if (!directory.exists() || !directory.isDirectory()) { - log.error("BigObject目录不存在: {}", directory.getAbsolutePath()); - throw new FileNotFoundException("BigObject目录不存在: " + directory.getAbsolutePath()); + throw new FileNotFoundException(api + "目录不存在: " + directory.getAbsolutePath()); } - + // 查找匹配的CSV文件(最近的) - File[] files = directory.listFiles((dir, name) -> - name.startsWith(api + "_") && name.endsWith(".csv")); - + File[] files = directory.listFiles((dir, name) -> name.endsWith(".csv")); + if (files == null || files.length == 0) { log.error("未找到API {} 的CSV文件", api); throw new FileNotFoundException("未找到API " + api + " 的CSV文件"); } - - // 选择最新的文件 - File csvFile = Arrays.stream(files) - .max(Comparator.comparingLong(File::lastModified)) - .orElse(null); - - if (csvFile == null) { - log.error("未找到API {} 的CSV文件", api); - throw new FileNotFoundException("未找到API " + api + " 的CSV文件"); - } - - String filePath = csvFile.getAbsolutePath(); - log.info("找到CSV文件: {}", filePath); - + log.info("开始导入BigObject数据到目标ORG: {}", api); - - // 创建Bulk API作业 - JobInfo job = BulkUtil.createJob(targetBulkConnect, api, OperationEnum.insert); - log.info("创建Bulk作业成功, 作业ID: {}", job.getId()); - - // 从CSV文件创建批次 - List batchInfos = BulkUtil.createBatchesFromCSVFile(targetBulkConnect, job, filePath); - log.info("创建批次成功, 批次数: {}", batchInfos.size()); - - // 等待作业完成 - BulkUtil.awaitCompletion(targetBulkConnect, job, batchInfos); - log.info("数据导入完成"); - - // 检查作业状态 - job = targetBulkConnect.getJobStatus(job.getId()); - if (job.getState() == JobStateEnum.Failed) { - log.error("Bulk作业失败, 作业ID: {}", job.getId()); - throw new RuntimeException("Bulk作业失败 "); + + // 使用Bulk V2 API导入所有文件 + for (File csvFile : files) { + String filePath = csvFile.getAbsolutePath(); + log.info("处理CSV文件: {}", filePath); + + // 创建Bulk V2 API作业 + String restEndpoint = targetBulkConnect.getConfig().getRestEndpoint(); + String url = restEndpoint.substring(0, restEndpoint.indexOf("services/")); + + // 创建Bulk V2作业 + String jobId = BulkUtil.createBulkV2InsertJob(url, api, targetBulkConnect); + log.info("创建Bulk V2作业成功, 作业ID: {}", jobId); + + // 上传CSV文件数据 + uploadCsvDataToJob(url, jobId, targetBulkConnect, filePath); + + // 关闭作业 + closeBulkV2Job(url, jobId, targetBulkConnect); + + // 等待作业完成 + BulkUtil.waitForBulkV2JobCompletion(url, jobId, targetBulkConnect); + log.info("文件 {} 数据导入完成", csvFile.getName()); } + + log.info("所有BigObject数据导入完成: {}", api); + return ReturnT.SUCCESS; + } + + /** + * 上传CSV数据到Bulk API 2.0作业 + * + * @param url Salesforce实例URL + * @param jobId 作业ID + * @param connection BulkConnection连接对象 + * @param csvFilePath CSV文件路径 + * @throws Exception 上传过程中发生错误时抛出异常 + */ + private void uploadCsvDataToJob(String url, String jobId, BulkConnection connection, String csvFilePath) throws Exception { + CloseableHttpClient httpClient = HttpClients.createDefault(); + String uploadUrl = url + "services/data/v56.0/jobs/ingest/" + jobId + "/batches/"; + HttpPut httpPut = new HttpPut(uploadUrl); + httpPut.setHeader("Content-Type", "text/csv"); + httpPut.setHeader("Authorization", "Bearer " + connection.getConfig().getSessionId()); - // 关闭作业 - BulkUtil.closeJob(targetBulkConnect, job.getId()); - log.info("关闭作业成功, 作业ID: {}", job.getId()); + FileEntity fileEntity = new FileEntity(new File(csvFilePath)); + httpPut.setEntity(fileEntity); - // 记录导入统计信息 - int totalRecords = batchInfos.stream() - .mapToInt(BatchInfo::getNumberRecordsProcessed) - .sum(); - log.info("成功导入 {} 条记录到 {}", totalRecords, api); + CloseableHttpResponse response = httpClient.execute(httpPut); + try { + int statusCode = response.getStatusLine().getStatusCode(); + if (statusCode != 201) { // 201 Created + String errorBody = EntityUtils.toString(response.getEntity()); + throw new RuntimeException("Failed to upload data to Bulk API 2.0 job, status code: " + statusCode + ", error: " + errorBody); + } + } finally { + response.close(); + httpClient.close(); + } + log.info("数据上传成功,作业ID: {}", jobId); + } + + /** + * 关闭Bulk API 2.0作业 + * + * @param url Salesforce实例URL + * @param jobId 作业ID + * @param connection BulkConnection连接对象 + * @throws Exception 关闭作业时发生错误时抛出异常 + */ + private void closeBulkV2Job(String url, String jobId, BulkConnection connection) throws Exception { + CloseableHttpClient httpClient = HttpClients.createDefault(); + String closeUrl = url + "services/data/v56.0/jobs/ingest/" + jobId; + HttpPatch httpPatch = new HttpPatch(closeUrl); + httpPatch.setHeader("Content-Type", "application/json"); + httpPatch.setHeader("Authorization", "Bearer " + connection.getConfig().getSessionId()); + + // 构建关闭作业的 JSON 请求体 + String jsonRequest = "{\"state\":\"UploadComplete\"}"; + StringEntity stringEntity = new StringEntity(jsonRequest, "UTF-8"); + httpPatch.setEntity(stringEntity); + + CloseableHttpResponse response = httpClient.execute(httpPatch); + try { + int statusCode = response.getStatusLine().getStatusCode(); + if (statusCode != 200) { + String errorBody = EntityUtils.toString(response.getEntity()); + throw new RuntimeException("Failed to close Bulk API 2.0 job, status code: " + statusCode + ", error: " + errorBody); + } + } finally { + response.close(); + httpClient.close(); + } + log.info("作业关闭成功,作业ID: {}", jobId); } private ReturnT manualBatchDump(SalesforceParam param, List> futures) throws InterruptedException { @@ -853,7 +870,7 @@ public class CommonBatchServiceImpl implements CommonBatchService { String url = restEndpoint.substring(0, restEndpoint.indexOf("services/")); - v2JobId = BulkUtil.createBulkV2Job(url,sql, bulkConnect); + v2JobId = BulkUtil.createBulkV2QueryJob(url,sql, bulkConnect); BulkUtil.waitForBulkV2JobCompletion(url,v2JobId, bulkConnect); @@ -934,7 +951,7 @@ public class CommonBatchServiceImpl implements CommonBatchService { } } - }catch (Exception e){ + }catch (Exception e){ log.error("processQueryResults error", e); } return totalProcessedRecords; diff --git a/src/main/java/com/celnet/datadump/service/impl/CommonServiceImpl.java b/src/main/java/com/celnet/datadump/service/impl/CommonServiceImpl.java index c0dbd30..4dd4e33 100644 --- a/src/main/java/com/celnet/datadump/service/impl/CommonServiceImpl.java +++ b/src/main/java/com/celnet/datadump/service/impl/CommonServiceImpl.java @@ -1607,64 +1607,51 @@ public class CommonServiceImpl implements CommonService { public ReturnT createBigObjectApi(SalesforceParam param) throws Exception { String api = param.getApi(); + log.info("开始创建BigObject API: {}", api); PartnerConnection connection = salesforceConnect.createConnect(); + log.debug("Salesforce连接创建成功"); DescribeSObjectResult dsr = connection.describeSObject(api); + log.info("获取对象描述信息成功, API: {}, Label: {}", api, dsr.getLabel()); + + + // 查询data_field表,如果有数据则不执行,没有数据执行 + QueryWrapper queryWrapper = new QueryWrapper<>(); + queryWrapper.eq("api", api); + long count = dataFieldService.count(queryWrapper); + + if (count == 0) { + List fieldList = Lists.newArrayList(); + log.debug("开始处理字段信息,字段总数: {}", dsr.getFields().length); + for (Field field : dsr.getFields()) { + DataField dataField = new DataField(); + dataField.setApi(api); + dataField.setSfType(field.getType().toString()); + dataField.setField(field.getName()); + dataField.setName(field.getLabel()); + dataField.setIsCreateable(field.getCreateable()); + dataField.setIsUpdateable(field.getUpdateable()); + dataField.setIsNillable(field.getNillable()); + dataField.setIsDefaultedOnCreate(field.getDefaultedOnCreate()); + fieldList.add(dataField); + } + log.info("字段信息处理完成,共处理{}个字段", fieldList.size()); + + dataFieldService.saveBatch(fieldList); + log.info("字段信息批量保存完成"); + } else { + log.info("API {} 的字段信息已存在,跳过字段处理", api); + } DataObject dataObject = dataObjectService.getById(api); dataObject.setLabel(dsr.getLabel()); + dataObject.setLastUpdateDate(new Date()); dataObjectService.updateById(dataObject); - - Date startDate = dataObject.getLastUpdateDate() == null ? DataUtil.DEFAULT_BEGIN_DATE : dataObject.getLastUpdateDate(); - Date endCreateDate = DateUtils.parseDate(DateFormatUtils.format(new Date(), "yyyy-MM-dd"), "yyyy-MM-dd"); + log.info("数据对象信息更新完成"); - // 按年生成批次 - DataBatch one = new DataBatch(); - one.setFirstDbNum(0); - one.setFirstSfNum(0); - one.setDbNum(0); - one.setSfNum(0); - one.setName(api); - one.setLabel(dsr.getLabel()); - - // 按年份划分开始时间和结束时间创建批次 - Calendar startCalendar = Calendar.getInstance(); - startCalendar.setTime(startDate); - int startYear = startCalendar.get(Calendar.YEAR); - - Calendar endCalendar = Calendar.getInstance(); - endCalendar.setTime(endCreateDate); - int endYear = endCalendar.get(Calendar.YEAR); - - // 优化批次创建逻辑,确保正确的日期范围 - for (int year = startYear; year <= endYear; year++) { - Date batchStartDate; - Date batchEndDate; - - // 设置批次开始时间为当年1月1日 00:00:00 - batchStartDate = DateUtils.parseDate(year + "-01-01 00:00:00", "yyyy-MM-dd HH:mm:ss"); - - if (year == endYear) { - // 结束年份使用实际结束日期 - batchEndDate = endCreateDate; - } else { - int nextYear = year + 1; - // 其他年份结束时间为当年12月31日 00:00:00 - batchEndDate = DateUtils.parseDate(nextYear + "-01-01 00:00:00", "yyyy-MM-dd HH:mm:ss"); - } - - // 确保批次开始时间不早于实际开始时间 - if (batchStartDate.before(startDate)) { - batchStartDate = startDate; - } - - DataBatch dataBatch = one.clone(); - dataBatch.setSyncStartDate(batchStartDate); - dataBatch.setSyncEndDate(batchEndDate); - dataBatchService.save(dataBatch); - } + log.info("BigObject API创建完成: {}", api); return ReturnT.SUCCESS; } diff --git a/src/main/java/com/celnet/datadump/service/impl/DataFieldServiceImpl.java b/src/main/java/com/celnet/datadump/service/impl/DataFieldServiceImpl.java index 12e5ee2..0043868 100644 --- a/src/main/java/com/celnet/datadump/service/impl/DataFieldServiceImpl.java +++ b/src/main/java/com/celnet/datadump/service/impl/DataFieldServiceImpl.java @@ -100,7 +100,7 @@ public class DataFieldServiceImpl extends ServiceImpl> referenceMap = customMapper.list("new_id", reference, "new_id is not null limit 1"); if (referenceMap.isEmpty()){ + log.info("引用对象:" + reference + "不存在数据!请检查这个对象详情!!!"); QueryWrapper maxIndex = new QueryWrapper<>(); maxIndex.select("IFNULL(max(data_index),0) as data_index"); maxIndex.ne("name", api); @@ -386,7 +387,6 @@ public class DataImportServiceImpl implements DataImportService { //如果必填lookup字段没有值,跳过 dataObject.setDataIndex(Integer.parseInt(map.get("data_index").toString()) + 1); dataObjectService.updateById(dataObject); - log.info("引用对象:" + reference + "不存在数据!请检查这个对象详情!!!"); return; }else{ account.setField(dataField.getField(), referenceMap.get(0).get("new_id")); diff --git a/src/main/java/com/celnet/datadump/util/BulkUtil.java b/src/main/java/com/celnet/datadump/util/BulkUtil.java index ad6ed12..5a9541b 100644 --- a/src/main/java/com/celnet/datadump/util/BulkUtil.java +++ b/src/main/java/com/celnet/datadump/util/BulkUtil.java @@ -171,11 +171,13 @@ public class BulkUtil { /** * 创建Bulk API 2.0查询作业 - ** @param soqlQuery SOQL查询语句 + * @param url Salesforce实例URL + * @param soqlQuery SOQL查询语句 + * @param connection BulkConnection连接对象 * @return 作业ID * @throws Exception 创建作业时发生错误 */ - public static String createBulkV2Job(String url,String soqlQuery,BulkConnection connection) throws Exception { + public static String createBulkV2QueryJob(String url, String soqlQuery, BulkConnection connection) throws Exception { // 构建创建查询作业的 JSON 请求体 String jsonRequest = String.format("{\"operation\":\"query\", \"query\":\"%s\"}", soqlQuery.replace("\"", "\\\"")); @@ -192,9 +194,7 @@ public class BulkUtil { // 从响应体中解析出 jobId String responseBody = EntityUtils.toString(response.getEntity()); // 简单的JSON解析提取id,实际项目中建议使用JSON库如Jackson - String jobId = responseBody.split("\"id\":\"")[1].split("\"")[0]; - System.out.println("Bulk API 2.0 Query job created with ID: " + jobId); - return jobId; + return responseBody.split("\"id\":\"")[1].split("\"")[0]; } else { String errorBody = EntityUtils.toString(response.getEntity()); throw new RuntimeException("Failed to create Bulk API 2.0 query job, status code: " + statusCode + ", error: " + errorBody); @@ -206,11 +206,11 @@ public class BulkUtil { } /** - * Wait for a Bulk API 2.0 job to complete by polling the job status. + * 通过轮询作业状态等待 Bulk API 2.0 作业完成。 * - * @param jobId The ID of the job to check. - * @param connection BulkConnection used to get session info. - * @throws Exception If an error occurs during the polling process. + * @param jobId 要检查的作业ID。 + * @param connection 用于获取会话信息的BulkConnection。 + * @throws Exception 轮询过程中发生错误时抛出异常。 */ public static void waitForBulkV2JobCompletion(String url, String jobId, BulkConnection connection) throws Exception { CloseableHttpClient httpClient = HttpClients.createDefault(); @@ -228,17 +228,17 @@ public class BulkUtil { String responseBody = EntityUtils.toString(response.getEntity()); // 简单解析响应体中的状态字段,实际项目中建议使用JSON库如Jackson if (responseBody.contains("\"state\":\"JobComplete\"")) { - System.out.println("Job completed successfully."); + System.out.println("作业成功完成。"); jobCompleted = true; } else if (responseBody.contains("\"state\":\"Failed\"") || responseBody.contains("\"state\":\"Aborted\"")) { - throw new RuntimeException("Job failed or was aborted. Response: " + responseBody); + throw new RuntimeException("作业失败或已被中止。响应内容: " + responseBody); } else { - System.out.println("Job is still running. Waiting..."); + System.out.println("作业仍在运行中。等待中..."); Thread.sleep(10000); // 等待10秒后再次检查 } } else { String errorBody = EntityUtils.toString(response.getEntity()); - throw new RuntimeException("Failed to get job status, status code: " + statusCode + ", error: " + errorBody); + throw new RuntimeException("获取作业状态失败,状态码: " + statusCode + ",错误: " + errorBody); } } finally { response.close(); @@ -256,10 +256,10 @@ public class BulkUtil { * @return 查询结果的字符串内容 * @throws Exception 如果在获取结果过程中发生错误 */ - public static String getBulkV2QueryJobResults(String url, String jobId, BulkConnection connection) throws Exception { + public static BulkV2QueryResult getBulkV2QueryJobResults(String url, String jobId, BulkConnection connection) throws Exception { CloseableHttpClient httpClient = HttpClients.createDefault(); - String resultsUrl = url + "/services/data/v65.0/jobs/query/" + jobId + "/results" ; - HttpPost httpGet = new HttpPost(resultsUrl); + String resultsUrl = url + "/services/data/v56.0/jobs/query/" + jobId + "/results"; + HttpGet httpGet = new HttpGet(resultsUrl); httpGet.setHeader("Authorization", "Bearer " + connection.getConfig().getSessionId()); httpGet.setHeader("Content-Type", "application/json"); httpGet.setHeader("Accept", "text/csv"); @@ -269,13 +269,12 @@ public class BulkUtil { try { int statusCode = response.getStatusLine().getStatusCode(); if (statusCode == 200) { - Header locatorHeader = response.getFirstHeader("Sforce-Locator"); - Header numRecordsHeader = response.getFirstHeader("Sforce-NumberOfRecords"); + Header numRecordsHeader = response.getFirstHeader("Sforce-NumberOfRecords"); String newLocator = (locatorHeader != null) ? locatorHeader.getValue() : null; int numberOfRecords = (numRecordsHeader != null) ? Integer.parseInt(numRecordsHeader.getValue()) : 0; - - return EntityUtils.toString(response.getEntity()); + String csvContent = EntityUtils.toString(response.getEntity()); + return new BulkV2QueryResult(csvContent, newLocator, numberOfRecords); } else { String errorBody = EntityUtils.toString(response.getEntity()); throw new RuntimeException("Failed to get query job results, status code: " + statusCode + ", error: " + errorBody); @@ -286,6 +285,41 @@ public class BulkUtil { } } + /** + * 创建Bulk API 2.0插入作业 + * @param url Salesforce实例URL + * @param objectName 对象名称 + * @param connection BulkConnection连接对象 + * @return 作业ID + * @throws Exception 创建作业时发生错误 + */ + public static String createBulkV2InsertJob(String url, String objectName, BulkConnection connection) throws Exception { + // 构建创建插入作业的 JSON 请求体 + String jsonRequest = String.format("{\"operation\":\"insert\", \"object\":\"%s\", \"contentType\":\"CSV\"}", objectName); + + CloseableHttpClient httpClient = HttpClients.createDefault(); + HttpPost httpPost = new HttpPost(url + "services/data/v56.0/jobs/ingest/"); + httpPost.setHeader("Content-Type", "application/json"); + httpPost.setHeader("Authorization", "Bearer " + connection.getConfig().getSessionId()); + httpPost.setEntity(new StringEntity(jsonRequest, "UTF-8")); + + CloseableHttpResponse response = httpClient.execute(httpPost); + try { + int statusCode = response.getStatusLine().getStatusCode(); + if (statusCode == 200) { + // 从响应体中解析出 jobId + String responseBody = EntityUtils.toString(response.getEntity()); + // 简单的JSON解析提取id,实际项目中建议使用JSON库如Jackson + return responseBody.split("\"id\":\"")[1].split("\"")[0]; + } else { + String errorBody = EntityUtils.toString(response.getEntity()); + throw new RuntimeException("Failed to create Bulk API 2.0 insert job, status code: " + statusCode + ", error: " + errorBody); + } + } finally { + response.close(); + httpClient.close(); + } + } /** diff --git a/src/main/java/com/celnet/datadump/util/DataUtil.java b/src/main/java/com/celnet/datadump/util/DataUtil.java index 57bd415..5eb4c6b 100644 --- a/src/main/java/com/celnet/datadump/util/DataUtil.java +++ b/src/main/java/com/celnet/datadump/util/DataUtil.java @@ -455,7 +455,7 @@ public class DataUtil { case "datetime": return date; default: - return "初始"+ UUID.randomUUID(); + return "初始"+ UUID.randomUUID().toString().substring(0, 8); } }