Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
package interview.guide.common.config;

import jakarta.validation.constraints.Min;
import lombok.Data;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.stereotype.Component;
import org.springframework.validation.annotation.Validated;

/** 知识库同步上传链路的单实例并发配置,修改后需重启应用。 */
@Data
@Component
@Validated
@ConfigurationProperties(prefix = "app.knowledge-base-upload")
public class KnowledgeBaseUploadProperties {

@Min(1)
private int maxConcurrent = 4;
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
package interview.guide.infrastructure.file;

import interview.guide.common.config.KnowledgeBaseUploadProperties;
import interview.guide.common.exception.BusinessException;
import interview.guide.common.exception.ErrorCode;
import java.util.concurrent.Semaphore;
import java.util.function.Supplier;
import org.springframework.stereotype.Component;
import org.springframework.util.Assert;

/**
* 限制一个应用实例中正在执行的知识库上传,不排队等待。
*
* <p>不替代请求频率限制,也不限制 multipart 接收或后台向量化任务。
*/
@Component
public class KnowledgeBaseUploadLimiter {

private final Semaphore permits;

public KnowledgeBaseUploadLimiter(KnowledgeBaseUploadProperties properties) {
Assert.isTrue(properties.getMaxConcurrent() > 0, "知识库上传并发上限必须大于 0");
permits = new Semaphore(properties.getMaxConcurrent());
}

public <T> T execute(Supplier<T> upload) {
if (!permits.tryAcquire()) {
throw new BusinessException(ErrorCode.RATE_LIMIT_EXCEEDED, "知识库上传繁忙,请稍后重试");
}
try {
return upload.get();
} finally {
permits.release();
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
import interview.guide.infrastructure.file.FileHashService;
import interview.guide.infrastructure.file.FileStorageService;
import interview.guide.infrastructure.file.FileValidationService;
import interview.guide.infrastructure.file.KnowledgeBaseUploadLimiter;
import interview.guide.modules.knowledgebase.listener.VectorizeStreamProducer;
import interview.guide.modules.knowledgebase.model.KnowledgeBaseEntity;
import interview.guide.modules.knowledgebase.model.VectorStatus;
Expand Down Expand Up @@ -34,6 +35,7 @@ public class KnowledgeBaseUploadService {
private final FileValidationService fileValidationService;
private final FileHashService fileHashService;
private final VectorizeStreamProducer vectorizeStreamProducer;
private final KnowledgeBaseUploadLimiter uploadLimiter;

private static final long MAX_FILE_SIZE = 50 * 1024 * 1024; // 50MB

Expand All @@ -46,6 +48,11 @@ public class KnowledgeBaseUploadService {
* @return 上传结果和存储信息(包含duplicate字段,表示是否为重复上传)
*/
public Map<String, Object> uploadKnowledgeBase(MultipartFile file, String name, String category) {
return uploadLimiter.execute(() -> doUploadKnowledgeBase(file, name, category));
}

private Map<String, Object> doUploadKnowledgeBase(
MultipartFile file, String name, String category) {
// 1. 验证文件
fileValidationService.validateFile(file, MAX_FILE_SIZE, "知识库");

Expand Down
4 changes: 4 additions & 0 deletions app/src/main/resources/application.yml
Original file line number Diff line number Diff line change
Expand Up @@ -206,6 +206,10 @@ app:
- application/vnd.openxmlformats-officedocument.wordprocessingml.document
- text/plain

# 单个应用实例中同步处理的知识库上传数,不替代请求频率限制
knowledge-base-upload:
max-concurrent: ${APP_KNOWLEDGE_BASE_UPLOAD_MAX_CONCURRENT:4}

# RustFS (S3兼容) 存储配置
storage:
endpoint: ${APP_STORAGE_ENDPOINT:http://localhost:9000}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
package interview.guide.common.config;

import interview.guide.common.exception.BusinessException;
import interview.guide.infrastructure.file.KnowledgeBaseUploadLimiter;
import java.io.IOException;
import java.util.Map;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.ValueSource;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.boot.env.YamlPropertySourceLoader;
import org.springframework.boot.context.properties.bind.validation.BindValidationException;
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Import;
import org.springframework.core.env.MapPropertySource;
import org.springframework.core.io.ClassPathResource;

import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;

@DisplayName("上传并发配置绑定与校验")
class KnowledgeBaseUploadPropertiesTest {

private final ApplicationContextRunner contextRunner = new ApplicationContextRunner()
.withUserConfiguration(UploadConfiguration.class);

@Test
@DisplayName("未设置配置时应用使用默认四个上传名额")
void shouldUseDefaultLimit() {
contextRunner.run(context -> {
assertThat(context).hasNotFailed();
assertThat(context.getBean(KnowledgeBaseUploadProperties.class).getMaxConcurrent())
.isEqualTo(4);
assertThat(context.getBean(KnowledgeBaseUploadLimiter.class).execute(() -> "uploaded"))
.isEqualTo("uploaded");
});
}

@Test
@DisplayName("自定义并发配置实际作用于共享上传名额")
void shouldBindConfiguredLimitToLimiter() {
contextRunner.withPropertyValues("app.knowledge-base-upload.max-concurrent=2").run(context -> {
assertThat(context).hasNotFailed();
KnowledgeBaseUploadLimiter limiter = context.getBean(KnowledgeBaseUploadLimiter.class);

limiter.execute(() -> limiter.execute(() -> {
assertThatThrownBy(() -> limiter.execute(() -> "excess upload"))
.isInstanceOf(BusinessException.class);
return "uploaded";
}));

assertThat(limiter.execute(() -> "uploaded")).isEqualTo("uploaded");
});
}

@Test
@DisplayName("application.yml 中的环境变量占位符可覆盖默认上传上限")
void shouldApplyEnvironmentOverrideFromApplicationYaml() throws IOException {
var sources = new YamlPropertySourceLoader()
.load("application", new ClassPathResource("application.yml"));
contextRunner.withInitializer(context -> {
context.getEnvironment().getPropertySources().addFirst(new MapPropertySource(
"upload-override", Map.of("APP_KNOWLEDGE_BASE_UPLOAD_MAX_CONCURRENT", "1")));
sources.forEach(source -> context.getEnvironment().getPropertySources().addLast(source));
}).run(context -> {
assertThat(context).hasNotFailed();
KnowledgeBaseUploadLimiter limiter = context.getBean(KnowledgeBaseUploadLimiter.class);

limiter.execute(() -> {
assertThatThrownBy(() -> limiter.execute(() -> "excess upload"))
.isInstanceOf(BusinessException.class);
return "uploaded";
});

assertThat(limiter.execute(() -> "uploaded")).isEqualTo("uploaded");
});
}

@ParameterizedTest
@ValueSource(ints = {0, -1})
@DisplayName("零或负数并发配置在启动时明确失败")
void shouldRejectInvalidLimit(int maximum) {
contextRunner.withPropertyValues("app.knowledge-base-upload.max-concurrent=" + maximum)
.run(context -> {
assertThat(context).hasFailed();
assertThat(context.getStartupFailure())
.hasRootCauseInstanceOf(BindValidationException.class);
});
}

@Configuration(proxyBeanMethods = false)
@EnableConfigurationProperties(KnowledgeBaseUploadProperties.class)
@Import(KnowledgeBaseUploadLimiter.class)
static class UploadConfiguration {
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
package interview.guide.infrastructure.file;

import interview.guide.common.config.KnowledgeBaseUploadProperties;
import interview.guide.common.exception.BusinessException;
import interview.guide.common.exception.ErrorCode;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.ValueSource;

import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;

@DisplayName("上传名额释放边界")
class KnowledgeBaseUploadLimiterTest {

@ParameterizedTest
@ValueSource(booleans = {false, true})
@DisplayName("业务异常与 Error 原样抛出,之后仍可再次执行上传")
void shouldReleasePermitAndPreserveFailure(boolean error) {
KnowledgeBaseUploadProperties properties = new KnowledgeBaseUploadProperties();
properties.setMaxConcurrent(1);
KnowledgeBaseUploadLimiter limiter = new KnowledgeBaseUploadLimiter(properties);
BusinessException businessFailure = new BusinessException(ErrorCode.STORAGE_UPLOAD_FAILED);
AssertionError unexpectedFailure = new AssertionError("unexpected failure");

assertThatThrownBy(() -> limiter.execute(() -> {
if (error) {
throw unexpectedFailure;
}
throw businessFailure;
})).isSameAs(error ? unexpectedFailure : businessFailure);

assertThat(limiter.execute(() -> "uploaded")).isEqualTo("uploaded");
}
}
Loading