task-98: 移除 similar-asin/appearance-patent 模块 Coze,状态机与共享组件改名 LLM
Build Backend JAR / build (push) Has been cancelled
Build Backend JAR / build (push) Has been cancelled
- similarasin/appearancepatent 模块全部 Coze 工作流调用改走 direct-LLM(已确认唯一运行路径) - 共享组件改名:CozeTaskQueueGate→TaskQueueGate、CozeGroupResultPropagator→GroupResultPropagator - 状态机改名:biz_task_scope_state 的 coze_* 列→llm_*、stateJson coze 键→llm(V100 迁移已应用生产) - 删除 biz_coze_credential 表、CozeCredential* 类、SimilarAsinCozeClient、AppearancePatentCozeClient→LlmClient - 前端 brand 页 Coze 文案→LLM;Python 后端删除 cozepy 依赖与死配置 - 修复 TaskResultFileJobWorker 启动失败:TaskFileJobConfig 注册 ResultFileJobHandlerRegistry 与 13 个 handler bean(含 validateCoverage 启动校验)
This commit is contained in:
@@ -19,7 +19,7 @@ public class CapacityPlanProperties {
|
||||
/** 数据库连接池(Hikari maximum-pool-size)。 */
|
||||
private int dbPoolMaxSize = 30;
|
||||
|
||||
/** 外部 HTTP 客户端(Coze/品牌/紫鸟)连接池容量。 */
|
||||
/** 外部 HTTP 客户端(LLM/品牌/紫鸟)连接池容量。 */
|
||||
private int httpClientPoolMaxSize = 32;
|
||||
|
||||
/** RustFS/MinIO OkHttp 连接池容量。 */
|
||||
|
||||
@@ -8,7 +8,7 @@ import java.time.Duration;
|
||||
|
||||
/**
|
||||
* Task 77:外部 HTTP 客户端统一连接复用池。
|
||||
* Coze / 品牌检查 / 紫鸟三个外部客户端共用同一个 java.net.http.HttpClient
|
||||
* LLM / 品牌检查 / 紫鸟三个外部客户端共用同一个 java.net.http.HttpClient
|
||||
* (内置 keep-alive 连接池),避免各自新建短命客户端导致连接无法复用、
|
||||
* 每次请求都重新建连。各客户端按自身超时创建独立的
|
||||
* JdkClientHttpRequestFactory(共享底层连接池),RestClient 单例懒加载。
|
||||
|
||||
@@ -4,6 +4,7 @@ import jakarta.servlet.FilterChain;
|
||||
import jakarta.servlet.ServletException;
|
||||
import jakarta.servlet.http.HttpServletRequest;
|
||||
import jakarta.servlet.http.HttpServletResponse;
|
||||
import java.util.Locale;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.beans.factory.annotation.Value;
|
||||
@@ -19,13 +20,15 @@ public class RequestTraceFilter extends OncePerRequestFilter {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(RequestTraceFilter.class);
|
||||
|
||||
private static final int MIN_REQUEST_BODY_CACHE_LIMIT_BYTES = 1024 * 1024;
|
||||
|
||||
private final InstanceMetadata instanceMetadata;
|
||||
private final int requestBodyCacheLimitBytes;
|
||||
|
||||
public RequestTraceFilter(InstanceMetadata instanceMetadata,
|
||||
@Value("${aiimage.instance-routing.request-body-cache-limit-bytes:104857600}") int requestBodyCacheLimitBytes) {
|
||||
@Value("${aiimage.instance-routing.request-body-cache-limit-bytes:1048576}") int requestBodyCacheLimitBytes) {
|
||||
this.instanceMetadata = instanceMetadata;
|
||||
this.requestBodyCacheLimitBytes = Math.max(1024 * 1024, requestBodyCacheLimitBytes);
|
||||
this.requestBodyCacheLimitBytes = Math.max(MIN_REQUEST_BODY_CACHE_LIMIT_BYTES, requestBodyCacheLimitBytes);
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -76,16 +79,23 @@ public class RequestTraceFilter extends OncePerRequestFilter {
|
||||
}
|
||||
}
|
||||
|
||||
private static HttpServletRequest wrapRequestIfNeeded(HttpServletRequest request, int requestBodyCacheLimitBytes) {
|
||||
static HttpServletRequest wrapRequestIfNeeded(HttpServletRequest request, int requestBodyCacheLimitBytes) {
|
||||
if (request instanceof ContentCachingRequestWrapper) {
|
||||
return request;
|
||||
}
|
||||
// multipart 不缓存:过滤器日志不读 body,缓存会整体复制上传流到内存
|
||||
String contentType = request.getContentType();
|
||||
if (contentType != null && contentType.trim().toLowerCase(Locale.ROOT).startsWith("multipart/")) {
|
||||
return request;
|
||||
}
|
||||
String method = request.getMethod();
|
||||
if ("POST".equalsIgnoreCase(method)
|
||||
|| "PUT".equalsIgnoreCase(method)
|
||||
|| "PATCH".equalsIgnoreCase(method)
|
||||
|| "DELETE".equalsIgnoreCase(method)) {
|
||||
return new ContentCachingRequestWrapper(request, requestBodyCacheLimitBytes);
|
||||
// 阈值统一钳到 1MB 下限:0/负数会令 ContentCachingRequestWrapper 构造抛错
|
||||
return new ContentCachingRequestWrapper(
|
||||
request, Math.max(MIN_REQUEST_BODY_CACHE_LIMIT_BYTES, requestBodyCacheLimitBytes));
|
||||
}
|
||||
return request;
|
||||
}
|
||||
|
||||
@@ -3,65 +3,38 @@ package com.nanri.aiimage.config;
|
||||
import lombok.Data;
|
||||
import org.springframework.boot.context.properties.ConfigurationProperties;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
@Data
|
||||
@ConfigurationProperties(prefix = "aiimage.similar-asin")
|
||||
public class SimilarAsinProperties {
|
||||
private String cozeBaseUrl = "https://api.coze.cn";
|
||||
private String cozeWorkflowPath = "/v1/workflow/run";
|
||||
private String cozeWorkflowHistoryPath = "/v1/workflows/{workflow_id}/run_histories/{execute_id}";
|
||||
private String cozeWorkflowId = "7635328462404583478";
|
||||
private String cozeToken = "";
|
||||
private List<CozeCredential> cozeCredentials = new ArrayList<>();
|
||||
private int cozeCredentialStripeSize = 0;
|
||||
/**
|
||||
* P0-1:单次提交 Coze 工作流的 row 数量。
|
||||
* P0-1:单次提交 LLM 批次的 row 数量。
|
||||
* 历史值 10,在含 puzzle 多图行的场景下频繁触发 720712008
|
||||
* "node executed out of limit: 1000"。降到 3 以避免节点上限被打爆。
|
||||
* 出现持续 720712008 时还会被 P1-1 滑窗自适应再降到 1。
|
||||
* 不影响 AppearancePatentProperties 的同名值。
|
||||
*/
|
||||
private int cozeBatchSize = 3;
|
||||
private int llmBatchSize = 3;
|
||||
/**
|
||||
* img_switch=false 时单次提交 Coze 的 row 数。
|
||||
* 不走图片检测时工作流压力小,恢复到 10 行一批以提高吞吐;开启图片检测时仍使用 cozeBatchSize。
|
||||
* img_switch=false 时单次提交 LLM 的 row 数。
|
||||
* 不走图片检测时压力小,恢复到 10 行一批以提高吞吐;开启图片检测时仍使用 llmBatchSize。
|
||||
*/
|
||||
private int cozeTextOnlyBatchSize = 10;
|
||||
private int cozeConnectTimeoutMillis = 10000;
|
||||
private int cozeReadTimeoutMillis = 60000;
|
||||
private int cozePollIntervalMillis = 30000;
|
||||
private int cozePollTimeoutMillis = 1800000;
|
||||
private int llmTextOnlyBatchSize = 10;
|
||||
private long dbTaskTouchIntervalMillis = 120000L;
|
||||
private long dbJobTouchIntervalMillis = 60000L;
|
||||
private int staleTimeoutMinutes = 30;
|
||||
private String staleFinalizeCron = "0 */2 * * * *";
|
||||
|
||||
/**
|
||||
* 同一 credential 两次提交之间的最小间隔(毫秒)。
|
||||
* 历史值硬编码 30000(持锁 sleep),导致单凭证仅 2 batch/分钟。
|
||||
* 几千行任务场景下成为提交吞吐瓶颈,下调到 5000ms 并改为锁外冷却。
|
||||
* 出现 Coze 限流加重时可通过 AIIMAGE_SIMILAR_ASIN_COZE_SUBMIT_MIN_INTERVAL_MILLIS 调高。
|
||||
*/
|
||||
private long cozeSubmitMinIntervalMillis = 5000L;
|
||||
|
||||
/**
|
||||
* 末尾零头 batch 的强制 flush 阈值(分钟):当不足 cozeBatchSize 的零头 row
|
||||
* 末尾零头 batch 的强制 flush 阈值(分钟):当不足 llmBatchSize 的零头 row
|
||||
* 长时间挂着(Python 慢回传)时触发提交。
|
||||
* 任务级实测:345 行 / 4h 总耗时中,约 2-3 小时是 batch 永远凑不满 batchSize 在等下一波回传,
|
||||
* 把阈值从 15 调到 1:最多 60s 后 1-2 行也强制提交,让 Coze 提交侧持续进票,
|
||||
* 总耗时降到与 Python 回传节奏接近。配合 cozeBatchSize=3、cozeSubmitMinIntervalMillis=5000,
|
||||
* 实际不会触发 Coze 限流。出现限流加重再调回 5/10。
|
||||
* 把阈值从 15 调到 1:最多 60s 后 1-2 行也强制提交,让 LLM 提交侧持续进票,
|
||||
* 总耗时降到与 Python 回传节奏接近。
|
||||
*/
|
||||
private int cozeFlushPendingMinutes = 1;
|
||||
private int llmFlushPendingMinutes = 1;
|
||||
|
||||
/**
|
||||
* 同 batch retry + split retry 共享的最大重试次数。原硬编码 5。
|
||||
* 图片下载、解码和缩放共享该池;4 核生产机默认 2,避免图片任务占满整机 CPU。
|
||||
*/
|
||||
private int cozeSubmitMaxRetryCount = 5;
|
||||
|
||||
/** 图片下载、解码和缩放共享该池;4 核生产机默认 2,避免图片任务占满整机 CPU。 */
|
||||
private int imageDownloadPoolSize = 2;
|
||||
|
||||
/**
|
||||
@@ -97,40 +70,11 @@ public class SimilarAsinProperties {
|
||||
private boolean imageDbCacheEnabled = false;
|
||||
|
||||
/**
|
||||
* 是否在 Coze 请求 parameters 中附带 api_key 字段。
|
||||
* 默认 true:线上 Coze 工作流将该字段视为必填,缺失会得到 4000
|
||||
* "Missing required parameters";前端传入的 api_key 必须透传到 coze。
|
||||
* 仅在工作流明确不再需要 api_key 时,可通过环境变量
|
||||
* AIIMAGE_SIMILAR_ASIN_COZE_INCLUDE_LEGACY_API_KEY=false 关闭。
|
||||
*/
|
||||
private boolean cozeIncludeLegacyApiKey = true;
|
||||
|
||||
/**
|
||||
* 是否使用旧的 item 字段顺序 {asin, sku, url, target_urls, title}。
|
||||
* 默认 false:当前实现使用 {asin, url, target_urls, title, sku}。
|
||||
* 出现兼容问题时可通过 AIIMAGE_SIMILAR_ASIN_COZE_USE_LEGACY_ITEM_ORDER=true
|
||||
* 切回旧顺序进行回归对比。
|
||||
*/
|
||||
private boolean cozeUseLegacyItemFieldOrder = false;
|
||||
|
||||
/**
|
||||
* 是否启用 P0-3 merge 增量缓冲:每个 batch DONE 时仅缓冲 cozeRows,
|
||||
* 是否启用 merge 增量缓冲:每个 batch DONE 时仅缓冲 llmRows,
|
||||
* 不立即合并到 chunk;finalize 前一次性按 chunkScopeHash 分组合并,
|
||||
* 把 OSS chunk 读写从 1000+ 次降到 chunk 数量级。
|
||||
* 仅作用于"正常 poll DONE"路径;失败 batch / 单 batch 任务 / 其他
|
||||
* 11 个 mergeCozeRowsIntoChunk 调用点保留原立即 merge 行为。
|
||||
* 出现问题时可通过 AIIMAGE_SIMILAR_ASIN_COZE_RESULT_BUFFER_ENABLED=false
|
||||
* 一键回滚到老路径。
|
||||
*/
|
||||
private boolean cozeResultBufferEnabled = true;
|
||||
|
||||
/**
|
||||
* P0-4:单 credential 抢 Coze 提交锁的最长等待时间(毫秒)。
|
||||
* 原硬编码 1000ms,在高并发 split retry 时大量抛 "Coze submit throttle lock timeout"
|
||||
* 并把整批行 markFailed。应与 cozeSubmitMinIntervalMillis(5000ms)保持 1.5-2 倍关系,
|
||||
* 默认 10000ms 给抢锁更多时间。
|
||||
*/
|
||||
private long cozeSubmitLockWaitMillis = 10000L;
|
||||
private boolean llmResultBufferEnabled = true;
|
||||
|
||||
/**
|
||||
* 解析接口返回的预览行/预览组数量上限。
|
||||
@@ -193,15 +137,9 @@ public class SimilarAsinProperties {
|
||||
*/
|
||||
private int imagePrefetchBudgetSeconds = 60;
|
||||
|
||||
/**
|
||||
* P0-4:抢 Coze 提交锁失败后下次重试间隔(毫秒)。
|
||||
* 原硬编码 500ms,会在指数退避算法中作为基础值(500/1000/2000/4000ms 上限 4000)。
|
||||
*/
|
||||
private long cozeSubmitLockRetryDelayMillis = 500L;
|
||||
|
||||
/**
|
||||
* 货源查询直连 LLM 模式开关(默认 true:新任务与存量 PENDING 批次都走直连 LLM,
|
||||
* 不再经过 Coze)。false 时回退到原 Coze 工作流链路(轮询/重试状态机保留)。
|
||||
* 不再经过工作流中转)。
|
||||
*/
|
||||
private boolean directLlmEnabled = true;
|
||||
|
||||
@@ -230,11 +168,4 @@ public class SimilarAsinProperties {
|
||||
|
||||
/** 拼接图/主图下载超时(秒),慢源图片较多时放大该值。 */
|
||||
private int llmImageDownloadTimeoutSeconds = 10;
|
||||
|
||||
@Data
|
||||
public static class CozeCredential {
|
||||
private String name;
|
||||
private String workflowId;
|
||||
private String token;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,5 +1,34 @@
|
||||
package com.nanri.aiimage.config;
|
||||
|
||||
import com.nanri.aiimage.modules.appearancepatent.service.AppearancePatentTaskService;
|
||||
import com.nanri.aiimage.modules.brand.service.BrandTaskService;
|
||||
import com.nanri.aiimage.modules.collectdata.service.CollectDataService;
|
||||
import com.nanri.aiimage.modules.deletebrand.service.DeleteBrandRunService;
|
||||
import com.nanri.aiimage.modules.patroldelete.service.PatrolDeleteTaskService;
|
||||
import com.nanri.aiimage.modules.pricetrack.service.PriceTrackTaskService;
|
||||
import com.nanri.aiimage.modules.productrisk.service.ProductRiskTaskService;
|
||||
import com.nanri.aiimage.modules.publish.service.PublishTaskService;
|
||||
import com.nanri.aiimage.modules.queryasin.service.QueryAsinTaskService;
|
||||
import com.nanri.aiimage.modules.shopdatacrawl.service.ShopDataCrawlTaskService;
|
||||
import com.nanri.aiimage.modules.shopmatch.service.ShopMatchTaskService;
|
||||
import com.nanri.aiimage.modules.similarasin.service.SimilarAsinTaskService;
|
||||
import com.nanri.aiimage.modules.task.service.AppearancePatentResultFileJobHandler;
|
||||
import com.nanri.aiimage.modules.task.service.BrandResultFileJobHandler;
|
||||
import com.nanri.aiimage.modules.task.service.CollectDataResultFileJobHandler;
|
||||
import com.nanri.aiimage.modules.task.service.DeleteBrandResultFileJobHandler;
|
||||
import com.nanri.aiimage.modules.task.service.PatrolDeleteResultFileJobHandler;
|
||||
import com.nanri.aiimage.modules.task.service.PriceTrackResultFileJobHandler;
|
||||
import com.nanri.aiimage.modules.task.service.ProductRiskResultFileJobHandler;
|
||||
import com.nanri.aiimage.modules.task.service.PublishResultFileJobHandler;
|
||||
import com.nanri.aiimage.modules.task.service.QueryAsinResultFileJobHandler;
|
||||
import com.nanri.aiimage.modules.task.service.ResultFileJobHandler;
|
||||
import com.nanri.aiimage.modules.task.service.ResultFileJobHandlerRegistry;
|
||||
import com.nanri.aiimage.modules.task.service.ShopDataCrawlResultFileJobHandler;
|
||||
import com.nanri.aiimage.modules.task.service.ShopMatchResultFileJobHandler;
|
||||
import com.nanri.aiimage.modules.task.service.SimilarAsinResultFileJobHandler;
|
||||
import com.nanri.aiimage.modules.task.service.TaskResultPayloadService;
|
||||
import com.nanri.aiimage.modules.task.service.WithdrawResultFileJobHandler;
|
||||
import com.nanri.aiimage.modules.withdraw.service.WithdrawTaskService;
|
||||
import io.micrometer.core.instrument.MeterRegistry;
|
||||
import org.springframework.beans.factory.ObjectProvider;
|
||||
import org.springframework.beans.factory.annotation.Value;
|
||||
@@ -9,6 +38,8 @@ import org.springframework.core.task.TaskExecutor;
|
||||
import org.springframework.scheduling.concurrent.ConcurrentTaskExecutor;
|
||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.Semaphore;
|
||||
@@ -16,6 +47,13 @@ import java.util.concurrent.Semaphore;
|
||||
@Configuration
|
||||
public class TaskFileJobConfig {
|
||||
|
||||
/** 结果文件 Job 支持的全部 moduleType(启动校验枚举源,见 ResultFileJobHandlerRegistry.validateCoverage) */
|
||||
public static final Set<String> RESULT_FILE_JOB_MODULE_TYPES = Set.of(
|
||||
"SHOP_MATCH", "PRICE_TRACK", "PRODUCT_RISK_RESOLVE",
|
||||
"PUBLISH", "QUERY_ASIN", "SHOP_DATA_CRAWL", "WITHDRAW",
|
||||
"PATROL_DELETE", "APPEARANCE_PATENT", "SIMILAR_ASIN",
|
||||
"DELETE_BRAND", "BRAND", "COLLECT_DATA");
|
||||
|
||||
@Bean("taskFileJobDispatchExecutor")
|
||||
public TaskExecutor taskFileJobDispatchExecutor(
|
||||
@Value("${aiimage.result-file-job.local-dispatch-pool-size:2}") int poolSize,
|
||||
@@ -33,24 +71,24 @@ public class TaskFileJobConfig {
|
||||
}
|
||||
|
||||
@Bean(destroyMethod = "shutdown")
|
||||
public ExecutorService cozeVirtualThreadExecutor() {
|
||||
public ExecutorService taskQueueVirtualThreadExecutor() {
|
||||
return Executors.newThreadPerTaskExecutor(Thread.ofVirtual()
|
||||
.name("coze-task-", 0)
|
||||
.name("task-queue-", 0)
|
||||
.factory());
|
||||
}
|
||||
|
||||
@Bean("cozeTaskExecutor")
|
||||
public TaskExecutor cozeTaskExecutor(
|
||||
ExecutorService cozeVirtualThreadExecutor,
|
||||
@Bean("taskQueueExecutor")
|
||||
public TaskExecutor taskQueueExecutor(
|
||||
ExecutorService taskQueueVirtualThreadExecutor,
|
||||
@Value("${aiimage.coze-task.max-concurrent:12}") int maxConcurrent,
|
||||
@Value("${aiimage.coze-task.max-waiting:1000}") int maxWaiting,
|
||||
ObjectProvider<MeterRegistry> meterRegistryProvider) {
|
||||
Semaphore semaphore = new Semaphore(Math.max(1, maxConcurrent));
|
||||
TaskExecutor semaphoreLimited = new ConcurrentTaskExecutor(command -> {
|
||||
if (command == null) {
|
||||
throw new IllegalArgumentException("coze 任务不能为 null");
|
||||
throw new IllegalArgumentException("task 不能为 null");
|
||||
}
|
||||
cozeVirtualThreadExecutor.execute(() -> {
|
||||
taskQueueVirtualThreadExecutor.execute(() -> {
|
||||
boolean acquired = false;
|
||||
try {
|
||||
semaphore.acquire();
|
||||
@@ -65,6 +103,87 @@ public class TaskFileJobConfig {
|
||||
}
|
||||
});
|
||||
});
|
||||
return new CozeTaskQueueGate(semaphoreLimited, maxWaiting, meterRegistryProvider);
|
||||
return new TaskQueueGate(semaphoreLimited, maxWaiting, meterRegistryProvider);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public ResultFileJobHandlerRegistry resultFileJobHandlerRegistry(List<ResultFileJobHandler> handlers) {
|
||||
ResultFileJobHandlerRegistry registry = new ResultFileJobHandlerRegistry(handlers);
|
||||
registry.validateCoverage(RESULT_FILE_JOB_MODULE_TYPES);
|
||||
return registry;
|
||||
}
|
||||
|
||||
@Bean
|
||||
public ResultFileJobHandler appearancePatentResultFileJobHandler(
|
||||
AppearancePatentTaskService appearancePatentTaskService) {
|
||||
return new AppearancePatentResultFileJobHandler(appearancePatentTaskService);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public ResultFileJobHandler brandResultFileJobHandler(
|
||||
BrandTaskService brandTaskService, TaskResultPayloadService taskResultPayloadService) {
|
||||
return new BrandResultFileJobHandler(brandTaskService, taskResultPayloadService);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public ResultFileJobHandler collectDataResultFileJobHandler(CollectDataService collectDataService) {
|
||||
return new CollectDataResultFileJobHandler(collectDataService);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public ResultFileJobHandler deleteBrandResultFileJobHandler(DeleteBrandRunService deleteBrandRunService) {
|
||||
return new DeleteBrandResultFileJobHandler(deleteBrandRunService);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public ResultFileJobHandler patrolDeleteResultFileJobHandler(
|
||||
PatrolDeleteTaskService patrolDeleteTaskService, TaskResultPayloadService taskResultPayloadService) {
|
||||
return new PatrolDeleteResultFileJobHandler(patrolDeleteTaskService, taskResultPayloadService);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public ResultFileJobHandler priceTrackResultFileJobHandler(
|
||||
PriceTrackTaskService priceTrackTaskService, TaskResultPayloadService taskResultPayloadService) {
|
||||
return new PriceTrackResultFileJobHandler(priceTrackTaskService, taskResultPayloadService);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public ResultFileJobHandler productRiskResultFileJobHandler(
|
||||
ProductRiskTaskService productRiskTaskService, TaskResultPayloadService taskResultPayloadService) {
|
||||
return new ProductRiskResultFileJobHandler(productRiskTaskService, taskResultPayloadService);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public ResultFileJobHandler publishResultFileJobHandler(PublishTaskService publishTaskService) {
|
||||
return new PublishResultFileJobHandler(publishTaskService);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public ResultFileJobHandler queryAsinResultFileJobHandler(
|
||||
QueryAsinTaskService queryAsinTaskService, TaskResultPayloadService taskResultPayloadService) {
|
||||
return new QueryAsinResultFileJobHandler(queryAsinTaskService, taskResultPayloadService);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public ResultFileJobHandler shopDataCrawlResultFileJobHandler(
|
||||
ShopDataCrawlTaskService shopDataCrawlTaskService, TaskResultPayloadService taskResultPayloadService) {
|
||||
return new ShopDataCrawlResultFileJobHandler(shopDataCrawlTaskService, taskResultPayloadService);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public ResultFileJobHandler shopMatchResultFileJobHandler(
|
||||
ShopMatchTaskService shopMatchTaskService, TaskResultPayloadService taskResultPayloadService) {
|
||||
return new ShopMatchResultFileJobHandler(shopMatchTaskService, taskResultPayloadService);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public ResultFileJobHandler similarAsinResultFileJobHandler(SimilarAsinTaskService similarAsinTaskService) {
|
||||
return new SimilarAsinResultFileJobHandler(similarAsinTaskService);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public ResultFileJobHandler withdrawResultFileJobHandler(
|
||||
WithdrawTaskService withdrawTaskService, TaskResultPayloadService taskResultPayloadService) {
|
||||
return new WithdrawResultFileJobHandler(withdrawTaskService, taskResultPayloadService);
|
||||
}
|
||||
}
|
||||
|
||||
+7
-7
@@ -11,7 +11,7 @@ import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
/**
|
||||
* Task 75:虚拟线程任务排队闸门。Coze 执行池的信号量只限制"正在执行"的
|
||||
* Task 75:虚拟线程任务排队闸门。任务执行池的信号量只限制"正在执行"的
|
||||
* 并发度,提交侧仍会在虚拟线程里无限排队。此闸门在提交时统计"已受理未启动"
|
||||
* 的等待数,达到上限立即拒绝并记录指标,防止等待队列无界堆积:
|
||||
* <ul>
|
||||
@@ -21,14 +21,14 @@ import java.util.concurrent.atomic.AtomicInteger;
|
||||
* </ul>
|
||||
*/
|
||||
@Slf4j
|
||||
public class CozeTaskQueueGate implements TaskExecutor {
|
||||
public class TaskQueueGate implements TaskExecutor {
|
||||
|
||||
private final TaskExecutor delegate;
|
||||
private final int maxWaiting;
|
||||
private final ObjectProvider<MeterRegistry> meterRegistryProvider;
|
||||
private final AtomicInteger waiting = new AtomicInteger();
|
||||
|
||||
public CozeTaskQueueGate(TaskExecutor delegate, int maxWaiting,
|
||||
public TaskQueueGate(TaskExecutor delegate, int maxWaiting,
|
||||
ObjectProvider<MeterRegistry> meterRegistryProvider) {
|
||||
this.delegate = delegate;
|
||||
this.maxWaiting = Math.max(1, maxWaiting);
|
||||
@@ -43,13 +43,13 @@ public class CozeTaskQueueGate implements TaskExecutor {
|
||||
public void execute(Runnable command) {
|
||||
if (command == null) {
|
||||
recordRejected("invalid-input");
|
||||
throw new IllegalArgumentException("coze 任务不能为 null");
|
||||
throw new IllegalArgumentException("task 不能为 null");
|
||||
}
|
||||
if (waiting.get() >= maxWaiting) {
|
||||
recordRejected("queue-full");
|
||||
log.warn("[coze-task][gate] waiting queue full, reject submit waiting={} limit={}",
|
||||
log.warn("[task-queue][gate] waiting queue full, reject submit waiting={} limit={}",
|
||||
waiting.get(), maxWaiting);
|
||||
throw new TaskRejectedException("coze 等待队列已满,limit=" + maxWaiting
|
||||
throw new TaskRejectedException("task 等待队列已满,limit=" + maxWaiting
|
||||
+ ", waiting=" + waiting.get());
|
||||
}
|
||||
waiting.incrementAndGet();
|
||||
@@ -68,7 +68,7 @@ public class CozeTaskQueueGate implements TaskExecutor {
|
||||
waiting.decrementAndGet();
|
||||
recordQueueWait(System.nanoTime() - queuedAt);
|
||||
recordRejected("delegate-rejected");
|
||||
log.warn("[coze-task][gate] delegate rejected submit waiting={} limit={} msg={}",
|
||||
log.warn("[task-queue][gate] delegate rejected submit waiting={} limit={} msg={}",
|
||||
waiting.get(), maxWaiting, ex.getMessage(), ex);
|
||||
throw ex;
|
||||
}
|
||||
Reference in New Issue
Block a user