feat(aihr): decouple batch upload from processing via async queue
批量上传治本:POST /doc/upload-async 只暂存+入队即秒回,后台 worker
(并发2)复用 saveDocument 链路逐条解析/归类/向量化;队列表
aihr_knowledge_upload_item 状态机 0/1/2/3,完成/失败 CAS 写入,同名
文件串行防撞 attach 唯一键,卡住30分钟清扫重置,失败暂存保留72h。
新增 GET /doc/upload-items 批次查询与 POST /doc/upload-items/{id}/retry
单文件重试。资料处理页批量导入改走异步:提交即返回,新增本批次
进度面板(排队/加工中/完成/失败+进度条+单条重试),轮询防过期回写。
同步接口保留给 SOP 页单文件即时预览。
This commit is contained in:
+23
-1
@@ -7,10 +7,13 @@ import org.dromara.aihr.domain.AihrSopDto.LocalImportTaskResponse;
|
||||
import org.dromara.aihr.domain.AihrSopDto.ProcessingOverviewResponse;
|
||||
import org.dromara.aihr.domain.AihrSopDto.SearchRequest;
|
||||
import org.dromara.aihr.domain.AihrSopDto.SearchResponse;
|
||||
import org.dromara.aihr.domain.AihrSopDto.UploadEnqueueResponse;
|
||||
import org.dromara.aihr.domain.AihrSopDto.UploadItemResponse;
|
||||
import org.dromara.aihr.domain.AihrSopDto.UploadResponse;
|
||||
import org.dromara.aihr.domain.AihrSopDto.VectorIndexStatusResponse;
|
||||
import org.dromara.aihr.domain.AihrSopDto.VectorizeResponse;
|
||||
import org.dromara.aihr.service.AihrSopSeedService;
|
||||
import org.dromara.aihr.service.AihrUploadQueueService;
|
||||
import org.dromara.common.core.domain.R;
|
||||
import org.springframework.web.bind.annotation.GetMapping;
|
||||
import org.springframework.web.bind.annotation.PathVariable;
|
||||
@@ -25,7 +28,7 @@ import org.springframework.web.multipart.MultipartFile;
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* SOP 知识库 seed API。
|
||||
* SOP 知识库 API。
|
||||
*/
|
||||
@RequiredArgsConstructor
|
||||
@RestController
|
||||
@@ -33,6 +36,25 @@ import java.util.List;
|
||||
public class AihrSopController {
|
||||
|
||||
private final AihrSopSeedService sopSeedService;
|
||||
private final AihrUploadQueueService uploadQueueService;
|
||||
|
||||
@PostMapping("/doc/upload-async")
|
||||
public R<UploadEnqueueResponse> uploadAsync(@RequestPart("file") MultipartFile file,
|
||||
@RequestParam(value = "category", required = false) String category,
|
||||
@RequestParam(value = "batchId", required = false) String batchId) {
|
||||
return R.ok(uploadQueueService.enqueue(file, category, batchId));
|
||||
}
|
||||
|
||||
@GetMapping("/doc/upload-items")
|
||||
public R<List<UploadItemResponse>> uploadItems(@RequestParam(value = "batchId", required = false) String batchId,
|
||||
@RequestParam(value = "limit", required = false, defaultValue = "50") int limit) {
|
||||
return R.ok(uploadQueueService.items(batchId, limit));
|
||||
}
|
||||
|
||||
@PostMapping("/doc/upload-items/{id}/retry")
|
||||
public R<UploadItemResponse> retryUploadItem(@PathVariable Long id) {
|
||||
return R.ok(uploadQueueService.retry(id));
|
||||
}
|
||||
|
||||
@PostMapping("/search")
|
||||
public R<SearchResponse> search(@RequestBody SearchRequest request) {
|
||||
|
||||
+17
@@ -129,4 +129,21 @@ public final class AihrSopDto {
|
||||
Integer failedTasks
|
||||
) {
|
||||
}
|
||||
|
||||
public record UploadEnqueueResponse(Long itemId, String batchId, String fileName, String status) {
|
||||
}
|
||||
|
||||
public record UploadItemResponse(
|
||||
Long id,
|
||||
String batchId,
|
||||
String fileName,
|
||||
String category,
|
||||
Integer status,
|
||||
String statusLabel,
|
||||
String error,
|
||||
String docId,
|
||||
Integer fragmentCount,
|
||||
String updatedAt
|
||||
) {
|
||||
}
|
||||
}
|
||||
|
||||
+23
-11
@@ -23,6 +23,7 @@ import org.dromara.aihr.domain.AihrSopDto.SnippetResponse;
|
||||
import org.dromara.aihr.domain.AihrSopDto.UploadResponse;
|
||||
import org.dromara.aihr.domain.AihrSopDto.VectorIndexStatusResponse;
|
||||
import org.dromara.aihr.domain.AihrSopDto.VectorizeResponse;
|
||||
import org.dromara.common.core.exception.ServiceException;
|
||||
import org.apache.tika.metadata.Metadata;
|
||||
import org.apache.tika.metadata.TikaCoreProperties;
|
||||
import org.apache.tika.parser.AutoDetectParser;
|
||||
@@ -122,19 +123,30 @@ public class AihrSopSeedService {
|
||||
|
||||
public UploadResponse uploadDoc(MultipartFile file, String category) {
|
||||
if (file == null || file.isEmpty()) {
|
||||
throw new IllegalArgumentException("上传文件不能为空");
|
||||
throw new ServiceException("上传文件不能为空");
|
||||
}
|
||||
String fileName = Optional.ofNullable(file.getOriginalFilename()).orElse("knowledge.txt").trim();
|
||||
if (!supportedFile(fileName)) {
|
||||
throw new IllegalArgumentException("仅支持 txt/md/markdown/pdf/doc/docx/xls/xlsx/ppt/pptx 文件");
|
||||
throw new ServiceException("仅支持 txt/md/markdown/pdf/doc/docx/xls/xlsx/ppt/pptx/图片文件");
|
||||
}
|
||||
if (file.getSize() > 100 * 1024 * 1024) {
|
||||
throw new IllegalArgumentException("文件不能超过 100MB");
|
||||
throw new ServiceException("文件不能超过 100MB");
|
||||
}
|
||||
|
||||
return saveDocument(fileName, category, () -> ossService.upload(file), () -> readContent(file, fileName), () -> fileFingerprint(file));
|
||||
}
|
||||
|
||||
/**
|
||||
* 供后台上传队列(AihrUploadQueueService)复用的完整加工链路:OSS 上传、解析、归类、切片、向量化。
|
||||
* 与目录导入同源,入参为暂存文件路径。
|
||||
*/
|
||||
public UploadResponse processStagedDocument(String fileName, String category, Path stagedFile) {
|
||||
return saveDocument(fileName, category,
|
||||
() -> ossService.upload(stagedFile.toFile()),
|
||||
() -> readContent(stagedFile, fileName),
|
||||
() -> fileFingerprint(stagedFile));
|
||||
}
|
||||
|
||||
public VectorizeResponse vectorizeMissing() {
|
||||
List<VectorizeTarget> targets = jdbcTemplate.query("""
|
||||
select f.knowledge_id,
|
||||
@@ -454,7 +466,7 @@ public class AihrSopSeedService {
|
||||
try {
|
||||
content = contentReader.get();
|
||||
if (content.isBlank()) {
|
||||
throw new IllegalArgumentException("文件内容不能为空");
|
||||
throw new ServiceException("文件内容不能为空");
|
||||
}
|
||||
} catch (RuntimeException e) {
|
||||
throw e;
|
||||
@@ -1967,7 +1979,7 @@ public class AihrSopSeedService {
|
||||
new AutoDetectParser().parse(input, handler, metadata, new ParseContext());
|
||||
return normalizeExtractedText(handler.toString());
|
||||
} catch (Exception e) {
|
||||
throw new IllegalArgumentException("文件解析失败");
|
||||
throw new ServiceException("文件解析失败");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1995,7 +2007,7 @@ public class AihrSopSeedService {
|
||||
new AutoDetectParser().parse(input, handler, metadata, new ParseContext());
|
||||
return normalizeExtractedText(handler.toString());
|
||||
} catch (Exception e) {
|
||||
throw new IllegalArgumentException("文件解析失败");
|
||||
throw new ServiceException("文件解析失败");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2003,7 +2015,7 @@ public class AihrSopSeedService {
|
||||
try (InputStream input = file.getInputStream()) {
|
||||
return fileFingerprint(input);
|
||||
} catch (IOException e) {
|
||||
throw new IllegalArgumentException("文件指纹计算失败");
|
||||
throw new ServiceException("文件指纹计算失败");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2011,7 +2023,7 @@ public class AihrSopSeedService {
|
||||
try (InputStream input = Files.newInputStream(file)) {
|
||||
return fileFingerprint(input);
|
||||
} catch (IOException e) {
|
||||
throw new IllegalArgumentException("文件指纹计算失败");
|
||||
throw new ServiceException("文件指纹计算失败");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2054,7 +2066,7 @@ public class AihrSopSeedService {
|
||||
try {
|
||||
return new String(file.getBytes(), StandardCharsets.UTF_8);
|
||||
} catch (Exception e) {
|
||||
throw new IllegalArgumentException("文件读取失败");
|
||||
throw new ServiceException("文件读取失败");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2062,7 +2074,7 @@ public class AihrSopSeedService {
|
||||
try {
|
||||
return Files.readString(file, StandardCharsets.UTF_8);
|
||||
} catch (Exception e) {
|
||||
throw new IllegalArgumentException("文件读取失败");
|
||||
throw new ServiceException("文件读取失败");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2140,7 +2152,7 @@ public class AihrSopSeedService {
|
||||
return AUTO_CATEGORY.equals(value) || "智能归类".equals(value) || "自动归类".equals(value);
|
||||
}
|
||||
|
||||
private static boolean supportedFile(String fileName) {
|
||||
static boolean supportedFile(String fileName) {
|
||||
String lower = fileName.toLowerCase();
|
||||
return textFile(lower)
|
||||
|| imageFile(lower)
|
||||
|
||||
+395
@@ -0,0 +1,395 @@
|
||||
package org.dromara.aihr.service;
|
||||
|
||||
import jakarta.annotation.PostConstruct;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.dromara.aihr.domain.AihrSopDto.UploadEnqueueResponse;
|
||||
import org.dromara.aihr.domain.AihrSopDto.UploadItemResponse;
|
||||
import org.dromara.aihr.domain.AihrSopDto.UploadResponse;
|
||||
import org.dromara.common.core.exception.ServiceException;
|
||||
import org.springframework.beans.factory.annotation.Value;
|
||||
import org.springframework.dao.DataAccessException;
|
||||
import org.springframework.jdbc.core.JdbcTemplate;
|
||||
import org.springframework.stereotype.Service;
|
||||
import org.springframework.web.multipart.MultipartFile;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.nio.file.StandardCopyOption;
|
||||
import java.sql.ResultSet;
|
||||
import java.sql.SQLException;
|
||||
import java.time.format.DateTimeFormatter;
|
||||
import java.util.List;
|
||||
import java.util.UUID;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.Semaphore;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
/**
|
||||
* 批量上传队列:上传请求只做「暂存文件 + 入队 + 秒回」,解析/归类/向量化由本服务的后台 worker
|
||||
* 复用 {@link AihrSopSeedService#processStagedDocument} 逐条加工。单文件失败不影响批次,支持按条重试。
|
||||
*/
|
||||
@Service
|
||||
@Slf4j
|
||||
public class AihrUploadQueueService {
|
||||
|
||||
private static final String TENANT_ID = "000000";
|
||||
private static final DateTimeFormatter TIME_FORMATTER = DateTimeFormatter.ofPattern("MM-dd HH:mm:ss");
|
||||
private static final long MAX_FILE_BYTES = 100L * 1024 * 1024;
|
||||
private static final int MAX_ERROR_CHARS = 480;
|
||||
/** 加工并发上限:链路含 LLM 归类与 embedding 外呼,并发过高会互相争抢配额并拖慢单文件耗时。 */
|
||||
private static final int WORKER_PERMITS = 2;
|
||||
/** 处理中超过该时长视为孤儿(服务重启/线程异常丢失),由清扫任务重置回待处理。需大于最慢单文件加工时长。 */
|
||||
private static final int STUCK_MINUTES = 30;
|
||||
/** 失败条目的暂存文件保留时长,超时由清扫回收磁盘(之后重试会提示重新上传)。 */
|
||||
private static final int STAGING_RETENTION_HOURS = 72;
|
||||
|
||||
private final AihrSopSeedService sopSeedService;
|
||||
private final JdbcTemplate jdbcTemplate;
|
||||
private final ScheduledExecutorService scheduledExecutorService;
|
||||
private final Path stagingRoot;
|
||||
private final Semaphore workerPermits = new Semaphore(WORKER_PERMITS);
|
||||
private volatile boolean tableReady;
|
||||
|
||||
public AihrUploadQueueService(AihrSopSeedService sopSeedService,
|
||||
JdbcTemplate jdbcTemplate,
|
||||
ScheduledExecutorService scheduledExecutorService,
|
||||
@Value("${aihr.upload.staging:./.data/staging}") String stagingRootConfig) {
|
||||
this.sopSeedService = sopSeedService;
|
||||
this.jdbcTemplate = jdbcTemplate;
|
||||
this.scheduledExecutorService = scheduledExecutorService;
|
||||
this.stagingRoot = Path.of(stagingRootConfig).toAbsolutePath().normalize();
|
||||
}
|
||||
|
||||
@PostConstruct
|
||||
void scheduleSweep() {
|
||||
scheduledExecutorService.scheduleWithFixedDelay(this::sweepAndTrigger, 20, 15, TimeUnit.SECONDS);
|
||||
}
|
||||
|
||||
public UploadEnqueueResponse enqueue(MultipartFile file, String category, String batchId) {
|
||||
if (file == null || file.isEmpty()) {
|
||||
throw new ServiceException("上传文件不能为空");
|
||||
}
|
||||
String fileName = sanitizeFileName(file.getOriginalFilename());
|
||||
if (!AihrSopSeedService.supportedFile(fileName)) {
|
||||
throw new ServiceException("仅支持 txt/md/markdown/pdf/doc/docx/xls/xlsx/ppt/pptx/图片文件");
|
||||
}
|
||||
if (file.getSize() > MAX_FILE_BYTES) {
|
||||
throw new ServiceException("文件不能超过 100MB");
|
||||
}
|
||||
ensureTable();
|
||||
String batch = sanitizeBatchId(batchId);
|
||||
Path staged = stagingPath(fileName);
|
||||
try (var input = file.getInputStream()) {
|
||||
Files.createDirectories(staged.getParent());
|
||||
Files.copy(input, staged, StandardCopyOption.REPLACE_EXISTING);
|
||||
} catch (IOException e) {
|
||||
throw new ServiceException("暂存文件写入失败:" + e.getMessage());
|
||||
}
|
||||
jdbcTemplate.update("""
|
||||
insert into aihr_knowledge_upload_item
|
||||
(tenant_id, batch_id, file_name, category, staging_path, status, create_time, update_time)
|
||||
values (?, ?, ?, ?, ?, 0, now(), now())
|
||||
""", TENANT_ID, batch, fileName, category == null ? "" : category.trim(), staged.toString());
|
||||
// staging_path 含 UUID 全局唯一,用它反查主键,避免 last_insert_id 的跨连接问题
|
||||
Long itemId = jdbcTemplate.query("""
|
||||
select id from aihr_knowledge_upload_item where tenant_id = ? and staging_path = ?
|
||||
""", rs -> rs.next() ? rs.getLong(1) : null, TENANT_ID, staged.toString());
|
||||
triggerProcessing();
|
||||
return new UploadEnqueueResponse(itemId, batch, fileName, "queued");
|
||||
}
|
||||
|
||||
public List<UploadItemResponse> items(String batchId, int limit) {
|
||||
ensureTable();
|
||||
int safeLimit = Math.max(1, Math.min(limit <= 0 ? 50 : limit, 200));
|
||||
String batch = sanitizeBatchId(batchId);
|
||||
if (batchId == null || batchId.isBlank()) {
|
||||
return jdbcTemplate.query("""
|
||||
select id, batch_id, file_name, category, status, error, doc_id, fragment_count, update_time
|
||||
from aihr_knowledge_upload_item
|
||||
where tenant_id = ?
|
||||
order by id desc
|
||||
limit %d
|
||||
""".formatted(safeLimit), this::mapItem, TENANT_ID);
|
||||
}
|
||||
return jdbcTemplate.query("""
|
||||
select id, batch_id, file_name, category, status, error, doc_id, fragment_count, update_time
|
||||
from aihr_knowledge_upload_item
|
||||
where tenant_id = ? and batch_id = ?
|
||||
order by id asc
|
||||
limit %d
|
||||
""".formatted(safeLimit), this::mapItem, TENANT_ID, batch);
|
||||
}
|
||||
|
||||
public UploadItemResponse retry(Long id) {
|
||||
ensureTable();
|
||||
if (id == null) {
|
||||
throw new ServiceException("重试条目不存在");
|
||||
}
|
||||
ItemRow row = jdbcTemplate.query("""
|
||||
select id, file_name, category, staging_path, status
|
||||
from aihr_knowledge_upload_item
|
||||
where tenant_id = ? and id = ?
|
||||
""", rs -> rs.next()
|
||||
? new ItemRow(rs.getLong("id"), rs.getString("file_name"), rs.getString("category"), rs.getString("staging_path"), rs.getInt("status"))
|
||||
: null, TENANT_ID, id);
|
||||
if (row == null) {
|
||||
throw new ServiceException("重试条目不存在");
|
||||
}
|
||||
if (row.status() != 3) {
|
||||
throw new ServiceException("仅失败状态的文件可重试");
|
||||
}
|
||||
if (!Files.isRegularFile(Path.of(row.stagingPath()))) {
|
||||
markFailed(id, "暂存文件已丢失,请重新上传该文件");
|
||||
throw new ServiceException("暂存文件已丢失,请重新上传该文件");
|
||||
}
|
||||
int updated = jdbcTemplate.update("""
|
||||
update aihr_knowledge_upload_item
|
||||
set status = 0, error = null, update_time = now()
|
||||
where tenant_id = ? and id = ? and status = 3
|
||||
""", TENANT_ID, id);
|
||||
if (updated == 0) {
|
||||
throw new ServiceException("仅失败状态的文件可重试");
|
||||
}
|
||||
triggerProcessing();
|
||||
return item(id);
|
||||
}
|
||||
|
||||
/** 领取待处理条目并在受限并发下加工;每领取一条占用一个许可,处理完释放后继续领取。 */
|
||||
private void triggerProcessing() {
|
||||
while (workerPermits.tryAcquire()) {
|
||||
Long claimed = claimNext();
|
||||
if (claimed == null) {
|
||||
workerPermits.release();
|
||||
return;
|
||||
}
|
||||
CompletableFuture.runAsync(() -> {
|
||||
try {
|
||||
processItem(claimed);
|
||||
} finally {
|
||||
workerPermits.release();
|
||||
triggerProcessing();
|
||||
}
|
||||
}, scheduledExecutorService);
|
||||
}
|
||||
}
|
||||
|
||||
private Long claimNext() {
|
||||
try {
|
||||
// 同名文件不并发领取:saveDocument 对同 knowledge+fileName 是替换语义,并发会撞 attach 唯一键
|
||||
List<Long> candidates = jdbcTemplate.query("""
|
||||
select i.id from aihr_knowledge_upload_item i
|
||||
where i.tenant_id = ? and i.status = 0
|
||||
and not exists (
|
||||
select 1 from (select file_name from aihr_knowledge_upload_item where tenant_id = ? and status = 1) p
|
||||
where p.file_name = i.file_name
|
||||
)
|
||||
order by i.id asc
|
||||
limit 5
|
||||
""", (rs, rowNum) -> rs.getLong(1), TENANT_ID, TENANT_ID);
|
||||
for (Long id : candidates) {
|
||||
int updated = jdbcTemplate.update("""
|
||||
update aihr_knowledge_upload_item
|
||||
set status = 1, update_time = now()
|
||||
where tenant_id = ? and id = ? and status = 0
|
||||
""", TENANT_ID, id);
|
||||
if (updated == 1) {
|
||||
return id;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
} catch (DataAccessException e) {
|
||||
log.warn("upload queue claim failed: {}", e.getMessage());
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
private void processItem(Long id) {
|
||||
ItemRow row = jdbcTemplate.query("""
|
||||
select id, file_name, category, staging_path, status
|
||||
from aihr_knowledge_upload_item
|
||||
where tenant_id = ? and id = ?
|
||||
""", rs -> rs.next()
|
||||
? new ItemRow(rs.getLong("id"), rs.getString("file_name"), rs.getString("category"), rs.getString("staging_path"), rs.getInt("status"))
|
||||
: null, TENANT_ID, id);
|
||||
if (row == null) {
|
||||
return;
|
||||
}
|
||||
Path staged = Path.of(row.stagingPath());
|
||||
if (!Files.isRegularFile(staged)) {
|
||||
markFailed(id, "暂存文件已丢失,请重新上传该文件");
|
||||
return;
|
||||
}
|
||||
try {
|
||||
UploadResponse result = sopSeedService.processStagedDocument(row.fileName(), row.category(), staged);
|
||||
// CAS:仅当仍是本工人持有的「处理中」才写完成,防止清扫重置后被后来的工人覆盖状态
|
||||
int updated = jdbcTemplate.update("""
|
||||
update aihr_knowledge_upload_item
|
||||
set status = 2, error = null, doc_id = ?, fragment_count = ?, update_time = now()
|
||||
where tenant_id = ? and id = ? and status = 1
|
||||
""", result.docId(), result.fragments(), TENANT_ID, id);
|
||||
if (updated == 1) {
|
||||
deleteQuietly(staged);
|
||||
} else {
|
||||
log.info("upload item {} finished but row was re-claimed, skip status write", id);
|
||||
}
|
||||
} catch (Exception e) {
|
||||
log.warn("upload item {} process failed: {}", id, e.getMessage());
|
||||
markFailed(id, e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
/** 清扫卡在「处理中」的孤儿条目(重启/异常导致),重置回待处理后再触发一轮加工;顺带回收过期暂存文件。 */
|
||||
private void sweepAndTrigger() {
|
||||
try {
|
||||
ensureTable();
|
||||
int reset = jdbcTemplate.update("""
|
||||
update aihr_knowledge_upload_item
|
||||
set status = 0, update_time = now()
|
||||
where tenant_id = ? and status = 1 and update_time < (now() - interval %d minute)
|
||||
""".formatted(STUCK_MINUTES), TENANT_ID);
|
||||
if (reset > 0) {
|
||||
log.info("upload queue reset {} stuck items", reset);
|
||||
}
|
||||
cleanExpiredStaging();
|
||||
triggerProcessing();
|
||||
} catch (Exception e) {
|
||||
log.debug("upload queue sweep skipped: {}", e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
/** 回收超过保留期的暂存文件,防失败条目/入库异常导致磁盘无限增长;活跃条目远不会到期。 */
|
||||
private void cleanExpiredStaging() {
|
||||
if (!Files.isDirectory(stagingRoot)) {
|
||||
return;
|
||||
}
|
||||
long cutoffMillis = System.currentTimeMillis() - STAGING_RETENTION_HOURS * 3600_000L;
|
||||
try (var stream = Files.list(stagingRoot)) {
|
||||
stream.filter(Files::isRegularFile)
|
||||
.filter(path -> {
|
||||
try {
|
||||
return Files.getLastModifiedTime(path).toMillis() < cutoffMillis;
|
||||
} catch (IOException e) {
|
||||
return false;
|
||||
}
|
||||
})
|
||||
.forEach(AihrUploadQueueService::deleteQuietly);
|
||||
} catch (IOException e) {
|
||||
log.debug("staging cleanup skipped: {}", e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
private void markFailed(Long id, String message) {
|
||||
String error = message == null || message.isBlank() ? "处理失败" : message;
|
||||
if (error.length() > MAX_ERROR_CHARS) {
|
||||
error = error.substring(0, MAX_ERROR_CHARS);
|
||||
}
|
||||
// 不覆盖已被清扫重置(0)或已完成(2)的行:处理中标失败,失败态允许更新错误信息(重试路径)
|
||||
jdbcTemplate.update("""
|
||||
update aihr_knowledge_upload_item
|
||||
set status = 3, error = ?, update_time = now()
|
||||
where tenant_id = ? and id = ? and status in (1, 3)
|
||||
""", error, TENANT_ID, id);
|
||||
}
|
||||
|
||||
private UploadItemResponse item(Long id) {
|
||||
List<UploadItemResponse> rows = jdbcTemplate.query("""
|
||||
select id, batch_id, file_name, category, status, error, doc_id, fragment_count, update_time
|
||||
from aihr_knowledge_upload_item
|
||||
where tenant_id = ? and id = ?
|
||||
""", this::mapItem, TENANT_ID, id);
|
||||
return rows.isEmpty() ? null : rows.get(0);
|
||||
}
|
||||
|
||||
private UploadItemResponse mapItem(ResultSet rs, int rowNum) throws SQLException {
|
||||
int status = rs.getInt("status");
|
||||
java.sql.Timestamp updateTime = rs.getTimestamp("update_time");
|
||||
return new UploadItemResponse(
|
||||
rs.getLong("id"),
|
||||
rs.getString("batch_id"),
|
||||
rs.getString("file_name"),
|
||||
rs.getString("category"),
|
||||
status,
|
||||
statusLabel(status),
|
||||
rs.getString("error"),
|
||||
rs.getString("doc_id"),
|
||||
rs.getObject("fragment_count", Integer.class),
|
||||
updateTime == null ? "" : updateTime.toLocalDateTime().format(TIME_FORMATTER)
|
||||
);
|
||||
}
|
||||
|
||||
private static String statusLabel(int status) {
|
||||
return switch (status) {
|
||||
case 0 -> "排队中";
|
||||
case 1 -> "加工中";
|
||||
case 2 -> "已完成";
|
||||
case 3 -> "失败";
|
||||
default -> "未知";
|
||||
};
|
||||
}
|
||||
|
||||
private Path stagingPath(String fileName) {
|
||||
return stagingRoot.resolve(UUID.randomUUID() + "_" + fileName);
|
||||
}
|
||||
|
||||
/** 去掉路径分隔与控制字符,防止暂存路径逃逸出 stagingRoot。 */
|
||||
private static String sanitizeFileName(String original) {
|
||||
String name = original == null || original.isBlank() ? "knowledge.txt" : original.trim();
|
||||
name = name.replaceAll("[\\\\/\\r\\n\\t]", "_");
|
||||
if (name.length() > 200) {
|
||||
name = name.substring(name.length() - 200);
|
||||
}
|
||||
return name;
|
||||
}
|
||||
|
||||
private static String sanitizeBatchId(String batchId) {
|
||||
if (batchId == null || batchId.isBlank()) {
|
||||
return "b-" + UUID.randomUUID();
|
||||
}
|
||||
String cleaned = batchId.trim().replaceAll("[^A-Za-z0-9_-]", "");
|
||||
return cleaned.isEmpty() || cleaned.length() > 64 ? "b-" + UUID.randomUUID() : cleaned;
|
||||
}
|
||||
|
||||
private static void deleteQuietly(Path path) {
|
||||
try {
|
||||
Files.deleteIfExists(path);
|
||||
} catch (IOException ignored) {
|
||||
// 暂存清理失败不影响业务,留给下次清扫。
|
||||
}
|
||||
}
|
||||
|
||||
private void ensureTable() {
|
||||
if (tableReady) {
|
||||
return;
|
||||
}
|
||||
synchronized (this) {
|
||||
if (tableReady) {
|
||||
return;
|
||||
}
|
||||
jdbcTemplate.execute("""
|
||||
CREATE TABLE IF NOT EXISTS `aihr_knowledge_upload_item` (
|
||||
`id` bigint NOT NULL AUTO_INCREMENT COMMENT '主键',
|
||||
`tenant_id` varchar(20) DEFAULT '000000' COMMENT '租户编号',
|
||||
`batch_id` varchar(64) NOT NULL COMMENT '上传批次',
|
||||
`file_name` varchar(255) NOT NULL COMMENT '文件名',
|
||||
`category` varchar(100) DEFAULT '' COMMENT '目标分类',
|
||||
`staging_path` varchar(500) NOT NULL COMMENT '暂存文件路径',
|
||||
`status` tinyint NOT NULL DEFAULT 0 COMMENT '0待处理 1处理中 2完成 3失败',
|
||||
`error` varchar(500) DEFAULT NULL COMMENT '失败原因',
|
||||
`doc_id` varchar(80) DEFAULT NULL COMMENT '入库文档ID',
|
||||
`fragment_count` int DEFAULT NULL COMMENT '片段数',
|
||||
`create_time` datetime DEFAULT NULL COMMENT '创建时间',
|
||||
`update_time` datetime DEFAULT NULL COMMENT '更新时间',
|
||||
PRIMARY KEY (`id`),
|
||||
KEY `idx_aihr_upload_item_batch` (`tenant_id`, `batch_id`, `id`),
|
||||
KEY `idx_aihr_upload_item_status` (`tenant_id`, `status`, `update_time`)
|
||||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='知识库批量上传队列';
|
||||
""");
|
||||
tableReady = true;
|
||||
}
|
||||
}
|
||||
|
||||
private record ItemRow(Long id, String fileName, String category, String stagingPath, int status) {
|
||||
}
|
||||
}
|
||||
@@ -124,3 +124,22 @@ VALUES
|
||||
ON DUPLICATE KEY UPDATE
|
||||
`content` = VALUES(`content`),
|
||||
`update_time` = NOW();
|
||||
|
||||
-- 知识库批量上传队列:上传秒回后由后台 worker 加工,支持单文件重试
|
||||
CREATE TABLE IF NOT EXISTS `aihr_knowledge_upload_item` (
|
||||
`id` bigint NOT NULL AUTO_INCREMENT COMMENT '主键',
|
||||
`tenant_id` varchar(20) DEFAULT '000000' COMMENT '租户编号',
|
||||
`batch_id` varchar(64) NOT NULL COMMENT '上传批次',
|
||||
`file_name` varchar(255) NOT NULL COMMENT '文件名',
|
||||
`category` varchar(100) DEFAULT '' COMMENT '目标分类',
|
||||
`staging_path` varchar(500) NOT NULL COMMENT '暂存文件路径',
|
||||
`status` tinyint NOT NULL DEFAULT 0 COMMENT '0待处理 1处理中 2完成 3失败',
|
||||
`error` varchar(500) DEFAULT NULL COMMENT '失败原因',
|
||||
`doc_id` varchar(80) DEFAULT NULL COMMENT '入库文档ID',
|
||||
`fragment_count` int DEFAULT NULL COMMENT '片段数',
|
||||
`create_time` datetime DEFAULT NULL COMMENT '创建时间',
|
||||
`update_time` datetime DEFAULT NULL COMMENT '更新时间',
|
||||
PRIMARY KEY (`id`),
|
||||
KEY `idx_aihr_upload_item_batch` (`tenant_id`, `batch_id`, `id`),
|
||||
KEY `idx_aihr_upload_item_status` (`tenant_id`, `status`, `update_time`)
|
||||
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='知识库批量上传队列';
|
||||
|
||||
Reference in New Issue
Block a user