集成redis

This commit is contained in:
luke 2025-06-07 01:53:50 +08:00
parent 4be29b3d85
commit 915fcd61ab
11 changed files with 440 additions and 83 deletions

View File

@ -68,6 +68,11 @@
<scope>test</scope>
<version>${spring-boot.version}</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-redis</artifactId>
<version>${spring-boot.version}</version>
</dependency>
<!-- #日志依赖,参考 https://cdn.jsdelivr.net/gh/lukeyewu/pictures@main/uPic/202411/YhMx1M.png -->
<!-- 日志_抽象API层 -->
@ -247,6 +252,10 @@
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter</artifactId>

View File

@ -1,7 +1,9 @@
package com.knowledge.base.domain.doc.impl;
import com.knowledge.base.domain.doc.iface.DocumentImporter;
import com.knowledge.base.infrastructure.util.LocalCacheUtil;
import com.knowledge.base.infrastructure.cache.FileCacheService;
import com.knowledge.base.infrastructure.util.SafeIdUtil;
import com.knowledge.base.infrastructure.util.ThreadPoolUtil;
import org.elasticsearch.action.delete.DeleteRequest;
import org.elasticsearch.action.get.GetRequest;
import org.elasticsearch.action.get.GetResponse;
@ -17,6 +19,7 @@ import java.io.IOException;
import java.nio.file.*;
import java.util.HashMap;
import java.util.Map;
import java.util.Optional;
public abstract class AbstractBaseFileImporter implements DocumentImporter {
@ -29,6 +32,9 @@ public abstract class AbstractBaseFileImporter implements DocumentImporter {
@Value("${exclude.file.path.prefix}")
private String excludePrefix;
@Autowired
private FileCacheService fileCacheService;
protected abstract String getDirectoryPath();
protected abstract String getFileSuffix();
@ -42,89 +48,82 @@ public abstract class AbstractBaseFileImporter implements DocumentImporter {
Files.walk(basePath)
.filter(p -> p.toString().toLowerCase().endsWith(getFileSuffix()))
.forEach(path -> {
try {
Path absPath = path.toAbsolutePath().normalize();
Path relativePathObj;
if (absPath.startsWith(excludeBase)) {
relativePathObj = excludeBase.relativize(absPath);
} else {
logger.warn("路径未匹配 exclude.prefix使用全路径: {}", absPath);
relativePathObj = absPath;
}
.forEach(path -> ThreadPoolUtil.execute(() -> processFile(path, excludeBase)));
String relativePath = relativePathObj.toString().replace("\\", "/");
String fileName = path.getFileName().toString();
if(LocalCacheUtil.get(fileName) == null){
// 将文件的路径缓存下来
LocalCacheUtil.put(fileName, relativePath);
}
Long localMTime = Files.getLastModifiedTime(path).toMillis();
logger.info("absPath: {}", absPath.toString().replace("\\", "/"));
logger.info("rerelativePath: {} localMTime {}", relativePath, localMTime);
// 优先从缓存中加载数据判断是否文件有修改
Long lastModifiedTime = (Long)LocalCacheUtil.get(relativePath);
if(localMTime.equals(lastModifiedTime)) {
logger.info("文件未变动,跳过导入: {}", relativePath);
return;
}
String content = extractContent(path);
if (content == null || content.isBlank()) {
logger.warn("跳过空内容文件: {}", path);
return;
}
// relativePath 作为文档 ID
GetRequest getRequest = new GetRequest(INDEX_NAME, relativePath);
boolean fileHasChanged = false;
if (esClient.exists(getRequest, RequestOptions.DEFAULT)) {
GetResponse existing = esClient.get(getRequest, RequestOptions.DEFAULT);
Map<String, Object> existingSource = existing.getSourceAsMap();
Object esMtime = existingSource.get("mtime");
if (esMtime != null && Long.parseLong(esMtime.toString()) == localMTime) {
logger.info("文件未变动,跳过导入: {}", relativePath);
if(LocalCacheUtil.get(relativePath) == null) {
// 缓存过期
LocalCacheUtil.put(relativePath, localMTime);
}
return;
} else {
fileHasChanged = true;
}
}
if (fileHasChanged) {
// 删除旧版本
DeleteRequest deleteRequest = new DeleteRequest(INDEX_NAME, relativePath);
esClient.delete(deleteRequest, RequestOptions.DEFAULT);
logger.info("已删除旧版本文件: {}", relativePath);
}
// === 构建文档 ===
Map<String, Object> doc = new HashMap<>();
doc.put("filename", fileName);
doc.put("filepath", relativePath);
doc.put("content", content);
doc.put("mtime", localMTime);
IndexRequest request = new IndexRequest(INDEX_NAME)
.id(relativePath)
.source(doc);
esClient.index(request, RequestOptions.DEFAULT);
logger.info("导入成功: {}", relativePath);
// 信息入缓存
LocalCacheUtil.put(relativePath, localMTime);
} catch (Exception e) {
logger.error("导入失败: {}", path, e);
}
});
ThreadPoolUtil.shutdownAndAwait(); // 等待任务完成
}
private void processFile(Path path, Path excludeBase) {
try {
Path absPath = path.toAbsolutePath().normalize();
Path relativePathObj;
if (absPath.startsWith(excludeBase)) {
relativePathObj = excludeBase.relativize(absPath);
} else {
logger.warn("路径未匹配 exclude.prefix使用全路径: {}", absPath);
relativePathObj = absPath;
}
String relativePath = relativePathObj.toString().replace("\\", "/");
String fileName = path.getFileName().toString();
Optional<String> pathOpt = fileCacheService.getPath(fileName);
if (pathOpt.isEmpty()) {
fileCacheService.cachePath(fileName, relativePath, 24 * 60);
}
Long localMTime = Files.getLastModifiedTime(path).toMillis();
logger.info("relativePath: {} localMTime {}", relativePath, localMTime);
// 优先从缓存中加载数据判断是否文件有修改
Optional<Long> lastModifiedTime = fileCacheService.getMTime(relativePath);
if(lastModifiedTime.isPresent() && localMTime.equals(lastModifiedTime.get())) {
logger.info("文件未变动,跳过导入: {}", relativePath);
return;
}
String content = extractContent(path);
if (content == null || content.isBlank()) {
logger.warn("跳过空内容文件: {}", path);
return;
}
String docId = SafeIdUtil.encode(relativePath);
GetRequest getRequest = new GetRequest(INDEX_NAME, docId);
if (esClient.exists(getRequest, RequestOptions.DEFAULT)) {
GetResponse existing = esClient.get(getRequest, RequestOptions.DEFAULT);
Map<String, Object> existingSource = existing.getSourceAsMap();
Object esMtime = existingSource.get("mtime");
if (esMtime != null && Long.parseLong(esMtime.toString()) == localMTime) {
logger.info("文件未变动,跳过导入: {}", relativePath);
fileCacheService.cacheMTime(relativePath, localMTime, 24 * 60);
return;
}
}
DeleteRequest deleteRequest = new DeleteRequest(INDEX_NAME, docId);
esClient.delete(deleteRequest, RequestOptions.DEFAULT);
logger.info("已删除旧版本文件: {}", relativePath);
Map<String, Object> doc = new HashMap<>();
doc.put("filename", fileName);
doc.put("filepath", relativePath);
doc.put("content", content);
doc.put("mtime", localMTime);
IndexRequest request = new IndexRequest(INDEX_NAME)
.id(docId)
.source(doc);
esClient.index(request, RequestOptions.DEFAULT);
logger.info("导入成功: {}", relativePath);
// 信息入缓存
fileCacheService.cacheMTime(relativePath, localMTime, 24 * 60);
} catch (Exception e) {
logger.error("导入失败: {}", path, e);
}
}
@Override
public String getType() {
return getDocTypeCode();

View File

@ -0,0 +1,61 @@
package com.knowledge.base.infrastructure.cache;
import java.util.Optional;
/**
* 文件缓存服务接口支持 Redis 或本地实现
*
* @author Luke.ye
* @date 2025/6/7
*/
public interface FileCacheService {
/**
* 获取文件名对应的相对路径
*
* @param filename 文件名
* @return 相对路径可能为 null
*/
Optional<String> getPath(String filename);
/**
* 缓存文件名与相对路径的映射
*
* @param filename 文件名
* @param relativePath 相对路径
* @param expireMinutes 过期时间分钟
*/
void cachePath(String filename, String relativePath, long expireMinutes);
/**
* 获取指定路径的上次修改时间用于变更检测
*
* @param relativePath 文件相对路径
* @return 本地上次修改时间可能为 null
*/
Optional<Long> getMTime(String relativePath);
/**
* 缓存指定路径的修改时间
*
* @param relativePath 文件路径
* @param mtime 修改时间
* @param expireMinutes 过期时间分钟
*/
void cacheMTime(String relativePath, Long mtime, long expireMinutes);
/**
* 清除文件名对应的路径缓存
*/
void clearPath(String filename);
/**
* 清除路径对应的修改时间缓存
*/
void clearMTime(String relativePath);
/**
* 清空所有缓存注意某些实现可能未实现
*/
void clearAll();
}

View File

@ -0,0 +1,67 @@
package com.knowledge.base.infrastructure.cache;
import com.knowledge.base.infrastructure.util.LocalCacheUtil;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.stereotype.Service;
import java.util.Optional;
/**
* 使用本地缓存实现的文件缓存服务 Redis
*
* @author Luke.ye
* @date 2025/6/7
*/
@Slf4j
@Service
@ConditionalOnProperty(prefix = "knowledge.base.redis", name = "enable", havingValue = "false")
public class LocalFileCacheService implements FileCacheService {
@Value("${knowledge.base.local-cache.expire-minutes:1440}")
private long expireMinutes;
private String key(String prefix, String key) {
return String.format("kb:file:%s:%s", prefix, key);
}
@Override
public Optional<String> getPath(String filename) {
Object val = LocalCacheUtil.get(key("path", filename));
return val instanceof String ? Optional.of((String) val) : Optional.empty();
}
@Override
public void cachePath(String filename, String relativePath, long expireMinutes) {
LocalCacheUtil.put(key("path", filename), relativePath);
}
@Override
public Optional<Long> getMTime(String relativePath) {
Object val = LocalCacheUtil.get(key("mtime", relativePath));
if (val instanceof Long) return Optional.of((Long) val);
if (val instanceof Integer) return Optional.of(((Integer) val).longValue());
return Optional.empty();
}
@Override
public void cacheMTime(String relativePath, Long mtime, long expireMinutes) {
LocalCacheUtil.put(key("mtime", relativePath), mtime);
}
@Override
public void clearPath(String filename) {
LocalCacheUtil.remove(key("path", filename));
}
@Override
public void clearMTime(String relativePath) {
LocalCacheUtil.remove(key("mtime", relativePath));
}
@Override
public void clearAll() {
LocalCacheUtil.clearAll();
}
}

View File

@ -0,0 +1,72 @@
package com.knowledge.base.infrastructure.cache;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Service;
import java.util.Optional;
import java.util.concurrent.TimeUnit;
/**
* 使用 Redis 实现的文件缓存服务
*
* @author Luke.ye
* @date 2025/6/7
*/
@Slf4j
@Service
@RequiredArgsConstructor
@ConditionalOnProperty(prefix = "knowledge.base.redis", name = "enable", havingValue = "true", matchIfMissing = true)
public class RedisFileCacheService implements FileCacheService {
private final StringRedisTemplate redisTemplate;
private String key(String prefix, String key) {
return String.format("kb:file:%s:%s", prefix, key);
}
@Override
public Optional<String> getPath(String filename) {
String value = redisTemplate.opsForValue().get(key("path", filename));
return Optional.ofNullable(value);
}
@Override
public void cachePath(String filename, String relativePath, long expireMinutes) {
redisTemplate.opsForValue().set(key("path", filename), relativePath, expireMinutes, TimeUnit.MINUTES);
}
@Override
public Optional<Long> getMTime(String relativePath) {
String val = redisTemplate.opsForValue().get(key("mtime", relativePath));
try {
return val != null ? Optional.of(Long.parseLong(val)) : Optional.empty();
} catch (NumberFormatException e) {
log.warn("mtime 解析失败: {}", val, e);
return Optional.empty();
}
}
@Override
public void cacheMTime(String relativePath, Long mtime, long expireMinutes) {
redisTemplate.opsForValue().set(key("mtime", relativePath), String.valueOf(mtime), expireMinutes, TimeUnit.MINUTES);
}
@Override
public void clearPath(String filename) {
redisTemplate.delete(key("path", filename));
}
@Override
public void clearMTime(String relativePath) {
redisTemplate.delete(key("mtime", relativePath));
}
@Override
public void clearAll() {
// 不清理 Redis 所有数据仅为接口保留占位符
log.warn("清空 Redis 文件缓存未实现,请手动清理前缀为 kb:file 的数据。");
}
}

View File

@ -0,0 +1,42 @@
package com.knowledge.base.infrastructure.config;
import lombok.Data;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.Configuration;
/**
* Redis 配置属性支持开关控制
*
* @author Luke.ye
* @date 2025/6/7
*/
@Data
@Configuration
@ConfigurationProperties(prefix = "knowledge.base.redis")
public class RedisConfig {
/**
* 是否启用 Redis 缓存
*/
private boolean enable = true;
/**
* Redis 服务器地址
*/
private String host = "localhost";
/**
* Redis 端口号
*/
private int port = 6379;
/**
* Redis 密码可选
*/
private String password;
/**
* 默认缓存时间分钟
*/
private long defaultExpireMinutes = 1440;
}

View File

@ -0,0 +1,32 @@
package com.knowledge.base.infrastructure.util;
import java.nio.charset.StandardCharsets;
import java.util.Base64;
/**
* 工具类将字符串转换为 Elasticsearch 安全的 _id 字符串
* 使用 Base64 URL 安全编码避免路径中出现特殊字符导致冲突
*
* @author Luke
*/
public class SafeIdUtil {
/**
* 将路径或任意字符串编码为 Base64 URL 安全形式可作为 Elasticsearch _id 使用
* @param original 原始字符串如文件相对路径
* @return 安全的 ID 字符串
*/
public static String encode(String original) {
return Base64.getUrlEncoder().withoutPadding()
.encodeToString(original.getBytes(StandardCharsets.UTF_8));
}
/**
* 将安全的 ID 解码还原为原始字符串如恢复文件路径
* @param encoded 已编码的安全 ID
* @return 原始字符串
*/
public static String decode(String encoded) {
return new String(Base64.getUrlDecoder().decode(encoded), StandardCharsets.UTF_8);
}
}

View File

@ -0,0 +1,52 @@
package com.knowledge.base.infrastructure.util;
import java.util.concurrent.*;
/**
* 通用线程池工具类
* @author Luke
*/
public class ThreadPoolUtil {
private static final int CORE_POOL_SIZE = Runtime.getRuntime().availableProcessors();
private static final int MAX_POOL_SIZE = CORE_POOL_SIZE * 2;
private static final int QUEUE_CAPACITY = 500;
private static final long KEEP_ALIVE_TIME = 60L;
private static final ThreadPoolExecutor THREAD_POOL_EXECUTOR = new ThreadPoolExecutor(
CORE_POOL_SIZE,
MAX_POOL_SIZE,
KEEP_ALIVE_TIME,
TimeUnit.SECONDS,
new LinkedBlockingQueue<>(QUEUE_CAPACITY),
new ThreadFactory() {
private final ThreadFactory defaultFactory = Executors.defaultThreadFactory();
private int count = 1;
@Override
public Thread newThread(Runnable r) {
Thread thread = defaultFactory.newThread(r);
thread.setName("knowledge-base-default-pool-" + count++);
thread.setDaemon(false);
return thread;
}
},
new ThreadPoolExecutor.CallerRunsPolicy() // 拒绝策略主线程执行
);
public static void execute(Runnable task) {
THREAD_POOL_EXECUTOR.execute(task);
}
public static void shutdownAndAwait() {
THREAD_POOL_EXECUTOR.shutdown();
try {
if (!THREAD_POOL_EXECUTOR.awaitTermination(30, TimeUnit.MINUTES)) {
THREAD_POOL_EXECUTOR.shutdownNow();
}
} catch (InterruptedException e) {
THREAD_POOL_EXECUTOR.shutdownNow();
Thread.currentThread().interrupt();
}
}
}

View File

@ -38,3 +38,15 @@ spring.datasource.driver-class-name=com.mysql.cj.jdbc.Driver
mybatis-plus.configuration.log-impl=org.apache.ibatis.logging.stdout.StdOutImpl
mybatis-plus.global-config.db-config.logic-delete-field=deleted
mybatis-plus.global-config.db-config.id-type=auto
# Redis 开关控制
knowledge.base.redis.enable=true
# Redis 服务配置
knowledge.base.redis.host=localhost
knowledge.base.redis.port=6379
knowledge.base.redis.password= # 可留空
# Redis 默认缓存时间(分钟)
knowledge.base.redis.default-expire-minutes=1440

View File

@ -39,3 +39,12 @@ spring.datasource.driver-class-name=com.mysql.cj.jdbc.Driver
mybatis-plus.configuration.log-impl=org.apache.ibatis.logging.stdout.StdOutImpl
mybatis-plus.global-config.db-config.logic-delete-field=deleted
mybatis-plus.global-config.db-config.id-type=auto
# Redis 开关控制
knowledge.base.redis.enable=true
# Redis 服务配置
knowledge.base.redis.host=host.docker.internal
knowledge.base.redis.port=6379
knowledge.base.redis.password= # 可留空
# Redis 默认缓存时间(分钟)
knowledge.base.redis.default-expire-minutes=1440

View File

@ -23,4 +23,6 @@ import.schedule.enabled=true
import.schedule.cron=0 0 * * * *
# Redis 开关控制
knowledge.base.redis.enable=false