From b0ea0f5ffac642ae78dc9952c29f08613e9f80cb Mon Sep 17 00:00:00 2001 From: luke Date: Fri, 20 Jun 2025 01:05:21 +0800 Subject: [PATCH] fix --- .../doc/service/FileImporterDispatcher.java | 45 +++++ .../doc/service/iface/DocumentImporter.java | 11 ++ .../impl/AbstractBaseFileImporter.java | 176 +++++++++++------- .../doc/service/impl/ExcelImporter.java | 5 +- .../doc/service/impl/MarkdownImporter.java | 5 +- .../domain/doc/service/impl/PdfImporter.java | 5 +- .../domain/doc/service/impl/WordImporter.java | 5 +- .../cache/iface/FileCacheService.java | 6 + .../cache/impl/RedisFileCacheServiceImpl.java | 5 + .../infrastructure/config/ConstantConfig.java | 2 +- .../infrastructure/config/DynamicConfig.java | 22 +-- .../config/ObjectStorageProperties.java | 2 +- .../north/controller/FileWriteController.java | 33 ++-- .../util/RateLimiterManager.java | 32 ++++ .../resources/application-dev-mac.properties | 2 +- .../application-dev-windows.properties | 2 +- .../resources/application-docker.properties | 2 +- 17 files changed, 254 insertions(+), 106 deletions(-) create mode 100644 src/main/java/com/knowledge/base/domain/doc/service/FileImporterDispatcher.java create mode 100644 src/main/java/com/knowledge/base/infrastructure/util/RateLimiterManager.java diff --git a/src/main/java/com/knowledge/base/domain/doc/service/FileImporterDispatcher.java b/src/main/java/com/knowledge/base/domain/doc/service/FileImporterDispatcher.java new file mode 100644 index 0000000..13318b4 --- /dev/null +++ b/src/main/java/com/knowledge/base/domain/doc/service/FileImporterDispatcher.java @@ -0,0 +1,45 @@ +package com.knowledge.base.domain.doc.service; + +import com.knowledge.base.domain.doc.service.impl.AbstractBaseFileImporter; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Component; + +import javax.annotation.PostConstruct; +import java.nio.file.Path; +import java.util.*; + +/** + * 文件导入分派器,可根据文件后缀名的不同,执行不同的导入逻辑 + */ +@Component +public class FileImporterDispatcher { + + @Autowired + private List fileImporters; + + private final Map suffixToImporter = new HashMap<>(); + + @PostConstruct + public void init() { + for (AbstractBaseFileImporter importer : fileImporters) { + for (String suffix : importer.getFileSuffixes()) { + suffixToImporter.put(suffix.toLowerCase(), importer); + } + } + } + + public boolean importSingleFile(Path filePath, Path excludePrefix, Map extInfo) { + String fileName = filePath.getFileName().toString().toLowerCase(); + Optional matchedSuffix = suffixToImporter.keySet().stream() + .filter(fileName::endsWith) + .findFirst(); + + if (matchedSuffix.isEmpty()) { + throw new IllegalArgumentException("不支持的文件后缀: " + fileName); + } + + AbstractBaseFileImporter importer = suffixToImporter.get(matchedSuffix.get()); + return importer.insertOrUpdateOneFileIntoES(filePath, excludePrefix, extInfo); + } +} + diff --git a/src/main/java/com/knowledge/base/domain/doc/service/iface/DocumentImporter.java b/src/main/java/com/knowledge/base/domain/doc/service/iface/DocumentImporter.java index 4ad35e6..684acac 100644 --- a/src/main/java/com/knowledge/base/domain/doc/service/iface/DocumentImporter.java +++ b/src/main/java/com/knowledge/base/domain/doc/service/iface/DocumentImporter.java @@ -2,6 +2,9 @@ package com.knowledge.base.domain.doc.service.iface; import com.knowledge.base.domain.common.enums.DocTypeEnum; +import java.nio.file.Path; +import java.util.Map; + /** * @author Luke.ye * @date 2025/5/20 09:02 @@ -14,4 +17,12 @@ public interface DocumentImporter { * @return */ String getType(); + + /** + * 将文件插入/更新到ES中 + * @param absoluteFilePath + * @param excludeFilePrefix + * @return + */ + default boolean insertOrUpdateOneFileIntoES(Path absoluteFilePath, Path excludeFilePrefix, Map extInfo) { return true; } } diff --git a/src/main/java/com/knowledge/base/domain/doc/service/impl/AbstractBaseFileImporter.java b/src/main/java/com/knowledge/base/domain/doc/service/impl/AbstractBaseFileImporter.java index 0231a4d..92f6d4b 100644 --- a/src/main/java/com/knowledge/base/domain/doc/service/impl/AbstractBaseFileImporter.java +++ b/src/main/java/com/knowledge/base/domain/doc/service/impl/AbstractBaseFileImporter.java @@ -1,8 +1,11 @@ package com.knowledge.base.domain.doc.service.impl; +import cn.hutool.core.map.MapUtil; import cn.hutool.core.util.StrUtil; import cn.hutool.json.JSON; import cn.hutool.json.JSONUtil; +import com.google.common.collect.Lists; +import com.google.common.collect.Maps; import com.knowledge.base.domain.common.enums.DocMetaPropEnum; import com.knowledge.base.domain.doc.model.OSRecordDO; import com.knowledge.base.domain.doc.service.iface.DocumentImporter; @@ -10,10 +13,7 @@ import com.knowledge.base.domain.doc.service.iface.FileDomainService; import com.knowledge.base.infrastructure.cache.iface.FileCacheService; import com.knowledge.base.infrastructure.config.ConstantConfig; import com.knowledge.base.infrastructure.config.ThreadPoolConfig; -import com.knowledge.base.infrastructure.util.CacheUtil; -import com.knowledge.base.infrastructure.util.DateUtil; -import com.knowledge.base.infrastructure.util.SafeIdUtil; -import com.knowledge.base.infrastructure.util.ThreadPoolUtil; +import com.knowledge.base.infrastructure.util.*; import org.elasticsearch.action.delete.DeleteRequest; import org.elasticsearch.action.get.GetRequest; import org.elasticsearch.action.get.GetResponse; @@ -29,6 +29,7 @@ import java.nio.file.*; import java.util.HashMap; import java.util.Map; import java.util.Optional; +import java.util.Set; public abstract class AbstractBaseFileImporter implements DocumentImporter { @@ -44,9 +45,12 @@ public abstract class AbstractBaseFileImporter implements DocumentImporter { @Autowired private FileDomainService fileDomainService;; + @Autowired + private RateLimiterManager rateLimiterManager; + protected abstract String getDirectoryPath(); - protected abstract String getFileSuffix(); + public abstract Set getFileSuffixes(); protected abstract String getDocTypeCode(); @@ -58,50 +62,72 @@ public abstract class AbstractBaseFileImporter implements DocumentImporter { Path excludeBase = Paths.get(getExcludePrefix()).toAbsolutePath().normalize(); Files.walk(basePath) - .filter(p -> p.toString().toLowerCase().endsWith(getFileSuffix())) + .filter(Files::isRegularFile) + .filter(path -> { + String fileName = path.getFileName().toString().toLowerCase(); + return getFileSuffixes().stream().anyMatch(fileName::endsWith); + }) .forEach(path -> { - ThreadPoolUtil.execute(() -> processFile(path, excludeBase), ThreadPoolConfig.IMPORT_DOC_POOL); - try { - Thread.sleep(100); - } catch (InterruptedException e) { - throw new RuntimeException(e); - } + rateLimiterManager.getRateLimiter().acquire(); + ThreadPoolUtil.execute(() -> insertOrUpdateOneFileIntoES(path, excludeBase, Maps.newHashMap()), ThreadPoolConfig.IMPORT_DOC_POOL); }); ThreadPoolUtil.shutdownAndAwait(); } - private void processFile(Path path, Path excludeBase) { + /** + * 将单个文件导入或更新到 Elasticsearch 中,并缓存元信息至 Redis。 + * + * 方法逻辑流程如下: + * 1. 获取文件相对路径及最后修改时间; + * 2. 检查 Redis 中的缓存元信息是否存在且未过期; + * 3. 检查 ES 中是否已有相同文档,且未变动; + * 4. 若有变化或首次导入,则提取文件内容、构建文档; + * 5. 将文档写入 ES,并缓存元信息; + * 6. 支持通过 extInfo 参数手动提供 uploader、url、expireTime 信息(用于无缓存情况)。 + * + * @param absoluteFilePath 文件绝对路径 + * @param excludeFilePrefix 排除前缀(用于计算相对路径) + * @param extInfo 可选附加元信息(当缓存未命中时使用) + * @return 导入是否成功 + */ + @Override + public boolean insertOrUpdateOneFileIntoES(Path absoluteFilePath, Path excludeFilePrefix, Map extInfo) { try { - Path absPath = path.toAbsolutePath().normalize(); - Path relativePathObj; - if (absPath.startsWith(excludeBase)) { - relativePathObj = excludeBase.relativize(absPath); - } else { + // === Step 1: 计算相对路径 & 获取文件信息 === + Path absPath = absoluteFilePath.toAbsolutePath().normalize(); + Path relativePathObj = absPath.startsWith(excludeFilePrefix) + ? excludeFilePrefix.relativize(absPath) + : absPath; + if (!absPath.startsWith(excludeFilePrefix)) { logger.warn("路径未匹配 exclude.prefix,使用全路径: {}", absPath); - relativePathObj = absPath; } - Long localMTime = Files.getLastModifiedTime(path).toMillis(); + Long localMTime = Files.getLastModifiedTime(absoluteFilePath).toMillis(); String localRelaFilePath = relativePathObj.toString().replace("\\", "/"); - String fileNameWithSuffix = path.getFileName().toString(); + String fileNameWithSuffix = absoluteFilePath.getFileName().toString(); + + // === Step 2: 检查 Redis 缓存是否已是最新 === Optional metaJsonOpt = fileCacheService.getMeta(localRelaFilePath); - if(metaJsonOpt.isPresent()) { - // 有文件信息 + if (metaJsonOpt.isPresent()) { JSON metaJson = JSONUtil.parse(metaJsonOpt.get()); Long uploadTime = metaJson.getByPath(DocMetaPropEnum.UPLOAD_TIME.code, Long.class); - if(localMTime.equals(uploadTime)) { + if (localMTime.equals(uploadTime)) { logger.info("文件未变动,跳过导入: {}", localRelaFilePath); - return; + return true; } + fileCacheService.removeMetaCache(localRelaFilePath); + logger.info("[Redis] 已移除旧版本文件缓存Meta信息: {}", localRelaFilePath); } - String content = extractContent(path); + // === Step 3: 提取文件内容 === + String content = extractContent(absoluteFilePath); if (StrUtil.isBlank(content)) { - logger.warn("跳过空内容文件: {}", path); - return; + logger.warn("跳过空内容文件: {}", absoluteFilePath); + return true; } + // === Step 4: 检查 ES 是否已有未变动版本 === String docId = SafeIdUtil.encode(localRelaFilePath); GetRequest getRequest = new GetRequest(INDEX_NAME, docId); if (esClient.exists(getRequest, RequestOptions.DEFAULT)) { @@ -110,73 +136,93 @@ public abstract class AbstractBaseFileImporter implements DocumentImporter { Object esMtime = existingSource.get("mtime"); if (esMtime != null && Long.parseLong(esMtime.toString()) == localMTime) { logger.info("文件未变动,跳过导入: {}", localRelaFilePath); - String metaJsonStr = transToMetaJson(existingSource); - fileCacheService.cacheMeta(localRelaFilePath, metaJsonStr, ConstantConfig.FILE_META_CACHE_EXPIRED_MINUTES); - return; + fileCacheService.cacheMeta(localRelaFilePath, transToMetaJson(existingSource), ConstantConfig.FILE_META_CACHE_EXPIRED_MINUTES); + return true; } - DeleteRequest deleteRequest = new DeleteRequest(INDEX_NAME, docId); - esClient.delete(deleteRequest, RequestOptions.DEFAULT); - logger.info("已删除旧版本文件: {}", localRelaFilePath); + esClient.delete(new DeleteRequest(INDEX_NAME, docId), RequestOptions.DEFAULT); + logger.info("[ES] 已删除旧版本文件: {}", localRelaFilePath); } logger.info("[开始进行文件导入......] localRelaFilePath: {}, localMTime: {}", localRelaFilePath, localMTime); - Map doc = new HashMap<>(); - doc.put("filename", fileNameWithSuffix); - doc.put("filepath", localRelaFilePath); - doc.put("content", content); - doc.put("mtime", localMTime); - if(metaJsonOpt.isPresent()) { - Map metaMap = JSONUtil.toBean(metaJsonOpt.get(), Map.class); - doc.put("uploader", CacheUtil.getFileMetaProp(metaMap, DocMetaPropEnum.UPLOADER.code, String.class, ConstantConfig.DEFAULT_UPLOADER)); - String cacheFileUrl = CacheUtil.getFileMetaProp(metaMap, DocMetaPropEnum.ACCESS_URL.code, String.class, StrUtil.EMPTY); - doc.put("url", buildAccessUrl(cacheFileUrl, localRelaFilePath, fileNameWithSuffix)); - doc.put("expireTime", CacheUtil.getFileMetaProp(metaMap, DocMetaPropEnum.EXPIRE_TIME.code, - Long.class, DateUtil.toMillis(ConstantConfig.LONG_TERM_EXPIRE_TIME))); - } else { - doc.put("uploader",ConstantConfig.DEFAULT_UPLOADER); - String url = buildAccessUrl(StrUtil.EMPTY, localRelaFilePath, fileNameWithSuffix); - doc.put("url", url); - doc.put("expireTime", DateUtil.toMillis(ConstantConfig.LONG_TERM_EXPIRE_TIME)); - } - - IndexRequest request = new IndexRequest(INDEX_NAME) - .id(docId) - .source(doc); + // === Step 5: 构建待写入文档 === + Map doc = buildDocument(fileNameWithSuffix, localRelaFilePath, content, localMTime, extInfo); + // === Step 6: 写入 Elasticsearch 并更新缓存 === + IndexRequest request = new IndexRequest(INDEX_NAME).id(docId).source(doc); esClient.index(request, RequestOptions.DEFAULT); logger.info("导入成功: {}", localRelaFilePath); - // 信息入缓存 fileCacheService.cacheMeta(localRelaFilePath, transToMetaJson(doc), ConstantConfig.FILE_META_CACHE_EXPIRED_MINUTES); + return true; + } catch (Exception e) { - logger.error("导入失败: {}", path, e); + logger.error("导入失败: {}", absoluteFilePath, e); + return false; } } - private String buildAccessUrl(String originUrl, String localRelaFilePath, String fileNameWithSuffix) { - // 优先使用原始链接 - if(StrUtil.isNotBlank(originUrl)) { - return originUrl; + private Map buildDocument(String filename, String filepath, String content, Long mtime, Map extInfo) { + Map doc = new HashMap<>(); + doc.put("filename", filename); + doc.put("filepath", filepath); + doc.put("content", content); + doc.put("mtime", mtime); + + Optional metaJsonOpt = fileCacheService.getMeta(filepath); + if (metaJsonOpt.isPresent()) { + Map metaMap = JSONUtil.toBean(metaJsonOpt.get(), Map.class); + doc.put("uploader", CacheUtil.getFileMetaProp(metaMap, DocMetaPropEnum.UPLOADER.code, String.class, ConstantConfig.DEFAULT_UPLOADER)); + doc.put("url", CacheUtil.getFileMetaProp(metaMap, DocMetaPropEnum.ACCESS_URL.code, String.class, StrUtil.EMPTY)); + doc.put("expireTime", CacheUtil.getFileMetaProp(metaMap, DocMetaPropEnum.EXPIRE_TIME.code, Long.class, DateUtil.toMillis(ConstantConfig.LONG_TERM_EXPIRE_TIME))); + } else { + extInfo = MapUtil.isEmpty(extInfo) ? MapUtil.empty() : extInfo; + boolean hasAllMeta = extInfo.keySet().containsAll(Lists.newArrayList( + DocMetaPropEnum.UPLOADER.code, DocMetaPropEnum.ACCESS_URL.code, DocMetaPropEnum.EXPIRE_TIME.code)); + if (hasAllMeta) { + doc.put("uploader", extInfo.get(DocMetaPropEnum.UPLOADER.code)); + doc.put("url", extInfo.get(DocMetaPropEnum.ACCESS_URL.code)); + doc.put("expireTime", extInfo.get(DocMetaPropEnum.EXPIRE_TIME.code)); + } else { + Map props = buildDocMetaProps(filepath); + doc.put("uploader", props.get(DocMetaPropEnum.UPLOADER.code)); + doc.put("url", props.get(DocMetaPropEnum.ACCESS_URL.code)); + doc.put("expireTime", props.get(DocMetaPropEnum.EXPIRE_TIME.code)); + } } + return doc; + } + + + + private Map buildDocMetaProps(String localRelaFilePath) { + Map metaMap = new HashMap<>(Map.of( + DocMetaPropEnum.UPLOADER.code, ConstantConfig.DEFAULT_UPLOADER, + DocMetaPropEnum.ACCESS_URL.code, StrUtil.EMPTY, + DocMetaPropEnum.EXPIRE_TIME.code, DateUtil.toMillis(ConstantConfig.LONG_TERM_EXPIRE_TIME) + )); if(StrUtil.isBlank(localRelaFilePath)) { - return StrUtil.EMPTY; + return metaMap; } - String accessUrl = originUrl; + String accessUrl = localRelaFilePath; if(localRelaFilePath.endsWith(".md")) { // markdown直接拼接http链接 accessUrl = String.format("%s/%s", ConstantConfig.SHARE_BASE_URL, localRelaFilePath); accessUrl = accessUrl.substring(0, accessUrl.length() - 3) + ".html"; + metaMap.put(DocMetaPropEnum.ACCESS_URL.code, accessUrl); } else { // 其它文件从对象存储的DB中获取 Optional latestRecordOpt = fileDomainService.getLatestRecordByRelaPath(localRelaFilePath); if(latestRecordOpt.isPresent()) { OSRecordDO latestRecord = latestRecordOpt.get(); accessUrl = latestRecord.getUrl(); + metaMap.put(DocMetaPropEnum.ACCESS_URL.code, accessUrl); + metaMap.put(DocMetaPropEnum.UPLOADER.code, latestRecord.getUploader()); + metaMap.put(DocMetaPropEnum.EXPIRE_TIME.code, latestRecord.getExpireTime()); } } - return accessUrl; + return metaMap; } private String transToMetaJson(Map esExistingSource) { diff --git a/src/main/java/com/knowledge/base/domain/doc/service/impl/ExcelImporter.java b/src/main/java/com/knowledge/base/domain/doc/service/impl/ExcelImporter.java index 6feedf0..c84bb80 100644 --- a/src/main/java/com/knowledge/base/domain/doc/service/impl/ExcelImporter.java +++ b/src/main/java/com/knowledge/base/domain/doc/service/impl/ExcelImporter.java @@ -13,6 +13,7 @@ import java.io.FileInputStream; import java.io.IOException; import java.nio.file.Path; import java.util.Iterator; +import java.util.Set; @Component public class ExcelImporter extends AbstractBaseFileImporter { @@ -26,8 +27,8 @@ public class ExcelImporter extends AbstractBaseFileImporter { } @Override - protected String getFileSuffix() { - return ".xlsx"; + public Set getFileSuffixes() { + return Set.of(".xlsx", ".xls"); } @Override diff --git a/src/main/java/com/knowledge/base/domain/doc/service/impl/MarkdownImporter.java b/src/main/java/com/knowledge/base/domain/doc/service/impl/MarkdownImporter.java index 6052538..ba5775f 100644 --- a/src/main/java/com/knowledge/base/domain/doc/service/impl/MarkdownImporter.java +++ b/src/main/java/com/knowledge/base/domain/doc/service/impl/MarkdownImporter.java @@ -7,6 +7,7 @@ import org.springframework.stereotype.Component; import java.io.IOException; import java.nio.charset.StandardCharsets; import java.nio.file.*; +import java.util.Set; /** * @author Luke.ye @@ -27,8 +28,8 @@ public class MarkdownImporter extends AbstractBaseFileImporter { } @Override - protected String getFileSuffix() { - return ".md"; + public Set getFileSuffixes() { + return Set.of(".md"); } @Override diff --git a/src/main/java/com/knowledge/base/domain/doc/service/impl/PdfImporter.java b/src/main/java/com/knowledge/base/domain/doc/service/impl/PdfImporter.java index 482a423..2f58e35 100644 --- a/src/main/java/com/knowledge/base/domain/doc/service/impl/PdfImporter.java +++ b/src/main/java/com/knowledge/base/domain/doc/service/impl/PdfImporter.java @@ -9,6 +9,7 @@ import org.springframework.beans.factory.annotation.Value; import java.io.IOException; import java.nio.file.*; +import java.util.Set; /** * @author Luke.ye @@ -26,8 +27,8 @@ public class PdfImporter extends AbstractBaseFileImporter { } @Override - protected String getFileSuffix() { - return ".pdf"; + public Set getFileSuffixes() { + return Set.of(".pdf"); } @Override diff --git a/src/main/java/com/knowledge/base/domain/doc/service/impl/WordImporter.java b/src/main/java/com/knowledge/base/domain/doc/service/impl/WordImporter.java index fcc6555..37ee013 100644 --- a/src/main/java/com/knowledge/base/domain/doc/service/impl/WordImporter.java +++ b/src/main/java/com/knowledge/base/domain/doc/service/impl/WordImporter.java @@ -9,6 +9,7 @@ import org.springframework.beans.factory.annotation.Value; import java.io.FileInputStream; import java.io.IOException; import java.nio.file.*; +import java.util.Set; /** @@ -27,8 +28,8 @@ public class WordImporter extends AbstractBaseFileImporter { } @Override - protected String getFileSuffix() { - return ".docx"; + public Set getFileSuffixes() { + return Set.of(".doc", ".docx"); } @Override diff --git a/src/main/java/com/knowledge/base/infrastructure/cache/iface/FileCacheService.java b/src/main/java/com/knowledge/base/infrastructure/cache/iface/FileCacheService.java index 338401c..9ec0480 100644 --- a/src/main/java/com/knowledge/base/infrastructure/cache/iface/FileCacheService.java +++ b/src/main/java/com/knowledge/base/infrastructure/cache/iface/FileCacheService.java @@ -25,6 +25,12 @@ public interface FileCacheService { */ default Optional getMeta(String fnWithRelativePath) { return null; }; + /** + * 移除指定key对应对的缓存 + * @param fnWithRelativePath + */ + default void removeMetaCache(String fnWithRelativePath) {}; + /** * 清空所有缓存(注意:某些实现可能未实现) */ diff --git a/src/main/java/com/knowledge/base/infrastructure/cache/impl/RedisFileCacheServiceImpl.java b/src/main/java/com/knowledge/base/infrastructure/cache/impl/RedisFileCacheServiceImpl.java index 70a6a09..c8e1ab9 100644 --- a/src/main/java/com/knowledge/base/infrastructure/cache/impl/RedisFileCacheServiceImpl.java +++ b/src/main/java/com/knowledge/base/infrastructure/cache/impl/RedisFileCacheServiceImpl.java @@ -43,6 +43,11 @@ public class RedisFileCacheServiceImpl implements FileCacheService { return Optional.ofNullable(value); } + @Override + public void removeMetaCache(String fnWithRelativePath) { + redisTemplate.delete(key("meta", fnWithRelativePath)); + } + @Override public void clearAll() { log.warn("正在清空 Redis 文件缓存,前缀: {}", FILE_CACHE_PREFIX); diff --git a/src/main/java/com/knowledge/base/infrastructure/config/ConstantConfig.java b/src/main/java/com/knowledge/base/infrastructure/config/ConstantConfig.java index d12d6d2..5857221 100644 --- a/src/main/java/com/knowledge/base/infrastructure/config/ConstantConfig.java +++ b/src/main/java/com/knowledge/base/infrastructure/config/ConstantConfig.java @@ -19,7 +19,7 @@ public class ConstantConfig { public static final LocalDateTime LONG_TERM_EXPIRE_TIME = LocalDateTime.of(9999, 12, 31, 23, 59, 59); - public static final long FILE_META_CACHE_EXPIRED_MINUTES = Duration.ofDays(2).toMinutes(); + public static final long FILE_META_CACHE_EXPIRED_MINUTES = Duration.ofDays(7).toMinutes(); public static final long USER_CACHE_EXPIRED_MINUTES = Duration.ofDays(10).toMinutes(); diff --git a/src/main/java/com/knowledge/base/infrastructure/config/DynamicConfig.java b/src/main/java/com/knowledge/base/infrastructure/config/DynamicConfig.java index 45742d1..dec22bb 100644 --- a/src/main/java/com/knowledge/base/infrastructure/config/DynamicConfig.java +++ b/src/main/java/com/knowledge/base/infrastructure/config/DynamicConfig.java @@ -1,15 +1,19 @@ package com.knowledge.base.infrastructure.config; +import lombok.Getter; import org.springframework.beans.factory.annotation.Value; import org.springframework.cloud.context.config.annotation.RefreshScope; import org.springframework.stereotype.Component; +import java.util.List; + /** * @author Luke.ye * @date 2025/5/8 09:59 */ @Component @RefreshScope +@Getter public class DynamicConfig { //是否记录请求和响应信息 Y-记录 N-不记录 默认记录 @Value("${micro.saas.doc.parser.recordMsgBody:Y}") @@ -25,19 +29,9 @@ public class DynamicConfig { @Value("${token.expire.time:86400000}") private long tokenExpireTime; - public String getRecordMsgBody() { - return recordMsgBody; - } + @Value("${os.supported.searchable.file.suffix: pdf,doc,docx,xls,xlsx}") + private String supportedSearchFileSuffix; - public String getImportScheduleCron() { - return importScheduleCron; - } - - public String getCookieDomainName() { - return cookieDomainName; - } - - public long getTokenExpireTime() { - return tokenExpireTime; - } + @Value("${file.import.rate.limit: 10}") + private String fileImportRateLimit; } diff --git a/src/main/java/com/knowledge/base/infrastructure/config/ObjectStorageProperties.java b/src/main/java/com/knowledge/base/infrastructure/config/ObjectStorageProperties.java index 63fe5fc..81a66b2 100644 --- a/src/main/java/com/knowledge/base/infrastructure/config/ObjectStorageProperties.java +++ b/src/main/java/com/knowledge/base/infrastructure/config/ObjectStorageProperties.java @@ -20,5 +20,5 @@ public class ObjectStorageProperties { /** * 可被检索文件的本地路径 */ - private String localSearchablePath; + private String localSearchablePathPrefix; } \ No newline at end of file diff --git a/src/main/java/com/knowledge/base/infrastructure/north/controller/FileWriteController.java b/src/main/java/com/knowledge/base/infrastructure/north/controller/FileWriteController.java index 363943c..9c6aebb 100644 --- a/src/main/java/com/knowledge/base/infrastructure/north/controller/FileWriteController.java +++ b/src/main/java/com/knowledge/base/infrastructure/north/controller/FileWriteController.java @@ -7,13 +7,14 @@ import cn.hutool.json.JSONUtil; import com.knowledge.base.application.service.DocAppService; import com.knowledge.base.application.service.UserAppService; import com.knowledge.base.domain.common.enums.DocMetaPropEnum; +import com.knowledge.base.domain.doc.service.FileImporterDispatcher; import com.knowledge.base.domain.doc.service.iface.DocumentImporter; import com.knowledge.base.infrastructure.cache.iface.FileCacheService; import com.knowledge.base.infrastructure.config.ConstantConfig; +import com.knowledge.base.infrastructure.config.DynamicConfig; import com.knowledge.base.infrastructure.config.ObjectStorageProperties; import com.knowledge.base.infrastructure.north.dto.doc.OSRecordDTO; import com.knowledge.base.infrastructure.north.dto.user.UserDTO; -import com.knowledge.base.infrastructure.north.dto.user.UserTokenDTO; import com.knowledge.base.infrastructure.south.ObjectStorageGateway; import com.knowledge.base.infrastructure.util.DateUtil; import com.knowledge.base.infrastructure.util.ThreadPoolUtil; @@ -35,6 +36,7 @@ import java.nio.file.Path; import java.nio.file.Paths; import java.time.LocalDateTime; import java.util.*; +import java.util.stream.Collectors; @RestController @RequestMapping("/api/v1/doc") @@ -47,12 +49,16 @@ public class FileWriteController { private final FileCacheService fileCacheService; private static final String INDEX_NAME = "documents"; + private final DynamicConfig dynamicConfig; + private final UserAppService userAppService; private final DocAppService docAppService; private final ObjectStorageGateway objectStorageGateway; private final ObjectStorageProperties objectStorageProperties; + private final FileImporterDispatcher dispatcher; + @Autowired private List importers; @@ -173,7 +179,8 @@ public class FileWriteController { String uploader = userDTO.get().getUsername(); LocalDateTime now = LocalDateTime.now(); String dateFolder = now.toLocalDate().toString(); - Set supportedTypes = Set.of("pdf", "doc", "docx", "xls", "xlsx"); + String[] supportedSuffix = dynamicConfig.getSupportedSearchFileSuffix().split(","); + Set supportedTypes = Arrays.stream(supportedSuffix).map(String::trim).collect(Collectors.toSet()); boolean allowSaveToLocal = searchable && supportedTypes.contains(suffix); String localRelaFilePath = String.format("%s/%s", dateFolder, originFileNameWithSuffix); @@ -198,7 +205,8 @@ public class FileWriteController { // 允许保存到本地 if (allowSaveToLocal) { - String localDir = Paths.get(objectStorageProperties.getLocalSearchablePath(), dateFolder).toString(); + String excludePrefix = objectStorageProperties.getLocalSearchablePathPrefix(); + String localDir = Paths.get(excludePrefix, dateFolder).toString(); File localTargetDir = new File(localDir); if (!localTargetDir.exists()) { localTargetDir.mkdirs(); @@ -207,18 +215,15 @@ public class FileWriteController { Path targetPath = Paths.get(localDir, originFileNameWithSuffix); Files.write(targetPath, fileBytes); - long mtime = Files.getLastModifiedTime(targetPath).toMillis(); + // 直接导入ES + Map extInfo = Map.of( + DocMetaPropEnum.UPLOADER.code, uploader, + DocMetaPropEnum.ACCESS_URL.code, url, + DocMetaPropEnum.EXPIRE_TIME.code, DateUtil.toMillis(expireTime) + ); + boolean res = dispatcher.importSingleFile(targetPath.toAbsolutePath(), Paths.get(excludePrefix), extInfo); - Map cacheMeta = new LinkedHashMap<>(); - cacheMeta.put(DocMetaPropEnum.UPLOADER.code, uploader); - cacheMeta.put(DocMetaPropEnum.ACCESS_URL.code, url); - cacheMeta.put(DocMetaPropEnum.UPLOAD_TIME.code, mtime); - cacheMeta.put(DocMetaPropEnum.EXPIRE_TIME.code, DateUtil.toMillis(expireTime)); - - String metaJson = JSONUtil.toJsonStr(cacheMeta); - fileCacheService.cacheMeta(localRelaFilePath, metaJson, ConstantConfig.FILE_META_CACHE_EXPIRED_MINUTES); - - logger.info("本地文件信息已保存并写入缓存: key = {}, value = {}", localRelaFilePath, metaJson); + logger.info("本地文件信息已保存并导入ES: localRelaPath = {}, ESImportRes = {}", localRelaFilePath, res); } } catch (Exception e) { logger.error("异步写入失败: {}", originFileNameWithSuffix, e); diff --git a/src/main/java/com/knowledge/base/infrastructure/util/RateLimiterManager.java b/src/main/java/com/knowledge/base/infrastructure/util/RateLimiterManager.java new file mode 100644 index 0000000..5ebc672 --- /dev/null +++ b/src/main/java/com/knowledge/base/infrastructure/util/RateLimiterManager.java @@ -0,0 +1,32 @@ +package com.knowledge.base.infrastructure.util; + +import com.google.common.util.concurrent.RateLimiter; +import com.knowledge.base.infrastructure.config.DynamicConfig; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Component; + +@Component +public class RateLimiterManager { + + @Autowired + private DynamicConfig dynamicConfig; + + private volatile RateLimiter rateLimiter; + private volatile double lastRate = -1; + + public RateLimiter getRateLimiter() { + double currentRate = Double.valueOf(dynamicConfig.getFileImportRateLimit()); + + // 如果速率发生变化,则更新限速器 + if (rateLimiter == null || currentRate != lastRate) { + synchronized (this) { + if (rateLimiter == null || currentRate != lastRate) { + rateLimiter = RateLimiter.create(currentRate); + lastRate = currentRate; + } + } + } + + return rateLimiter; + } +} diff --git a/src/main/resources/application-dev-mac.properties b/src/main/resources/application-dev-mac.properties index 965fb6a..a30a2b4 100644 --- a/src/main/resources/application-dev-mac.properties +++ b/src/main/resources/application-dev-mac.properties @@ -27,7 +27,7 @@ markdown.path=/Users/admin/Desktop/Archived/micro-saas pdf.path=/Users/admin/Desktop/Archived/micro-saas word.path=/Users/admin/Desktop/Archived/micro-saas excel.path=/Users/admin/Desktop/Archived/micro-saas -object.storage.local-searchable-path=/Users/admin/Desktop/Archived/micro-saas +object.storage.local-searchable-path-prefix=/Users/admin/Desktop/Archived/micro-saas # mysql spring.datasource.url=jdbc:mysql://localhost:3306/kbase?useUnicode=true&characterEncoding=UTF-8&serverTimezone=Asia/Shanghai diff --git a/src/main/resources/application-dev-windows.properties b/src/main/resources/application-dev-windows.properties index 02c8edf..8f21468 100644 --- a/src/main/resources/application-dev-windows.properties +++ b/src/main/resources/application-dev-windows.properties @@ -27,7 +27,7 @@ markdown.path=D:/02-documents/01-ahnx-share-src-public/public pdf.path=D:/02-documents/01-os-uploaded-searchable-files word.path=D:/02-documents/01-os-uploaded-searchable-files excel.path=D:/02-documents/01-os-uploaded-searchable-files -object.storage.local-searchable-path=D:/02-documents/01-os-uploaded-searchable-files +object.storage.local-searchable-path-prefix=D:/02-documents/01-os-uploaded-searchable-files # mysql diff --git a/src/main/resources/application-docker.properties b/src/main/resources/application-docker.properties index 1f40a91..a576f66 100644 --- a/src/main/resources/application-docker.properties +++ b/src/main/resources/application-docker.properties @@ -26,7 +26,7 @@ markdown.path=/app/import-data pdf.path=/app/os-uploaded-searchable-files word.path=/app/os-uploaded-searchable-files excel.path=/app/os-uploaded-searchable-files -object.storage.local-searchable-path=/app/os-uploaded-searchable-files +object.storage.local-searchable-path-prefix=/app/os-uploaded-searchable-files # 相对路径裁剪 exclude.file.path.prefix=/app/import-data