|
|
|
|
@@ -0,0 +1,950 @@
|
|
|
|
|
package tech.easyflow.datacenter.excel.service.impl;
|
|
|
|
|
|
|
|
|
|
import com.alibaba.fastjson2.JSONObject;
|
|
|
|
|
import com.mybatisflex.core.keygen.impl.SnowFlakeIDKeyGenerator;
|
|
|
|
|
import com.mybatisflex.core.paginate.Page;
|
|
|
|
|
import com.mybatisflex.core.row.Row;
|
|
|
|
|
import org.apache.poi.ss.usermodel.Cell;
|
|
|
|
|
import org.apache.poi.ss.usermodel.DataFormatter;
|
|
|
|
|
import org.apache.poi.ss.usermodel.Sheet;
|
|
|
|
|
import org.apache.poi.ss.usermodel.Workbook;
|
|
|
|
|
import org.apache.poi.ss.usermodel.WorkbookFactory;
|
|
|
|
|
import org.apache.poi.xssf.streaming.SXSSFWorkbook;
|
|
|
|
|
import org.springframework.stereotype.Service;
|
|
|
|
|
import org.springframework.transaction.annotation.Transactional;
|
|
|
|
|
import org.springframework.util.CollectionUtils;
|
|
|
|
|
import org.springframework.web.multipart.MultipartFile;
|
|
|
|
|
import tech.easyflow.common.constant.enums.EnumFieldType;
|
|
|
|
|
import tech.easyflow.common.entity.LoginAccount;
|
|
|
|
|
import tech.easyflow.common.web.exceptions.BusinessException;
|
|
|
|
|
import tech.easyflow.datacenter.adapter.DbHandleManager;
|
|
|
|
|
import tech.easyflow.datacenter.entity.DatacenterTable;
|
|
|
|
|
import tech.easyflow.datacenter.entity.DatacenterTableField;
|
|
|
|
|
import tech.easyflow.datacenter.excel.model.DatacenterExcelDeriveRequest;
|
|
|
|
|
import tech.easyflow.datacenter.excel.model.DatacenterExcelExportRequest;
|
|
|
|
|
import tech.easyflow.datacenter.excel.model.DatacenterExcelMergeRequest;
|
|
|
|
|
import tech.easyflow.datacenter.excel.model.DatacenterExcelSplitRequest;
|
|
|
|
|
import tech.easyflow.datacenter.excel.service.DatacenterExcelImportService;
|
|
|
|
|
import tech.easyflow.datacenter.execution.model.DatacenterQueryFilter;
|
|
|
|
|
import tech.easyflow.datacenter.execution.model.DatacenterQueryRequest;
|
|
|
|
|
import tech.easyflow.datacenter.execution.model.DatasetRef;
|
|
|
|
|
import tech.easyflow.datacenter.execution.service.DatacenterDatasetQueryService;
|
|
|
|
|
import tech.easyflow.datacenter.mapper.DatacenterDatasetVersionMapper;
|
|
|
|
|
import tech.easyflow.datacenter.mapper.DatacenterDerivedTableMapper;
|
|
|
|
|
import tech.easyflow.datacenter.mapper.DatacenterImportJobMapper;
|
|
|
|
|
import tech.easyflow.datacenter.meta.entity.DatacenterCatalog;
|
|
|
|
|
import tech.easyflow.datacenter.meta.entity.DatacenterDatasetVersion;
|
|
|
|
|
import tech.easyflow.datacenter.meta.entity.DatacenterDerivedTable;
|
|
|
|
|
import tech.easyflow.datacenter.meta.entity.DatacenterImportJob;
|
|
|
|
|
import tech.easyflow.datacenter.meta.entity.DatacenterSource;
|
|
|
|
|
import tech.easyflow.datacenter.meta.enums.DatacenterImportStatus;
|
|
|
|
|
import tech.easyflow.datacenter.meta.enums.DatacenterSourceType;
|
|
|
|
|
import tech.easyflow.datacenter.meta.enums.DatacenterTableKind;
|
|
|
|
|
import tech.easyflow.datacenter.meta.model.DatacenterTableDetailMeta;
|
|
|
|
|
import tech.easyflow.datacenter.meta.service.DatacenterDatasetRegistryService;
|
|
|
|
|
import tech.easyflow.datacenter.meta.service.DatacenterSourceService;
|
|
|
|
|
|
|
|
|
|
import javax.annotation.Resource;
|
|
|
|
|
import java.io.File;
|
|
|
|
|
import java.io.FileOutputStream;
|
|
|
|
|
import java.io.InputStream;
|
|
|
|
|
import java.math.BigInteger;
|
|
|
|
|
import java.nio.file.Files;
|
|
|
|
|
import java.nio.file.Path;
|
|
|
|
|
import java.time.LocalDateTime;
|
|
|
|
|
import java.time.format.DateTimeFormatter;
|
|
|
|
|
import java.util.ArrayList;
|
|
|
|
|
import java.util.Date;
|
|
|
|
|
import java.util.HashMap;
|
|
|
|
|
import java.util.HashSet;
|
|
|
|
|
import java.util.LinkedHashMap;
|
|
|
|
|
import java.util.LinkedHashSet;
|
|
|
|
|
import java.util.List;
|
|
|
|
|
import java.util.Locale;
|
|
|
|
|
import java.util.Map;
|
|
|
|
|
import java.util.Objects;
|
|
|
|
|
import java.util.Set;
|
|
|
|
|
import java.util.UUID;
|
|
|
|
|
|
|
|
|
|
@Service
|
|
|
|
|
public class DatacenterExcelImportServiceImpl implements DatacenterExcelImportService {
|
|
|
|
|
|
|
|
|
|
private static final long QUERY_BATCH_SIZE = 500L;
|
|
|
|
|
private static final DateTimeFormatter EXPORT_TIME_FORMAT = DateTimeFormatter.ofPattern("yyyyMMddHHmmss");
|
|
|
|
|
|
|
|
|
|
@Resource
|
|
|
|
|
private DatacenterSourceService sourceService;
|
|
|
|
|
@Resource
|
|
|
|
|
private DatacenterDatasetRegistryService registryService;
|
|
|
|
|
@Resource
|
|
|
|
|
private DatacenterImportJobMapper importJobMapper;
|
|
|
|
|
@Resource
|
|
|
|
|
private DatacenterDatasetVersionMapper datasetVersionMapper;
|
|
|
|
|
@Resource
|
|
|
|
|
private DatacenterDerivedTableMapper derivedTableMapper;
|
|
|
|
|
@Resource
|
|
|
|
|
private DatacenterDatasetQueryService queryService;
|
|
|
|
|
@Resource
|
|
|
|
|
private DbHandleManager dbHandleManager;
|
|
|
|
|
|
|
|
|
|
@Override
|
|
|
|
|
@Transactional(rollbackFor = Exception.class)
|
|
|
|
|
public DatacenterImportJob importWorkbook(MultipartFile file, LoginAccount account) throws Exception {
|
|
|
|
|
if (file == null || file.isEmpty()) {
|
|
|
|
|
throw new BusinessException("Excel 文件不能为空");
|
|
|
|
|
}
|
|
|
|
|
String workbookName = extractWorkbookName(file.getOriginalFilename());
|
|
|
|
|
DatacenterSource source = new DatacenterSource();
|
|
|
|
|
source.setSourceName(workbookName);
|
|
|
|
|
source.setSourceCode("EXCEL_" + UUID.randomUUID());
|
|
|
|
|
source.setSourceType(DatacenterSourceType.EXCEL.name());
|
|
|
|
|
source.setAccessMode("READ_WRITE");
|
|
|
|
|
source.setBuiltinFlag(0);
|
|
|
|
|
source.setConfigJson(Map.of("originFileName", file.getOriginalFilename()));
|
|
|
|
|
source = sourceService.saveSource(source, account);
|
|
|
|
|
DatacenterCatalog catalog = registryService.ensureCatalog(source, workbookName, account);
|
|
|
|
|
|
|
|
|
|
DatacenterImportJob job = createJob("EXCEL_IMPORT", source.getId(), catalog.getId(), null,
|
|
|
|
|
file.getOriginalFilename(), Map.of("operation", "import"), account);
|
|
|
|
|
|
|
|
|
|
long totalRows = 0L;
|
|
|
|
|
long successRows = 0L;
|
|
|
|
|
List<BigInteger> createdTableIds = new ArrayList<>();
|
|
|
|
|
try (InputStream inputStream = file.getInputStream(); Workbook workbook = WorkbookFactory.create(inputStream)) {
|
|
|
|
|
DataFormatter formatter = new DataFormatter();
|
|
|
|
|
for (int sheetIndex = 0; sheetIndex < workbook.getNumberOfSheets(); sheetIndex++) {
|
|
|
|
|
Sheet sheet = workbook.getSheetAt(sheetIndex);
|
|
|
|
|
org.apache.poi.ss.usermodel.Row headerRow = sheet.getRow(sheet.getFirstRowNum());
|
|
|
|
|
if (headerRow == null) {
|
|
|
|
|
continue;
|
|
|
|
|
}
|
|
|
|
|
List<DatacenterTableField> fields = buildFields(headerRow, formatter);
|
|
|
|
|
if (fields.isEmpty()) {
|
|
|
|
|
continue;
|
|
|
|
|
}
|
|
|
|
|
DatacenterTable table = new DatacenterTable();
|
|
|
|
|
table.setTableName(uniqueTableName(source.getId(), catalog.getId(), sheet.getSheetName()));
|
|
|
|
|
table.setTableDesc(sheet.getSheetName());
|
|
|
|
|
table.setActualTable(buildMaterializedTableName(source.getId(), sheetIndex));
|
|
|
|
|
table.setMaterializedTable(table.getActualTable());
|
|
|
|
|
table.setTableKind(DatacenterTableKind.EXCEL_MATERIALIZED.name());
|
|
|
|
|
table.setAccessMode("READ_WRITE");
|
|
|
|
|
table.setVersioningEnabled(1);
|
|
|
|
|
table.setCapabilitiesJson(defaultExcelCapabilities());
|
|
|
|
|
table.setFields(fields);
|
|
|
|
|
|
|
|
|
|
DatacenterTableDetailMeta detail = new DatacenterTableDetailMeta();
|
|
|
|
|
detail.setTable(table);
|
|
|
|
|
detail.setFields(fields);
|
|
|
|
|
|
|
|
|
|
dbHandleManager.getDbHandler().createTable(table);
|
|
|
|
|
DatacenterTable savedTable = registryService.registerTable(source, catalog, detail, account);
|
|
|
|
|
savedTable.setFields(registryService.getFields(savedTable.getId()));
|
|
|
|
|
createdTableIds.add(savedTable.getId());
|
|
|
|
|
|
|
|
|
|
for (int rowIndex = sheet.getFirstRowNum() + 1; rowIndex <= sheet.getLastRowNum(); rowIndex++) {
|
|
|
|
|
org.apache.poi.ss.usermodel.Row row = sheet.getRow(rowIndex);
|
|
|
|
|
if (row == null) {
|
|
|
|
|
continue;
|
|
|
|
|
}
|
|
|
|
|
JSONObject payload = new JSONObject();
|
|
|
|
|
boolean hasValue = false;
|
|
|
|
|
for (int cellIndex = 0; cellIndex < savedTable.getFields().size(); cellIndex++) {
|
|
|
|
|
String value = formatter.formatCellValue(row.getCell(cellIndex));
|
|
|
|
|
if (value != null && !value.isBlank()) {
|
|
|
|
|
hasValue = true;
|
|
|
|
|
}
|
|
|
|
|
payload.put(savedTable.getFields().get(cellIndex).getFieldName(), value);
|
|
|
|
|
}
|
|
|
|
|
if (!hasValue) {
|
|
|
|
|
continue;
|
|
|
|
|
}
|
|
|
|
|
dbHandleManager.getDbHandler().saveValue(savedTable, payload, account);
|
|
|
|
|
totalRows++;
|
|
|
|
|
successRows++;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
createVersion(savedTable, "initial-import", Map.of(
|
|
|
|
|
"sheetName", sheet.getSheetName(),
|
|
|
|
|
"sourceId", source.getId(),
|
|
|
|
|
"originFileName", file.getOriginalFilename()
|
|
|
|
|
), account);
|
|
|
|
|
}
|
|
|
|
|
job.setTableId(createdTableIds.isEmpty() ? null : createdTableIds.get(0));
|
|
|
|
|
finishJobSuccess(job, totalRows, successRows, Map.of("tableIds", createdTableIds));
|
|
|
|
|
return job;
|
|
|
|
|
} catch (Exception ex) {
|
|
|
|
|
finishJobFailure(job, ex);
|
|
|
|
|
throw ex;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Override
|
|
|
|
|
@Transactional(rollbackFor = Exception.class)
|
|
|
|
|
public DatacenterImportJob splitWorkbook(DatacenterExcelSplitRequest request, LoginAccount account) {
|
|
|
|
|
DatasetRef datasetRef = request == null ? null : request.getDatasetRef();
|
|
|
|
|
DatacenterImportJob job = createJob("EXCEL_SPLIT",
|
|
|
|
|
request == null ? null : request.getSourceId(),
|
|
|
|
|
request == null ? null : request.getCatalogId(),
|
|
|
|
|
datasetRef == null ? null : datasetRef.getTableId(),
|
|
|
|
|
null,
|
|
|
|
|
buildPayload("request", request),
|
|
|
|
|
account);
|
|
|
|
|
try {
|
|
|
|
|
String splitMode = normalizeMode(request == null ? null : request.getSplitMode(), "BY_ROW_COUNT");
|
|
|
|
|
return switch (splitMode) {
|
|
|
|
|
case "BY_SHEET" -> doSplitBySheet(request, account, job);
|
|
|
|
|
case "BY_FIELD_VALUE" -> doSplitByFieldValue(request, account, job);
|
|
|
|
|
default -> doSplitByRowCount(request, account, job);
|
|
|
|
|
};
|
|
|
|
|
} catch (Exception ex) {
|
|
|
|
|
finishJobFailure(job, ex);
|
|
|
|
|
throw ex;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Override
|
|
|
|
|
@Transactional(rollbackFor = Exception.class)
|
|
|
|
|
public DatacenterImportJob mergeWorkbook(DatacenterExcelMergeRequest request, LoginAccount account) {
|
|
|
|
|
if (request == null || CollectionUtils.isEmpty(request.getDatasetRefs())) {
|
|
|
|
|
throw new BusinessException("合并数据集不能为空");
|
|
|
|
|
}
|
|
|
|
|
DatacenterImportJob job = createJob("EXCEL_MERGE", null, null, null, null, buildPayload("request", request), account);
|
|
|
|
|
try {
|
|
|
|
|
String mergeMode = normalizeMode(request.getMergeMode(), "VERTICAL");
|
|
|
|
|
return switch (mergeMode) {
|
|
|
|
|
case "HORIZONTAL" -> doHorizontalMerge(request, account, job);
|
|
|
|
|
default -> doVerticalMerge(request, account, job);
|
|
|
|
|
};
|
|
|
|
|
} catch (Exception ex) {
|
|
|
|
|
finishJobFailure(job, ex);
|
|
|
|
|
throw ex;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Override
|
|
|
|
|
@Transactional(rollbackFor = Exception.class)
|
|
|
|
|
public DatacenterImportJob deriveWorkbook(DatacenterExcelDeriveRequest request, LoginAccount account) {
|
|
|
|
|
if (request == null || request.getDatasetRef() == null) {
|
|
|
|
|
throw new BusinessException("派生数据集不能为空");
|
|
|
|
|
}
|
|
|
|
|
DatacenterImportJob job = createJob("EXCEL_DERIVE",
|
|
|
|
|
request.getDatasetRef().getSourceId(),
|
|
|
|
|
request.getDatasetRef().getCatalogId(),
|
|
|
|
|
request.getDatasetRef().getTableId(),
|
|
|
|
|
null,
|
|
|
|
|
buildPayload("request", request),
|
|
|
|
|
account);
|
|
|
|
|
try {
|
|
|
|
|
DatacenterTable sourceTable = resolveTable(request.getDatasetRef());
|
|
|
|
|
DatacenterSource source = registryService.getSourceRequired(sourceTable.getSourceId());
|
|
|
|
|
DatacenterCatalog catalog = requireCatalog(sourceTable.getCatalogId());
|
|
|
|
|
List<String> selectedColumns = resolveSelectedColumns(sourceTable, request.getSelectedColumns());
|
|
|
|
|
List<DatacenterTableField> targetFields = buildDerivedFields(sourceTable, request);
|
|
|
|
|
DatacenterTable targetTable = createDerivedTable(source, catalog, targetFields,
|
|
|
|
|
request.getTargetTableName(), "DERIVE", Map.of("sourceTableId", sourceTable.getId()), account);
|
|
|
|
|
|
|
|
|
|
DatacenterQueryRequest queryRequest = new DatacenterQueryRequest();
|
|
|
|
|
queryRequest.setDatasetRef(registryService.resolveDatasetRef(sourceTable.getId()));
|
|
|
|
|
queryRequest.setSelectedColumns(selectedColumns);
|
|
|
|
|
queryRequest.setFilters(request.getFilters());
|
|
|
|
|
|
|
|
|
|
long successRows = copyRows(queryRequest, rows -> mapDerivedRow(rows, selectedColumns, request), targetTable, account);
|
|
|
|
|
createLineage(sourceTable.getId(), targetTable.getId(), "DERIVE", Map.of("request", request), account);
|
|
|
|
|
finishJobSuccess(job, successRows, successRows, Map.of("derivedTableId", targetTable.getId()));
|
|
|
|
|
job.setTableId(targetTable.getId());
|
|
|
|
|
return job;
|
|
|
|
|
} catch (Exception ex) {
|
|
|
|
|
finishJobFailure(job, ex);
|
|
|
|
|
throw ex;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Override
|
|
|
|
|
@Transactional(rollbackFor = Exception.class)
|
|
|
|
|
public DatacenterImportJob exportWorkbook(DatacenterExcelExportRequest request, LoginAccount account) throws Exception {
|
|
|
|
|
List<DatacenterTable> tables = resolveExportTables(request);
|
|
|
|
|
if (tables.isEmpty()) {
|
|
|
|
|
throw new BusinessException("没有可导出的 Excel 数据集");
|
|
|
|
|
}
|
|
|
|
|
String fileName = buildExportFileName(request == null ? null : request.getFileName());
|
|
|
|
|
DatacenterImportJob job = createJob("EXCEL_EXPORT",
|
|
|
|
|
request == null ? null : request.getSourceId(),
|
|
|
|
|
request == null ? null : request.getCatalogId(),
|
|
|
|
|
null,
|
|
|
|
|
fileName,
|
|
|
|
|
buildPayload("request", request),
|
|
|
|
|
account);
|
|
|
|
|
Path exportDir = ensureExportDir();
|
|
|
|
|
Path exportFile = exportDir.resolve(fileName);
|
|
|
|
|
long totalRows = 0L;
|
|
|
|
|
try (SXSSFWorkbook workbook = new SXSSFWorkbook(200); FileOutputStream outputStream = new FileOutputStream(exportFile.toFile())) {
|
|
|
|
|
workbook.setCompressTempFiles(true);
|
|
|
|
|
Set<String> usedSheetNames = new HashSet<>();
|
|
|
|
|
for (DatacenterTable table : tables) {
|
|
|
|
|
String sheetName = uniqueSheetName(table.getTableName(), usedSheetNames);
|
|
|
|
|
org.apache.poi.ss.usermodel.Sheet sheet = workbook.createSheet(sheetName);
|
|
|
|
|
writeHeaderRow(sheet, table.getFields());
|
|
|
|
|
DatacenterQueryRequest queryRequest = new DatacenterQueryRequest();
|
|
|
|
|
queryRequest.setDatasetRef(registryService.resolveDatasetRef(table.getId()));
|
|
|
|
|
queryRequest.setSelectedColumns(table.getFields().stream().map(DatacenterTableField::getFieldName).toList());
|
|
|
|
|
final int[] rowIndex = {1};
|
|
|
|
|
totalRows += iterateRows(queryRequest, row -> {
|
|
|
|
|
org.apache.poi.ss.usermodel.Row excelRow = sheet.createRow(rowIndex[0]++);
|
|
|
|
|
for (int i = 0; i < table.getFields().size(); i++) {
|
|
|
|
|
Cell cell = excelRow.createCell(i);
|
|
|
|
|
Object value = row.get(table.getFields().get(i).getFieldName());
|
|
|
|
|
cell.setCellValue(value == null ? "" : String.valueOf(value));
|
|
|
|
|
}
|
|
|
|
|
});
|
|
|
|
|
}
|
|
|
|
|
workbook.write(outputStream);
|
|
|
|
|
workbook.dispose();
|
|
|
|
|
job.setStoragePath(exportFile.toAbsolutePath().toString());
|
|
|
|
|
finishJobSuccess(job, totalRows, totalRows, Map.of("storagePath", job.getStoragePath(), "fileName", fileName));
|
|
|
|
|
return job;
|
|
|
|
|
} catch (Exception ex) {
|
|
|
|
|
finishJobFailure(job, ex);
|
|
|
|
|
throw ex;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Override
|
|
|
|
|
public DatacenterImportJob getImportJobDetail(BigInteger jobId) {
|
|
|
|
|
DatacenterImportJob job = importJobMapper.selectOneById(jobId);
|
|
|
|
|
if (job == null) {
|
|
|
|
|
throw new BusinessException("导入任务不存在: " + jobId);
|
|
|
|
|
}
|
|
|
|
|
return job;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Override
|
|
|
|
|
public List<DatacenterImportJob> listJobs(BigInteger sourceId, BigInteger tableId) {
|
|
|
|
|
var wrapper = com.mybatisflex.core.query.QueryWrapper.create();
|
|
|
|
|
if (sourceId != null) {
|
|
|
|
|
wrapper.eq(DatacenterImportJob::getSourceId, sourceId);
|
|
|
|
|
}
|
|
|
|
|
if (tableId != null) {
|
|
|
|
|
wrapper.eq(DatacenterImportJob::getTableId, tableId);
|
|
|
|
|
}
|
|
|
|
|
wrapper.orderBy("created desc");
|
|
|
|
|
wrapper.limit(20L);
|
|
|
|
|
return importJobMapper.selectListByQuery(wrapper);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private DatacenterImportJob doSplitBySheet(DatacenterExcelSplitRequest request, LoginAccount account, DatacenterImportJob job) {
|
|
|
|
|
BigInteger sourceId = request.getSourceId();
|
|
|
|
|
BigInteger catalogId = request.getCatalogId();
|
|
|
|
|
if (sourceId == null && request.getDatasetRef() != null) {
|
|
|
|
|
sourceId = request.getDatasetRef().getSourceId();
|
|
|
|
|
catalogId = request.getDatasetRef().getCatalogId();
|
|
|
|
|
}
|
|
|
|
|
if (sourceId == null) {
|
|
|
|
|
throw new BusinessException("按 sheet 拆分需要 sourceId");
|
|
|
|
|
}
|
|
|
|
|
DatacenterSource source = registryService.getSourceRequired(sourceId);
|
|
|
|
|
DatacenterCatalog catalog = requireCatalog(catalogId);
|
|
|
|
|
List<DatacenterTable> sourceTables = registryService.listManagedTables(sourceId, catalogId);
|
|
|
|
|
if (sourceTables.isEmpty()) {
|
|
|
|
|
throw new BusinessException("当前 workbook 下没有可拆分的 sheet 表");
|
|
|
|
|
}
|
|
|
|
|
List<BigInteger> derivedIds = new ArrayList<>();
|
|
|
|
|
long successRows = 0L;
|
|
|
|
|
for (DatacenterTable sourceTable : sourceTables) {
|
|
|
|
|
DatacenterTable fullTable = registryService.getTableWithFields(sourceTable.getId());
|
|
|
|
|
DatacenterTable targetTable = createDerivedTable(source, catalog, cloneFields(fullTable.getFields()),
|
|
|
|
|
resolveSplitPrefix(request, fullTable.getTableName()) + "_copy", "SPLIT_BY_SHEET",
|
|
|
|
|
Map.of("sourceTableId", fullTable.getId()), account);
|
|
|
|
|
DatacenterQueryRequest queryRequest = new DatacenterQueryRequest();
|
|
|
|
|
queryRequest.setDatasetRef(registryService.resolveDatasetRef(fullTable.getId()));
|
|
|
|
|
queryRequest.setSelectedColumns(fullTable.getFields().stream().map(DatacenterTableField::getFieldName).toList());
|
|
|
|
|
successRows += copyRows(queryRequest, this::mapRow, targetTable, account);
|
|
|
|
|
createLineage(fullTable.getId(), targetTable.getId(), "SPLIT_BY_SHEET", Map.of("sourceTableId", fullTable.getId()), account);
|
|
|
|
|
derivedIds.add(targetTable.getId());
|
|
|
|
|
}
|
|
|
|
|
finishJobSuccess(job, successRows, successRows, Map.of("derivedTableIds", derivedIds));
|
|
|
|
|
return job;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private DatacenterImportJob doSplitByRowCount(DatacenterExcelSplitRequest request, LoginAccount account, DatacenterImportJob job) {
|
|
|
|
|
if (request == null || request.getDatasetRef() == null) {
|
|
|
|
|
throw new BusinessException("按行数拆分需要数据集");
|
|
|
|
|
}
|
|
|
|
|
int rowBatchSize = request.getRowBatchSize() == null || request.getRowBatchSize() < 1 ? 1000 : request.getRowBatchSize();
|
|
|
|
|
DatacenterTable sourceTable = resolveTable(request.getDatasetRef());
|
|
|
|
|
DatacenterSource source = registryService.getSourceRequired(sourceTable.getSourceId());
|
|
|
|
|
DatacenterCatalog catalog = requireCatalog(sourceTable.getCatalogId());
|
|
|
|
|
String baseName = resolveSplitPrefix(request, sourceTable.getTableName());
|
|
|
|
|
List<BigInteger> derivedIds = new ArrayList<>();
|
|
|
|
|
final Holder holder = new Holder();
|
|
|
|
|
long totalRows = iterateRows(buildFullQuery(sourceTable), row -> {
|
|
|
|
|
if (holder.targetTable == null || holder.currentSize >= rowBatchSize) {
|
|
|
|
|
holder.batchNo++;
|
|
|
|
|
holder.targetTable = createDerivedTable(source, catalog, cloneFields(sourceTable.getFields()),
|
|
|
|
|
baseName + "_part_" + holder.batchNo, "SPLIT_BY_ROW_COUNT",
|
|
|
|
|
Map.of("sourceTableId", sourceTable.getId(), "batchNo", holder.batchNo, "rowBatchSize", rowBatchSize), account);
|
|
|
|
|
createLineage(sourceTable.getId(), holder.targetTable.getId(), "SPLIT_BY_ROW_COUNT",
|
|
|
|
|
Map.of("sourceTableId", sourceTable.getId(), "batchNo", holder.batchNo), account);
|
|
|
|
|
derivedIds.add(holder.targetTable.getId());
|
|
|
|
|
holder.currentSize = 0;
|
|
|
|
|
}
|
|
|
|
|
saveToTable(holder.targetTable, mapRow(row), account);
|
|
|
|
|
holder.currentSize++;
|
|
|
|
|
});
|
|
|
|
|
finishJobSuccess(job, totalRows, totalRows, Map.of("derivedTableIds", derivedIds));
|
|
|
|
|
return job;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private DatacenterImportJob doSplitByFieldValue(DatacenterExcelSplitRequest request, LoginAccount account, DatacenterImportJob job) {
|
|
|
|
|
if (request == null || request.getDatasetRef() == null || request.getFieldName() == null || request.getFieldName().isBlank()) {
|
|
|
|
|
throw new BusinessException("按字段值拆分需要数据集和字段名");
|
|
|
|
|
}
|
|
|
|
|
DatacenterTable sourceTable = resolveTable(request.getDatasetRef());
|
|
|
|
|
DatacenterSource source = registryService.getSourceRequired(sourceTable.getSourceId());
|
|
|
|
|
DatacenterCatalog catalog = requireCatalog(sourceTable.getCatalogId());
|
|
|
|
|
DatacenterTableField splitField = sourceTable.getFields().stream()
|
|
|
|
|
.filter(field -> request.getFieldName().equals(field.getFieldName()))
|
|
|
|
|
.findFirst()
|
|
|
|
|
.orElseThrow(() -> new BusinessException("拆分字段不存在: " + request.getFieldName()));
|
|
|
|
|
String prefix = resolveSplitPrefix(request, sourceTable.getTableName());
|
|
|
|
|
Map<String, DatacenterTable> targets = new LinkedHashMap<>();
|
|
|
|
|
List<BigInteger> derivedIds = new ArrayList<>();
|
|
|
|
|
long totalRows = iterateRows(buildFullQuery(sourceTable), row -> {
|
|
|
|
|
String fieldValue = stringify(row.get(splitField.getFieldName()));
|
|
|
|
|
String bucket = fieldValue == null || fieldValue.isBlank() ? "empty" : fieldValue;
|
|
|
|
|
DatacenterTable targetTable = targets.get(bucket);
|
|
|
|
|
if (targetTable == null) {
|
|
|
|
|
targetTable = createDerivedTable(source, catalog, cloneFields(sourceTable.getFields()),
|
|
|
|
|
prefix + "_" + normalizeIdentifier(bucket), "SPLIT_BY_FIELD_VALUE",
|
|
|
|
|
Map.of("sourceTableId", sourceTable.getId(), "fieldName", splitField.getFieldName(), "fieldValue", bucket), account);
|
|
|
|
|
createLineage(sourceTable.getId(), targetTable.getId(), "SPLIT_BY_FIELD_VALUE",
|
|
|
|
|
Map.of("sourceTableId", sourceTable.getId(), "fieldName", splitField.getFieldName(), "fieldValue", bucket), account);
|
|
|
|
|
targets.put(bucket, targetTable);
|
|
|
|
|
derivedIds.add(targetTable.getId());
|
|
|
|
|
}
|
|
|
|
|
saveToTable(targetTable, mapRow(row), account);
|
|
|
|
|
});
|
|
|
|
|
finishJobSuccess(job, totalRows, totalRows, Map.of("derivedTableIds", derivedIds));
|
|
|
|
|
return job;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private DatacenterImportJob doVerticalMerge(DatacenterExcelMergeRequest request, LoginAccount account, DatacenterImportJob job) {
|
|
|
|
|
List<DatacenterTable> tables = request.getDatasetRefs().stream().map(this::resolveTable).toList();
|
|
|
|
|
DatacenterTable firstTable = tables.get(0);
|
|
|
|
|
DatacenterSource source = registryService.getSourceRequired(firstTable.getSourceId());
|
|
|
|
|
DatacenterCatalog catalog = requireCatalog(firstTable.getCatalogId());
|
|
|
|
|
assertSameCatalog(tables);
|
|
|
|
|
assertSameFields(tables);
|
|
|
|
|
DatacenterTable targetTable = createDerivedTable(source, catalog, cloneFields(firstTable.getFields()),
|
|
|
|
|
request.getTargetTableName(), "MERGE_VERTICAL", Map.of("sourceTableIds", tables.stream().map(DatacenterTable::getId).toList()), account);
|
|
|
|
|
long successRows = 0L;
|
|
|
|
|
for (DatacenterTable table : tables) {
|
|
|
|
|
successRows += copyRows(buildFullQuery(table), this::mapRow, targetTable, account);
|
|
|
|
|
createLineage(table.getId(), targetTable.getId(), "MERGE_VERTICAL", Map.of("sourceTableId", table.getId()), account);
|
|
|
|
|
}
|
|
|
|
|
job.setTableId(targetTable.getId());
|
|
|
|
|
finishJobSuccess(job, successRows, successRows, Map.of("derivedTableId", targetTable.getId()));
|
|
|
|
|
return job;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private DatacenterImportJob doHorizontalMerge(DatacenterExcelMergeRequest request, LoginAccount account, DatacenterImportJob job) {
|
|
|
|
|
if (request.getJoinKey() == null || request.getJoinKey().isBlank()) {
|
|
|
|
|
throw new BusinessException("横向合并必须指定 joinKey");
|
|
|
|
|
}
|
|
|
|
|
List<DatacenterTable> tables = request.getDatasetRefs().stream().map(this::resolveTable).toList();
|
|
|
|
|
DatacenterTable firstTable = tables.get(0);
|
|
|
|
|
DatacenterSource source = registryService.getSourceRequired(firstTable.getSourceId());
|
|
|
|
|
DatacenterCatalog catalog = requireCatalog(firstTable.getCatalogId());
|
|
|
|
|
assertSameCatalog(tables);
|
|
|
|
|
|
|
|
|
|
List<DatacenterTableField> mergedFields = new ArrayList<>();
|
|
|
|
|
Set<String> usedFieldNames = new LinkedHashSet<>();
|
|
|
|
|
Map<BigInteger, Map<String, String>> fieldMappings = new LinkedHashMap<>();
|
|
|
|
|
for (DatacenterTable table : tables) {
|
|
|
|
|
Map<String, String> mapping = new LinkedHashMap<>();
|
|
|
|
|
for (DatacenterTableField field : table.getFields()) {
|
|
|
|
|
String targetFieldName;
|
|
|
|
|
if (field.getFieldName().equals(request.getJoinKey())) {
|
|
|
|
|
targetFieldName = field.getFieldName();
|
|
|
|
|
} else {
|
|
|
|
|
targetFieldName = field.getFieldName();
|
|
|
|
|
if (usedFieldNames.contains(targetFieldName)) {
|
|
|
|
|
targetFieldName = normalizeIdentifier(table.getTableName()) + "_" + targetFieldName;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
if (!usedFieldNames.contains(targetFieldName)) {
|
|
|
|
|
usedFieldNames.add(targetFieldName);
|
|
|
|
|
mergedFields.add(cloneField(field, targetFieldName, field.getFieldDesc()));
|
|
|
|
|
}
|
|
|
|
|
mapping.put(field.getFieldName(), targetFieldName);
|
|
|
|
|
}
|
|
|
|
|
fieldMappings.put(table.getId(), mapping);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
DatacenterTable targetTable = createDerivedTable(source, catalog, mergedFields,
|
|
|
|
|
request.getTargetTableName(), "MERGE_HORIZONTAL", Map.of("sourceTableIds", tables.stream().map(DatacenterTable::getId).toList(), "joinKey", request.getJoinKey()), account);
|
|
|
|
|
Map<String, JSONObject> mergedRows = new LinkedHashMap<>();
|
|
|
|
|
for (DatacenterTable table : tables) {
|
|
|
|
|
Map<String, String> mapping = fieldMappings.get(table.getId());
|
|
|
|
|
iterateRows(buildFullQuery(table), row -> {
|
|
|
|
|
String joinValue = stringify(row.get(request.getJoinKey()));
|
|
|
|
|
if (joinValue == null || joinValue.isBlank()) {
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
JSONObject target = mergedRows.computeIfAbsent(joinValue, key -> new JSONObject());
|
|
|
|
|
mapping.forEach((sourceField, targetField) -> target.put(targetField, row.get(sourceField)));
|
|
|
|
|
});
|
|
|
|
|
createLineage(table.getId(), targetTable.getId(), "MERGE_HORIZONTAL", Map.of("sourceTableId", table.getId(), "joinKey", request.getJoinKey()), account);
|
|
|
|
|
}
|
|
|
|
|
for (JSONObject row : mergedRows.values()) {
|
|
|
|
|
saveToTable(targetTable, row, account);
|
|
|
|
|
}
|
|
|
|
|
job.setTableId(targetTable.getId());
|
|
|
|
|
finishJobSuccess(job, (long) mergedRows.size(), (long) mergedRows.size(), Map.of("derivedTableId", targetTable.getId()));
|
|
|
|
|
return job;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private DatacenterImportJob createJob(String jobType, BigInteger sourceId, BigInteger catalogId, BigInteger tableId,
|
|
|
|
|
String fileName, Map<String, Object> payload, LoginAccount account) {
|
|
|
|
|
DatacenterImportJob job = new DatacenterImportJob();
|
|
|
|
|
job.setSourceId(sourceId);
|
|
|
|
|
job.setCatalogId(catalogId);
|
|
|
|
|
job.setTableId(tableId);
|
|
|
|
|
job.setTenantId(account == null || account.getTenantId() == null ? BigInteger.ZERO : account.getTenantId());
|
|
|
|
|
job.setDeptId(account == null || account.getDeptId() == null ? BigInteger.ZERO : account.getDeptId());
|
|
|
|
|
job.setJobType(jobType);
|
|
|
|
|
job.setFileName(fileName);
|
|
|
|
|
job.setStatus(DatacenterImportStatus.RUNNING.name());
|
|
|
|
|
job.setPayloadJson(payload == null ? new LinkedHashMap<>() : new LinkedHashMap<>(payload));
|
|
|
|
|
job.setStartedAt(new Date());
|
|
|
|
|
job.setCreated(new Date());
|
|
|
|
|
job.setModified(new Date());
|
|
|
|
|
job.setCreatedBy(account == null ? BigInteger.ZERO : account.getId());
|
|
|
|
|
job.setModifiedBy(account == null ? BigInteger.ZERO : account.getId());
|
|
|
|
|
importJobMapper.insert(job);
|
|
|
|
|
return job;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private Map<String, Object> buildPayload(String key, Object value) {
|
|
|
|
|
Map<String, Object> payload = new LinkedHashMap<>();
|
|
|
|
|
payload.put(key, value);
|
|
|
|
|
return payload;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private void finishJobSuccess(DatacenterImportJob job, long totalRows, long successRows, Map<String, Object> payload) {
|
|
|
|
|
job.setStatus(DatacenterImportStatus.SUCCESS.name());
|
|
|
|
|
job.setTotalRows(totalRows);
|
|
|
|
|
job.setSuccessRows(successRows);
|
|
|
|
|
job.setErrorRows(Math.max(0L, totalRows - successRows));
|
|
|
|
|
if (payload != null) {
|
|
|
|
|
job.setPayloadJson(new LinkedHashMap<>(payload));
|
|
|
|
|
}
|
|
|
|
|
job.setFinishedAt(new Date());
|
|
|
|
|
job.setModified(new Date());
|
|
|
|
|
importJobMapper.update(job);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private void finishJobFailure(DatacenterImportJob job, Exception ex) {
|
|
|
|
|
job.setStatus(DatacenterImportStatus.FAILED.name());
|
|
|
|
|
job.setErrorSummary(ex.getMessage());
|
|
|
|
|
job.setFinishedAt(new Date());
|
|
|
|
|
job.setModified(new Date());
|
|
|
|
|
importJobMapper.update(job);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private DatacenterTable resolveTable(DatasetRef datasetRef) {
|
|
|
|
|
if (datasetRef == null || datasetRef.getTableId() == null) {
|
|
|
|
|
throw new BusinessException("缺少数据集 tableId");
|
|
|
|
|
}
|
|
|
|
|
return registryService.getTableWithFields(datasetRef.getTableId());
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private DatacenterCatalog requireCatalog(BigInteger catalogId) {
|
|
|
|
|
DatacenterCatalog catalog = registryService.getCatalogById(catalogId);
|
|
|
|
|
if (catalog == null) {
|
|
|
|
|
throw new BusinessException("目录不存在: " + catalogId);
|
|
|
|
|
}
|
|
|
|
|
return catalog;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private DatacenterTable createDerivedTable(DatacenterSource source, DatacenterCatalog catalog, List<DatacenterTableField> fields,
|
|
|
|
|
String tableName, String deriveType, Map<String, Object> config, LoginAccount account) {
|
|
|
|
|
DatacenterTable table = new DatacenterTable();
|
|
|
|
|
String resolvedName = uniqueTableName(source.getId(), catalog.getId(), normalizeLogicalName(tableName, deriveType));
|
|
|
|
|
table.setTableName(resolvedName);
|
|
|
|
|
table.setTableDesc(resolvedName);
|
|
|
|
|
table.setActualTable(buildMaterializedTableName(source.getId(), Math.abs(Objects.hash(resolvedName, deriveType))));
|
|
|
|
|
table.setMaterializedTable(table.getActualTable());
|
|
|
|
|
table.setTableKind(DatacenterTableKind.DERIVED_TABLE.name());
|
|
|
|
|
table.setAccessMode("READ_WRITE");
|
|
|
|
|
table.setVersioningEnabled(1);
|
|
|
|
|
table.setCapabilitiesJson(defaultExcelCapabilities());
|
|
|
|
|
table.setFields(fields);
|
|
|
|
|
|
|
|
|
|
DatacenterTableDetailMeta detail = new DatacenterTableDetailMeta();
|
|
|
|
|
detail.setTable(table);
|
|
|
|
|
detail.setFields(fields);
|
|
|
|
|
dbHandleManager.getDbHandler().createTable(table);
|
|
|
|
|
DatacenterTable savedTable = registryService.registerTable(source, catalog, detail, account);
|
|
|
|
|
savedTable.setFields(registryService.getFields(savedTable.getId()));
|
|
|
|
|
createVersion(savedTable, deriveType.toLowerCase(Locale.ROOT), config, account);
|
|
|
|
|
return savedTable;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private DatacenterDatasetVersion createVersion(DatacenterTable table, String versionLabel, Map<String, Object> snapshot, LoginAccount account) {
|
|
|
|
|
QueryWrapperWrapper wrapper = new QueryWrapperWrapper(table.getId());
|
|
|
|
|
DatacenterDatasetVersion version = new DatacenterDatasetVersion();
|
|
|
|
|
version.setTableId(table.getId());
|
|
|
|
|
version.setTenantId(table.getTenantId());
|
|
|
|
|
version.setDeptId(table.getDeptId());
|
|
|
|
|
version.setVersionNo(wrapper.nextVersionNo(datasetVersionMapper));
|
|
|
|
|
version.setVersionLabel(versionLabel);
|
|
|
|
|
version.setMaterializedTable(table.getMaterializedTable());
|
|
|
|
|
version.setSnapshotJson(snapshot == null ? new LinkedHashMap<>() : new LinkedHashMap<>(snapshot));
|
|
|
|
|
version.setStatus(0);
|
|
|
|
|
version.setCreated(new Date());
|
|
|
|
|
version.setModified(new Date());
|
|
|
|
|
version.setCreatedBy(account == null ? BigInteger.ZERO : account.getId());
|
|
|
|
|
version.setModifiedBy(account == null ? BigInteger.ZERO : account.getId());
|
|
|
|
|
datasetVersionMapper.insert(version);
|
|
|
|
|
return version;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private void createLineage(BigInteger sourceTableId, BigInteger derivedTableId, String deriveType, Map<String, Object> config, LoginAccount account) {
|
|
|
|
|
DatacenterDerivedTable relation = new DatacenterDerivedTable();
|
|
|
|
|
relation.setSourceTableId(sourceTableId);
|
|
|
|
|
relation.setDerivedTableId(derivedTableId);
|
|
|
|
|
relation.setDeriveType(deriveType);
|
|
|
|
|
relation.setDeriveConfigJson(config == null ? new LinkedHashMap<>() : new LinkedHashMap<>(config));
|
|
|
|
|
relation.setStatus(0);
|
|
|
|
|
relation.setTenantId(account == null || account.getTenantId() == null ? BigInteger.ZERO : account.getTenantId());
|
|
|
|
|
relation.setDeptId(account == null || account.getDeptId() == null ? BigInteger.ZERO : account.getDeptId());
|
|
|
|
|
relation.setCreated(new Date());
|
|
|
|
|
relation.setModified(new Date());
|
|
|
|
|
relation.setCreatedBy(account == null ? BigInteger.ZERO : account.getId());
|
|
|
|
|
relation.setModifiedBy(account == null ? BigInteger.ZERO : account.getId());
|
|
|
|
|
derivedTableMapper.insert(relation);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private long copyRows(DatacenterQueryRequest queryRequest, RowMapper mapper, DatacenterTable targetTable, LoginAccount account) {
|
|
|
|
|
return iterateRows(queryRequest, row -> saveToTable(targetTable, mapper.map(row), account));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private long iterateRows(DatacenterQueryRequest queryRequest, RowConsumer consumer) {
|
|
|
|
|
long total = 0L;
|
|
|
|
|
long pageNumber = 1L;
|
|
|
|
|
while (true) {
|
|
|
|
|
queryRequest.setPageNumber(pageNumber);
|
|
|
|
|
queryRequest.setPageSize(QUERY_BATCH_SIZE);
|
|
|
|
|
Page<Row> page = queryService.queryPage(queryRequest);
|
|
|
|
|
if (page.getRecords() == null || page.getRecords().isEmpty()) {
|
|
|
|
|
break;
|
|
|
|
|
}
|
|
|
|
|
for (Row row : page.getRecords()) {
|
|
|
|
|
consumer.accept(row);
|
|
|
|
|
total++;
|
|
|
|
|
}
|
|
|
|
|
if (page.getRecords().size() < QUERY_BATCH_SIZE) {
|
|
|
|
|
break;
|
|
|
|
|
}
|
|
|
|
|
pageNumber++;
|
|
|
|
|
}
|
|
|
|
|
return total;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private void saveToTable(DatacenterTable targetTable, JSONObject data, LoginAccount account) {
|
|
|
|
|
dbHandleManager.getDbHandler().saveValue(targetTable, data, account);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private DatacenterQueryRequest buildFullQuery(DatacenterTable table) {
|
|
|
|
|
DatacenterQueryRequest queryRequest = new DatacenterQueryRequest();
|
|
|
|
|
queryRequest.setDatasetRef(registryService.resolveDatasetRef(table.getId()));
|
|
|
|
|
queryRequest.setSelectedColumns(table.getFields().stream().map(DatacenterTableField::getFieldName).toList());
|
|
|
|
|
return queryRequest;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private JSONObject mapRow(Row row) {
|
|
|
|
|
JSONObject payload = new JSONObject();
|
|
|
|
|
row.forEach(payload::put);
|
|
|
|
|
payload.remove("id");
|
|
|
|
|
payload.remove("dept_id");
|
|
|
|
|
payload.remove("tenant_id");
|
|
|
|
|
payload.remove("created");
|
|
|
|
|
payload.remove("created_by");
|
|
|
|
|
payload.remove("modified");
|
|
|
|
|
payload.remove("modified_by");
|
|
|
|
|
payload.remove("remark");
|
|
|
|
|
return payload;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private JSONObject mapDerivedRow(Row row, List<String> selectedColumns, DatacenterExcelDeriveRequest request) {
|
|
|
|
|
JSONObject payload = new JSONObject();
|
|
|
|
|
for (String column : selectedColumns) {
|
|
|
|
|
String targetName = request.getRenameMappings().getOrDefault(column, column);
|
|
|
|
|
payload.put(targetName, row.get(column));
|
|
|
|
|
}
|
|
|
|
|
return payload;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private List<String> resolveSelectedColumns(DatacenterTable sourceTable, List<String> selectedColumns) {
|
|
|
|
|
if (CollectionUtils.isEmpty(selectedColumns)) {
|
|
|
|
|
return sourceTable.getFields().stream().map(DatacenterTableField::getFieldName).toList();
|
|
|
|
|
}
|
|
|
|
|
return selectedColumns;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private List<DatacenterTableField> buildDerivedFields(DatacenterTable sourceTable, DatacenterExcelDeriveRequest request) {
|
|
|
|
|
List<String> selectedColumns = resolveSelectedColumns(sourceTable, request.getSelectedColumns());
|
|
|
|
|
Map<String, DatacenterTableField> fieldMap = new LinkedHashMap<>();
|
|
|
|
|
for (DatacenterTableField field : sourceTable.getFields()) {
|
|
|
|
|
fieldMap.put(field.getFieldName(), field);
|
|
|
|
|
}
|
|
|
|
|
List<DatacenterTableField> fields = new ArrayList<>();
|
|
|
|
|
for (String column : selectedColumns) {
|
|
|
|
|
DatacenterTableField sourceField = fieldMap.get(column);
|
|
|
|
|
if (sourceField == null) {
|
|
|
|
|
throw new BusinessException("派生字段不存在: " + column);
|
|
|
|
|
}
|
|
|
|
|
String targetName = request.getRenameMappings().getOrDefault(column, column);
|
|
|
|
|
fields.add(cloneField(sourceField, targetName, targetName));
|
|
|
|
|
}
|
|
|
|
|
return fields;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private List<DatacenterTableField> cloneFields(List<DatacenterTableField> sourceFields) {
|
|
|
|
|
List<DatacenterTableField> fields = new ArrayList<>();
|
|
|
|
|
for (DatacenterTableField field : sourceFields) {
|
|
|
|
|
fields.add(cloneField(field, field.getFieldName(), field.getFieldDesc()));
|
|
|
|
|
}
|
|
|
|
|
return fields;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private DatacenterTableField cloneField(DatacenterTableField source, String fieldName, String fieldDesc) {
|
|
|
|
|
DatacenterTableField field = new DatacenterTableField();
|
|
|
|
|
field.setFieldName(fieldName);
|
|
|
|
|
field.setSourceColumnName(source.getSourceColumnName());
|
|
|
|
|
field.setFieldDesc(fieldDesc);
|
|
|
|
|
field.setFieldType(source.getFieldType());
|
|
|
|
|
field.setJdbcType(source.getJdbcType());
|
|
|
|
|
field.setPrecision(source.getPrecision());
|
|
|
|
|
field.setScale(source.getScale());
|
|
|
|
|
field.setRequired(source.getRequired());
|
|
|
|
|
field.setQueryable(source.getQueryable());
|
|
|
|
|
field.setSortable(source.getSortable());
|
|
|
|
|
field.setWritable(source.getWritable());
|
|
|
|
|
field.setIndexed(source.getIndexed());
|
|
|
|
|
field.setOptions(source.getOptions());
|
|
|
|
|
return field;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private Map<String, Object> defaultExcelCapabilities() {
|
|
|
|
|
return Map.of("capabilities", List.of("READ_QUERY", "WRITE_MUTATION", "MATERIALIZE", "EXPORT"));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private void assertSameCatalog(List<DatacenterTable> tables) {
|
|
|
|
|
Set<BigInteger> catalogIds = new HashSet<>();
|
|
|
|
|
Set<BigInteger> sourceIds = new HashSet<>();
|
|
|
|
|
for (DatacenterTable table : tables) {
|
|
|
|
|
catalogIds.add(table.getCatalogId());
|
|
|
|
|
sourceIds.add(table.getSourceId());
|
|
|
|
|
}
|
|
|
|
|
if (catalogIds.size() > 1 || sourceIds.size() > 1) {
|
|
|
|
|
throw new BusinessException("Excel 操作暂只支持同一 workbook/catalog 下的数据集");
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private void assertSameFields(List<DatacenterTable> tables) {
|
|
|
|
|
List<String> first = tables.get(0).getFields().stream().map(DatacenterTableField::getFieldName).toList();
|
|
|
|
|
for (int i = 1; i < tables.size(); i++) {
|
|
|
|
|
List<String> current = tables.get(i).getFields().stream().map(DatacenterTableField::getFieldName).toList();
|
|
|
|
|
if (!first.equals(current)) {
|
|
|
|
|
throw new BusinessException("纵向合并仅支持同结构表");
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private List<DatacenterTable> resolveExportTables(DatacenterExcelExportRequest request) {
|
|
|
|
|
List<DatacenterTable> tables = new ArrayList<>();
|
|
|
|
|
if (request != null && !CollectionUtils.isEmpty(request.getDatasetRefs())) {
|
|
|
|
|
for (DatasetRef datasetRef : request.getDatasetRefs()) {
|
|
|
|
|
tables.add(resolveTable(datasetRef));
|
|
|
|
|
}
|
|
|
|
|
return tables;
|
|
|
|
|
}
|
|
|
|
|
if (request == null || request.getSourceId() == null) {
|
|
|
|
|
throw new BusinessException("导出需要 sourceId 或 datasetRefs");
|
|
|
|
|
}
|
|
|
|
|
tables.addAll(registryService.listManagedTables(request.getSourceId(), request.getCatalogId()));
|
|
|
|
|
return tables.stream().map(table -> registryService.getTableWithFields(table.getId())).toList();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private void writeHeaderRow(org.apache.poi.ss.usermodel.Sheet sheet, List<DatacenterTableField> fields) {
|
|
|
|
|
org.apache.poi.ss.usermodel.Row headerRow = sheet.createRow(0);
|
|
|
|
|
for (int i = 0; i < fields.size(); i++) {
|
|
|
|
|
Cell cell = headerRow.createCell(i);
|
|
|
|
|
cell.setCellValue(fields.get(i).getFieldDesc());
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private Path ensureExportDir() throws Exception {
|
|
|
|
|
Path dir = Path.of(System.getProperty("java.io.tmpdir"), "easyflow-datacenter", "exports");
|
|
|
|
|
Files.createDirectories(dir);
|
|
|
|
|
return dir;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private String buildExportFileName(String rawFileName) {
|
|
|
|
|
String baseName = rawFileName == null || rawFileName.isBlank() ? "excel_export" : extractWorkbookName(rawFileName);
|
|
|
|
|
return normalizeIdentifier(baseName) + "_" + EXPORT_TIME_FORMAT.format(LocalDateTime.now()) + ".xlsx";
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private String uniqueSheetName(String rawName, Set<String> usedSheetNames) {
|
|
|
|
|
String base = rawName == null || rawName.isBlank() ? "Sheet" : rawName;
|
|
|
|
|
base = base.length() > 31 ? base.substring(0, 31) : base;
|
|
|
|
|
String result = base;
|
|
|
|
|
int suffix = 1;
|
|
|
|
|
while (usedSheetNames.contains(result)) {
|
|
|
|
|
String suffixText = "_" + suffix++;
|
|
|
|
|
int limit = Math.max(1, 31 - suffixText.length());
|
|
|
|
|
result = base.substring(0, Math.min(base.length(), limit)) + suffixText;
|
|
|
|
|
}
|
|
|
|
|
usedSheetNames.add(result);
|
|
|
|
|
return result;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private List<DatacenterTableField> buildFields(org.apache.poi.ss.usermodel.Row headerRow, DataFormatter formatter) {
|
|
|
|
|
List<DatacenterTableField> fields = new ArrayList<>();
|
|
|
|
|
Set<String> usedNames = new HashSet<>();
|
|
|
|
|
short lastCellNum = headerRow.getLastCellNum();
|
|
|
|
|
for (int cellIndex = 0; cellIndex < lastCellNum; cellIndex++) {
|
|
|
|
|
String header = formatter.formatCellValue(headerRow.getCell(cellIndex));
|
|
|
|
|
String fieldName = normalizeIdentifier(header, cellIndex, usedNames);
|
|
|
|
|
DatacenterTableField field = new DatacenterTableField();
|
|
|
|
|
field.setFieldName(fieldName);
|
|
|
|
|
field.setSourceColumnName(header);
|
|
|
|
|
field.setFieldDesc(header == null || header.isBlank() ? fieldName : header);
|
|
|
|
|
field.setFieldType(EnumFieldType.STRING.getCode());
|
|
|
|
|
field.setJdbcType("VARCHAR");
|
|
|
|
|
field.setPrecision(255);
|
|
|
|
|
field.setScale(0);
|
|
|
|
|
field.setRequired(0);
|
|
|
|
|
field.setQueryable(1);
|
|
|
|
|
field.setSortable(1);
|
|
|
|
|
field.setWritable(1);
|
|
|
|
|
field.setIndexed(0);
|
|
|
|
|
fields.add(field);
|
|
|
|
|
}
|
|
|
|
|
return fields;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private String normalizeMode(String value, String defaultValue) {
|
|
|
|
|
return value == null || value.isBlank() ? defaultValue : value.trim().toUpperCase(Locale.ROOT);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private String resolveSplitPrefix(DatacenterExcelSplitRequest request, String fallback) {
|
|
|
|
|
if (request != null && request.getTargetNamePrefix() != null && !request.getTargetNamePrefix().isBlank()) {
|
|
|
|
|
return request.getTargetNamePrefix();
|
|
|
|
|
}
|
|
|
|
|
return fallback + "_split";
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private String normalizeLogicalName(String tableName, String deriveType) {
|
|
|
|
|
if (tableName != null && !tableName.isBlank()) {
|
|
|
|
|
return tableName;
|
|
|
|
|
}
|
|
|
|
|
return deriveType.toLowerCase(Locale.ROOT) + "_" + System.currentTimeMillis();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private String uniqueTableName(BigInteger sourceId, BigInteger catalogId, String rawName) {
|
|
|
|
|
String baseName = rawName == null || rawName.isBlank() ? "dataset" : rawName;
|
|
|
|
|
baseName = baseName.trim();
|
|
|
|
|
String result = baseName;
|
|
|
|
|
int suffix = 1;
|
|
|
|
|
while (tableNameExists(sourceId, catalogId, result)) {
|
|
|
|
|
result = baseName + "_" + suffix++;
|
|
|
|
|
}
|
|
|
|
|
return result;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private boolean tableNameExists(BigInteger sourceId, BigInteger catalogId, String tableName) {
|
|
|
|
|
List<DatacenterTable> tables = registryService.listManagedTables(sourceId, catalogId);
|
|
|
|
|
return tables.stream().anyMatch(table -> tableName.equals(table.getTableName()));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private String extractWorkbookName(String originalFileName) {
|
|
|
|
|
if (originalFileName == null || originalFileName.isBlank()) {
|
|
|
|
|
return "excel_workbook";
|
|
|
|
|
}
|
|
|
|
|
int index = originalFileName.lastIndexOf('.');
|
|
|
|
|
return index > 0 ? originalFileName.substring(0, index) : originalFileName;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private String buildMaterializedTableName(BigInteger sourceId, int sheetIndex) {
|
|
|
|
|
long snowId = new SnowFlakeIDKeyGenerator().nextId();
|
|
|
|
|
return "tb_excel_" + sourceId + "_" + sheetIndex + "_" + snowId;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private String normalizeIdentifier(String raw) {
|
|
|
|
|
if (raw == null || raw.isBlank()) {
|
|
|
|
|
return "value";
|
|
|
|
|
}
|
|
|
|
|
String normalized = raw.trim().toLowerCase(Locale.ROOT).replaceAll("[^a-z0-9_\\u4e00-\\u9fa5]+", "_");
|
|
|
|
|
normalized = normalized.replaceAll("_+", "_");
|
|
|
|
|
if (normalized.isBlank()) {
|
|
|
|
|
return "value";
|
|
|
|
|
}
|
|
|
|
|
return normalized;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private String normalizeIdentifier(String raw, int index, Set<String> usedNames) {
|
|
|
|
|
String value = normalizeIdentifier(raw);
|
|
|
|
|
if (value.isBlank() || "value".equals(value)) {
|
|
|
|
|
value = "col_" + (index + 1);
|
|
|
|
|
}
|
|
|
|
|
if (Character.isDigit(value.charAt(0))) {
|
|
|
|
|
value = "col_" + value;
|
|
|
|
|
}
|
|
|
|
|
String result = value;
|
|
|
|
|
int suffix = 1;
|
|
|
|
|
while (usedNames.contains(result)) {
|
|
|
|
|
result = value + "_" + suffix++;
|
|
|
|
|
}
|
|
|
|
|
usedNames.add(result);
|
|
|
|
|
return result;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private String stringify(Object value) {
|
|
|
|
|
return value == null ? null : String.valueOf(value);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private static final class Holder {
|
|
|
|
|
private DatacenterTable targetTable;
|
|
|
|
|
private int batchNo;
|
|
|
|
|
private int currentSize;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private interface RowConsumer {
|
|
|
|
|
void accept(Row row);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private interface RowMapper {
|
|
|
|
|
JSONObject map(Row row);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private static final class QueryWrapperWrapper {
|
|
|
|
|
private final BigInteger tableId;
|
|
|
|
|
|
|
|
|
|
private QueryWrapperWrapper(BigInteger tableId) {
|
|
|
|
|
this.tableId = tableId;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private int nextVersionNo(DatacenterDatasetVersionMapper mapper) {
|
|
|
|
|
return mapper.selectListByQuery(com.mybatisflex.core.query.QueryWrapper.create()
|
|
|
|
|
.eq(DatacenterDatasetVersion::getTableId, tableId)
|
|
|
|
|
.orderBy("version_no desc"))
|
|
|
|
|
.stream()
|
|
|
|
|
.findFirst()
|
|
|
|
|
.map(version -> version.getVersionNo() + 1)
|
|
|
|
|
.orElse(1);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|