From bc1418599353465076500ff9085ccf7b9b978c40 Mon Sep 17 00:00:00 2001 From: Kris <2893855659@qq.com> Date: Wed, 10 Sep 2025 22:10:36 +0800 Subject: [PATCH] [[feat] 20250910 --- .../celnet/datadump/mapper/CustomMapper.java | 7 +- .../service/impl/CommonBatchServiceImpl.java | 118 ++++++++++++++- .../service/impl/CommonServiceImpl.java | 140 ++++++++++++++++-- .../impl/DataImportBatchServiceImpl.java | 2 +- .../impl/DataImportNewServiceImpl.java | 24 +-- .../service/impl/DataImportServiceImpl.java | 8 +- .../service/impl/FileServiceImpl.java | 4 +- src/main/resources/mapper/CustomMapper.xml | 3 + .../resources/mapper/SalesforceMapper.xml | 3 + 9 files changed, 279 insertions(+), 30 deletions(-) diff --git a/src/main/java/com/celnet/datadump/mapper/CustomMapper.java b/src/main/java/com/celnet/datadump/mapper/CustomMapper.java index a05895f..da35166 100644 --- a/src/main/java/com/celnet/datadump/mapper/CustomMapper.java +++ b/src/main/java/com/celnet/datadump/mapper/CustomMapper.java @@ -55,8 +55,11 @@ public interface CustomMapper { * @param tableName tableName * @param maps maps */ - public void createTable(@Param("tableName") String tableName, @Param("tableComment") String tableComment, @Param("maps") List> maps, @Param("index") List> index); - + void createTable(@Param("tableName") String tableName, + @Param("tableComment") String tableComment, + @Param("maps") List> maps, + @Param("index") List> index, + @Param("hasId") boolean hasId); /** * 创建RANGE分区表 * 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 66fa038..13ceb4e 100644 --- a/src/main/java/com/celnet/datadump/service/impl/CommonBatchServiceImpl.java +++ b/src/main/java/com/celnet/datadump/service/impl/CommonBatchServiceImpl.java @@ -2,7 +2,10 @@ package com.celnet.datadump.service.impl; import com.alibaba.fastjson2.JSON; +import com.alibaba.fastjson2.JSONArray; +import com.alibaba.fastjson2.JSONObject; import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper; +import com.baomidou.mybatisplus.core.conditions.update.UpdateWrapper; import com.celnet.datadump.config.SalesforceConnect; import com.celnet.datadump.config.SalesforceExecutor; import com.celnet.datadump.entity.*; @@ -117,7 +120,7 @@ public class CommonBatchServiceImpl implements CommonBatchService { if (!dataFieldService.hasCreatedDate(api)) { DataDumpSpecialParam dataDumpSpecialParam = new DataDumpSpecialParam(); dataDumpSpecialParam.setApi(api); - Future future = dataDumpSpecialService.getData(dataDumpSpecialParam, salesforceConnect.createConnect()); + Future future = getData(dataDumpSpecialParam, salesforceConnect.createBulkConnect()); // 等待当前所有线程执行完成 salesforceExecutor.waitForFutures(future); continue; @@ -222,6 +225,119 @@ public class CommonBatchServiceImpl implements CommonBatchService { return bulkConnect; } + public Future getData(DataDumpSpecialParam param, BulkConnection bulkConnect) { + String api = param.getApi(); + return salesforceExecutor.execute(() -> { + try { + // 检测表 不存在就生成 但不生成批次 + commonService.checkApi(api, true); + int count = 0; + boolean hasMore = true; + QueryWrapper wrapper = new QueryWrapper<>(); + wrapper.eq("api", api); + List list = dataFieldService.list(wrapper); + + Map map = Maps.newHashMap(); + SalesforceParam salesforceParam = new SalesforceParam(); + salesforceParam.setApi(api); + salesforceParam.setIdField(param.getField()); + String maxId = null; + + List fields = Lists.newArrayList(); + for (DataField field : list) { + if (!"base64".equalsIgnoreCase(field.getSfType()) && field.getSfType() != null){ + fields.add(field.getField()); + } + } + salesforceParam.setSelect(StringUtils.join(fields, ",")); + + while (hasMore) { + int sfData = 0; + salesforceParam.setMaxId(maxId); + salesforceParam.setLimit(10000); + map.put("param", salesforceParam); + JobInfo job = null; + try { + String sql = SqlUtil.showSql("com.celnet.datadump.mapper.SalesforceMapper.listOrderByIdNew", map); + + log.info("query sql: {}", sql); + XxlJobLogger.log("query 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()); + break; + } + + List> batchRecords = new ArrayList<>(); + + QueryResultList queryResultList = bulkConnect.getQueryResultList(job.getId(), batchInfo.getId()); + log.info("获取查询结果列表, 结果数量: {}", queryResultList.getResult().length); + + for (final String resultId : queryResultList.getResult()) { + InputStream resultStream = bulkConnect.getQueryResultStream(job.getId(), batchInfo.getId(), resultId); + + CSVReader rdr = new CSVReader(resultStream); + List headers = rdr.nextRecord(); + log.debug("处理结果ID: {}, 表头数量: {}", resultId, headers.size()); + + // 批量处理记录 + List record; + while ((record = rdr.nextRecord()) != null) { + Map recordMap = new HashMap<>(); + for (int i = 0; i < headers.size() && i < record.size(); i++) { + recordMap.put(headers.get(i), record.get(i)); + } + batchRecords.add(recordMap); + } + } + if (!batchRecords.isEmpty()){ + log.info("开始保存或更新数据, 记录数: {}", batchRecords.size()); + sfData = saveOrUpdate(api, batchRecords, true); + count += sfData; + maxId = batchRecords.get(batchRecords.size() - 1).get("Id").toString(); + log.info("处理结果: api:{}, sf数量:{}, 当前批次最大Id:{}", api, sfData, maxId); + } else { + log.info("批次无数据记录"); + } + + BulkUtil.closeJob(bulkConnect, job.getId()); + log.info("关闭作业成功, 作业ID: {}", job.getId()); + }catch (InterruptedException interruptedException){ + return ; + } catch (Exception e) { + log.error("处理查询结果时异常: ", e); + throw new AsyncApiException("处理查询结果时异常: " + e.getMessage(), AsyncExceptionCode.Unknown); + }finally { + if (sfData != 10000){ + hasMore = false; + log.info("当前批次数据不足10000条, 结束循环"); + } + } + } + UpdateWrapper updateWrapper = new UpdateWrapper<>(); + updateWrapper.eq("name", param.getApi()); + updateWrapper.set("first_db_num", count); + updateWrapper.set("first_sf_num", count); + dataBatchService.update(updateWrapper); + + } catch (Throwable throwable) { + String format = String.format("数据特殊表迁移 error, api name: %s, \nparam: %s, \ncause:\n%s", api, JSON.toJSONString(param), throwable); + EmailUtil.send("DataDump ERROR", format); + log.error("dataDumpSpecial error ", throwable); + throw new RuntimeException(throwable); + } + }, 0, 0); + } private int getAllBulkV1SfData(SalesforceParam param, BulkConnection bulkConnect) throws Throwable { 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 9935c95..0fd293d 100644 --- a/src/main/java/com/celnet/datadump/service/impl/CommonServiceImpl.java +++ b/src/main/java/com/celnet/datadump/service/impl/CommonServiceImpl.java @@ -85,6 +85,8 @@ public class CommonServiceImpl implements CommonService { private LinkConfigService linkConfigService; @Autowired private DataLogService dataLogService; + @Autowired + private CommonService commonService; @Override public ReturnT increment(SalesforceParam param) throws Exception { @@ -1019,7 +1021,7 @@ public class CommonServiceImpl implements CommonService { } } // picklist保存到picklist表 - if ("picklist".equalsIgnoreCase(sfType)) { + if ("picklist".equalsIgnoreCase(sfType) || "multipicklist".equalsIgnoreCase(sfType)) { join = "data_picklist"; PicklistEntry[] picklistValues = field.getPicklistValues(); if (ArrayUtils.isNotEmpty(picklistValues)) { @@ -1122,7 +1124,9 @@ public class CommonServiceImpl implements CommonService { map4.put("name", "error_message"); list.add(map4); - customMapper.createTable(apiName, label, list, index); + boolean hasId = list.stream().anyMatch(t -> "Id".equals(t.get("name").toString())); + + customMapper.createTable(apiName, label, list, index,hasId); // 生成字段映射 QueryWrapper delfieldQw = new QueryWrapper<>(); delfieldQw.eq("api", apiName); @@ -1137,7 +1141,7 @@ public class CommonServiceImpl implements CommonService { Date startDate = DataUtil.DEFAULT_BEGIN_DATE; // 取当天零点 作为最大时间 Date endCreateDate = DateUtils.parseDate(DateFormatUtils.format(now, "yyyy-MM-dd"), "yyyy-MM-dd"); - Date lastDay = null; + QueryWrapper delQw = new QueryWrapper<>(); delQw.eq("name", apiName); dataBatchService.remove(delQw); @@ -1153,14 +1157,7 @@ public class CommonServiceImpl implements CommonService { if (hasCreatedDate) { List> bachList = customMapper.list("code,value","system_config","code ='"+BATCH_TYPE+"'"); String batchType = bachList.get(0).get("value").toString(); - do { - lastDay = getLastDay(batchType, endCreateDate, startDate); - DataBatch dataBatch = one.clone(); - dataBatch.setSyncStartDate(startDate); - dataBatch.setSyncEndDate(lastDay); - dataBatchService.save(dataBatch); - startDate = lastDay; - } while (lastDay.compareTo(endCreateDate) < 0); + saveBatch(batchType,startDate,endCreateDate,apiName,connection,one); } else { // 没创建时间 开始结束时间为空 dataBatchService.save(one); @@ -1182,6 +1179,127 @@ public class CommonServiceImpl implements CommonService { } } + /** + * 保存批次信息,根据数据量自动调整批次粒度 + * @param batchType 批次类型(年、月、周) + * @param startDate 开始日期 + * @param endDate 结束日期 + * @param apiName API名称 + * @param connection Salesforce连接 + * @param one 批次对象模板 + */ + public void saveBatch(String batchType, Date startDate, Date endDate, String apiName, PartnerConnection connection, DataBatch one) { + // 添加递归深度限制,避免栈溢出 + saveBatchRecursive(batchType, startDate, endDate, apiName, connection, one, 0); + log.info("API {} 批次创建完成", apiName); + } + + /** + * 递归保存批次信息 + * @param batchType 批次类型 + * @param startDate 开始日期 + * @param endDate 结束日期 + * @param apiName API名称 + * @param connection Salesforce连接 + * @param one 批次对象模板 + * @param depth 递归深度 + */ + private void saveBatchRecursive(String batchType, Date startDate, Date endDate, String apiName, PartnerConnection connection, DataBatch one, int depth) { + // 限制递归深度,避免栈溢出 + if (depth > 10) { + log.warn("递归深度超过限制,直接保存批次: {}, startDate: {}, endDate: {}", apiName, startDate, endDate); + DataBatch dataBatch = one.clone(); + dataBatch.setSyncStartDate(startDate); + dataBatch.setSyncEndDate(endDate); + dataBatchService.save(dataBatch); + return; + } + + Date currentStartDate = startDate; + Date lastDay; + + do { + // 根据批次类型获取结束日期 + lastDay = getLastDay(batchType, endDate, currentStartDate); + log.debug("当前批次起始日期: {}, 计算得到的结束日期: {}", currentStartDate, lastDay); + + // 添加边界检查,避免无限循环 + if (lastDay.compareTo(currentStartDate) <= 0) { + log.warn("计算得到的结束日期({})早于或等于起始日期({}),跳过该批次", lastDay, currentStartDate); + break; + } + + // 查询salesforce当前批次数据量 + SalesforceParam salesforceParam = new SalesforceParam(); + salesforceParam.setApi(apiName); + salesforceParam.setBeginCreateDate(currentStartDate); + salesforceParam.setEndCreateDate(lastDay); + Integer totalNum = 0; + int retryCount = 0; + int maxRetries = 3; + + // 增强异常处理,添加重试机制 + while (retryCount < maxRetries) { + try { + totalNum = commonService.countSfNum(connection, salesforceParam); + log.debug("API {} 在 {} 到 {} 期间的数据量为: {}", apiName, currentStartDate, lastDay, totalNum); + break; // 成功获取数据量,跳出重试循环 + } catch (Exception e) { + retryCount++; + log.error("api {} count error, retry {}/{}", apiName, retryCount, maxRetries, e); + if (retryCount >= maxRetries) { + log.error("api {} count failed after {} retries", apiName, maxRetries, e); + // 如果重试失败,使用默认值继续处理 + totalNum = 0; + } else { + try { + Thread.sleep(1000 * retryCount); // 指数退避 + } catch (InterruptedException ie) { + Thread.currentThread().interrupt(); + throw new RuntimeException("Thread interrupted during retry delay", ie); + } + } + } + } + + // 如果当前批次数据量超过50万,则细化批次粒度并递归处理 + if (totalNum > 500000) { + log.info("API {} 当前批次数据量 {} 超过50万,需要细化批次粒度", apiName, totalNum); + if (BATCH_TYPE_YEAR.equals(batchType)) { + // 年批次超过50万,细化为月批次 + log.info("将年批次细化为月批次"); + saveBatchRecursive(BATCH_TYPE_MONTH, currentStartDate, lastDay, apiName, connection, one, depth + 1); + } else if (BATCH_TYPE_MONTH.equals(batchType)) { + // 月批次超过50万,细化为周批次 + log.info("将月批次细化为周批次"); + saveBatchRecursive(BATCH_TYPE_WEEK, currentStartDate, lastDay, apiName, connection, one, depth + 1); + } else { + // 已经是周粒度 still > 500000, 直接保存当前批次 + log.info("已经是周粒度但数据量仍超过50万,直接保存批次"); + DataBatch dataBatch = one.clone(); + dataBatch.setSyncStartDate(currentStartDate); + dataBatch.setSyncEndDate(lastDay); + dataBatch.setSfNum(totalNum); + dataBatchService.save(dataBatch); + } + } else { + // 数据量未超过50万,直接保存当前批次 + log.info("API {} 当前批次数据量 {} 未超过50万,直接保存批次", apiName, totalNum); + DataBatch dataBatch = one.clone(); + dataBatch.setSyncStartDate(currentStartDate); + dataBatch.setSyncEndDate(lastDay); + dataBatch.setSfNum(totalNum); + dataBatchService.save(dataBatch); + } + + currentStartDate = lastDay; + + // 添加安全检查,避免无限循环 + if (currentStartDate.compareTo(endDate) >= 0) { + break; + } + } while (lastDay.compareTo(endDate) < 0); + } /** * 创建分区表结构 * @param apiName 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 f48ae77..6b8ea60 100644 --- a/src/main/java/com/celnet/datadump/service/impl/DataImportBatchServiceImpl.java +++ b/src/main/java/com/celnet/datadump/service/impl/DataImportBatchServiceImpl.java @@ -400,7 +400,7 @@ public class DataImportBatchServiceImpl implements DataImportBatchService { continue; } } - if ("picklist".equals(dataField.getSfType())){ + if ("picklist".equals(dataField.getSfType()) || "multipicklist".equals(dataField.getSfType())){ List> pickList = customMapper.list("value", "data_picklist", "api = '"+api+"' and field = '"+dataField.getField()+"' limit 1"); account.put(dataField.getField(), pickList.get(0).get("value")); continue; diff --git a/src/main/java/com/celnet/datadump/service/impl/DataImportNewServiceImpl.java b/src/main/java/com/celnet/datadump/service/impl/DataImportNewServiceImpl.java index 70146ca..9a1985e 100644 --- a/src/main/java/com/celnet/datadump/service/impl/DataImportNewServiceImpl.java +++ b/src/main/java/com/celnet/datadump/service/impl/DataImportNewServiceImpl.java @@ -739,25 +739,26 @@ public class DataImportNewServiceImpl implements DataImportNewService { if (1 == param.getType()) { if (api.endsWith("Share")){ sql = "where RowCause = 'Manual' and is_update = 0 and new_id is not null and CreatedDate >= '" + beginDateStr + "' and CreatedDate < '" + endDateStr + "'"; - sql2 = "RowCause = 'Manual' and is_update = 0 and new_id is not null and CreatedDate >= '" + beginDateStr + "' and CreatedDate < '" + endDateStr + "' order by Id asc limit "; + sql2 = "RowCause = 'Manual' and is_update = 0 and new_id is not null and CreatedDate >= '" + beginDateStr + "' and CreatedDate < '" + endDateStr + "' order by Id asc limit 10000 "; }else { sql = "where new_id is not null and is_update = 0 and CreatedDate >= '" + beginDateStr + "' and CreatedDate < '" + endDateStr + "'"; - sql2 = "new_id is not null and is_update = 0 and CreatedDate >= '" + beginDateStr + "' and CreatedDate < '" + endDateStr + "' order by Id asc limit "; + sql2 = "new_id is not null and is_update = 0 and CreatedDate >= '" + beginDateStr + "' and CreatedDate < '" + endDateStr + "' order by Id asc limit 10000 "; } }else { String updateDateField = dataFieldService.returnUpdateDateField(api); if (api.endsWith("Share")){ sql = "where RowCause = 'Manual' and is_update = 0 and new_id is not null and "+updateDateField+" >= '" + beginDateStr + "' "; - sql2 = "RowCause = 'Manual' and is_update = 0 and new_id is not null and "+updateDateField+" >= '" + beginDateStr + "' order by Id asc limit "; + sql2 = "RowCause = 'Manual' and is_update = 0 and new_id is not null and "+updateDateField+" >= '" + beginDateStr + "' order by Id asc limit 10000 "; }else { - sql = "where new_id is not null and and is_update = 0 "+updateDateField+" >= '" + beginDateStr + "' "; - sql2 = "new_id is not null and is_update = 0 and "+updateDateField+" >= '" + beginDateStr + "' order by Id asc limit "; + sql = "where new_id is not null and is_update = 0 and "+updateDateField+" >= '" + beginDateStr + "' "; + sql2 = "new_id is not null and is_update = 0 and "+updateDateField+" >= '" + beginDateStr + "' order by Id asc limit 10000 "; } } //表内数据总量 Integer count = customMapper.countBySQL(api, sql); if(count == 0){ + log.info("没有需要更新的数据,API: {},开始时间: {},结束时间: {}", api, beginDateStr, endDateStr); return; } @@ -784,19 +785,20 @@ public class DataImportNewServiceImpl implements DataImportNewService { for (int z = 0; z < num; z++){ - List> mapsList = customMapper.list("*", api, sql2+ z * 10000 + ",10000"); + List> mapsList = customMapper.list("*", api, sql2); int size = mapsList.size(); - log.info("总Update数据 count:{};当前批次:{};-开始时间:{};-结束时间:{};-api:{};", count,z, beginDateStr, endDateStr, api); + log.info("总批次:{};当前批次:{};当前批次数据量:{};-开始时间:{};-结束时间:{};-api:{};", num,z,size, beginDateStr, endDateStr, api); //批量插入200一次 - int page = size%200 == 0 ? count/200 : (count/200) + 1; + int page = size%200 == 0 ? size/200 : (size/200) + 1; for (int i = 0; i < page; i++) { - DataLog dataLog = new DataLog(api, "开始时间:"+beginDateStr+";结束时间:"+endDateStr, new Date(), null, "数据更新,查询并组装第" + z*50 + i + "批数据", "QueryDB"); + DataLog dataLog = new DataLog(api, "开始时间:"+beginDateStr+";结束时间:"+endDateStr, new Date(), null, "数据更新,查询并组装第" + (z*50 + i) + "批数据", "QueryDB"); int startIndex = 200 * i; int endIndex = Math.min(startIndex + 200, size); if (startIndex >= endIndex) { + log.info("当前批次:{} 结束!! -开始时间:{};-结束时间:{};-api:{};", z, beginDateStr, endDateStr, api); break; } List> mapList = mapsList.subList(startIndex, endIndex); @@ -880,7 +882,7 @@ public class DataImportNewServiceImpl implements DataImportNewService { if (infoFlag != null && "1".equals(infoFlag.get("value"))){ printlnAccountsDetails(accounts,list); } - DataLog dataLog1 = new DataLog(api, "开始时间:"+beginDateStr+";结束时间:"+endDateStr, new Date(), null, "数据更新,更新SF第" + z*50 + i + "批数据", "UpdateSF"); + DataLog dataLog1 = new DataLog(api, "开始时间:"+beginDateStr+";结束时间:"+endDateStr, new Date(), null, "数据更新,更新SF第" + (z*50 + i) + "批数据", "UpdateSF"); SaveResult[] saveResults = partnerConnection.update(accounts); for (SaveResult saveResult : saveResults) { if (!saveResult.getSuccess()) { @@ -932,6 +934,8 @@ public class DataImportNewServiceImpl implements DataImportNewService { .eq("sync_end_date", endDate) .set("sf_update_num", targetCount); dataBatchService.update(updateQw2); + + log.info("数据更新完成,API: {},总更新数: {},开始时间: {},结束时间: {}", api, targetCount, beginDateStr, endDateStr); } /** diff --git a/src/main/java/com/celnet/datadump/service/impl/DataImportServiceImpl.java b/src/main/java/com/celnet/datadump/service/impl/DataImportServiceImpl.java index 2a33678..a1de704 100644 --- a/src/main/java/com/celnet/datadump/service/impl/DataImportServiceImpl.java +++ b/src/main/java/com/celnet/datadump/service/impl/DataImportServiceImpl.java @@ -316,7 +316,7 @@ public class DataImportServiceImpl implements DataImportService { if (startIndex >= endIndex) { break; } - DataLog dataLog = new DataLog(api, "开始时间:"+beginDateStr+";结束时间:"+endDateStr, new Date(), null, "获取new_id,查询并组装第" + z*50 + i + "批数据", "QueryDB"); + DataLog dataLog = new DataLog(api, "开始时间:"+beginDateStr+";结束时间:"+endDateStr, new Date(), null, "获取new_id,查询并组装第" + (z*50 + i) + "批数据", "QueryDB"); List> mapList = mapsList.subList(startIndex, endIndex); int sized = mapList.size(); @@ -392,7 +392,7 @@ public class DataImportServiceImpl implements DataImportService { continue; } } - if ("picklist".equals(dataField.getSfType())){ + if ("picklist".equals(dataField.getSfType()) || "multipicklist".equals(dataField.getSfType())){ List> pickList = customMapper.list("value", "data_picklist", "api = '"+api+"' and field = '"+dataField.getField()+"' limit 1"); account.setField(dataField.getField(), pickList.get(0).get("value")); continue; @@ -420,7 +420,7 @@ public class DataImportServiceImpl implements DataImportService { dataLog.setEndTime(new Date()); dataLogService.save(dataLog); - DataLog dataLog1 = new DataLog(api, "开始时间:"+beginDateStr+";结束时间:"+endDateStr, new Date(), null, "获取new_id,插入SF第" + z*50 + i + "批数据", "InserSF"); + DataLog dataLog1 = new DataLog(api, "开始时间:"+beginDateStr+";结束时间:"+endDateStr, new Date(), null, "获取new_id,插入SF第" + (z*50 + i) + "批数据", "InserSF"); if (infoFlag != null && "1".equals(infoFlag.get("value"))){ printAccountsDetails(accounts, list); } @@ -429,7 +429,7 @@ public class DataImportServiceImpl implements DataImportService { dataLog1.setEndTime(new Date()); dataLogService.save(dataLog1); - DataLog dataLog2 = new DataLog(api, "开始时间:"+beginDateStr+";结束时间:"+endDateStr, new Date(), null, "获取new_id,更新DB第" + z*50 + i + "批数据", "UpdateDB"); + DataLog dataLog2 = new DataLog(api, "开始时间:"+beginDateStr+";结束时间:"+endDateStr, new Date(), null, "获取new_id,更新DB第" + (z*50 + i) + "批数据", "UpdateDB"); int index = 0; for (SaveResult saveResult : saveResults){ diff --git a/src/main/java/com/celnet/datadump/service/impl/FileServiceImpl.java b/src/main/java/com/celnet/datadump/service/impl/FileServiceImpl.java index a3deac7..b12ab58 100644 --- a/src/main/java/com/celnet/datadump/service/impl/FileServiceImpl.java +++ b/src/main/java/com/celnet/datadump/service/impl/FileServiceImpl.java @@ -475,7 +475,7 @@ public class FileServiceImpl implements FileService { boolean index = false; String documentId = (String) documentMap.get("Id"); // 获取未存储的附件id - List> list = customMapper.list("Id, url, PathOnClient, Title, ContentDocumentId, VersionNumber", api, "ContentDocumentId = '" + documentId + "' and new_id is null ORDER BY VersionNumber ASC"); + List> list = customMapper.list("Id, url, PathOnClient, Title, ContentDocumentId, VersionNumber,Description", api, "ContentDocumentId = '" + documentId + "' and new_id is null ORDER BY VersionNumber ASC"); if (CollectionUtils.isEmpty(list)) { continue; } @@ -497,6 +497,7 @@ public class FileServiceImpl implements FileService { String fileName = (String) map.get("PathOnClient"); String title = (String) map.get("Title"); String oldDocumentId = (String) map.get("ContentDocumentId"); + String description = (String) map.get("Description"); int versionNumber = Integer.parseInt(map.get("VersionNumber").toString()); log.info("文件名称:" + fileName); @@ -521,6 +522,7 @@ public class FileServiceImpl implements FileService { if (newDocumentId != null) { credentialsJsonParam.put("ContentDocumentId", newDocumentId); } + credentialsJsonParam.put("description", description); MultipartEntityBuilder builder = MultipartEntityBuilder.create(); builder.addTextBody("entity_content", credentialsJsonParam.toJSONString(), ContentType.APPLICATION_JSON); builder.addBinaryBody("VersionData", new FileInputStream(file), ContentType.APPLICATION_OCTET_STREAM, fileName); diff --git a/src/main/resources/mapper/CustomMapper.xml b/src/main/resources/mapper/CustomMapper.xml index cadc0f2..626f1df 100644 --- a/src/main/resources/mapper/CustomMapper.xml +++ b/src/main/resources/mapper/CustomMapper.xml @@ -21,6 +21,9 @@ `${map.name}` ${map.type} comment '${map.comment}', + + `id` INT NOT NULL AUTO_INCREMENT, + diff --git a/src/main/resources/mapper/SalesforceMapper.xml b/src/main/resources/mapper/SalesforceMapper.xml index 4c5e66d..687e4be 100644 --- a/src/main/resources/mapper/SalesforceMapper.xml +++ b/src/main/resources/mapper/SalesforceMapper.xml @@ -202,6 +202,9 @@ order by ${param.idField} asc + + limit #{param.limit} +