feat(game-cloud): RocketMQ 异步 gen 队列(生产者+有界消费≤15+CAS 幂等)(切片一 阶段〇 C3)
入队发 AIGC_GEN_TOPIC 信号替 @Scheduled 轮询触发;消费者 consumeThreadMax=15(在跑并发权威)+ 既有 claimQueuedTask CAS 认领幂等;派发失败落 failed 不留悬空 running;@Scheduled 降级兜底(捞 stale queued + 超时 running)。发 MQ best-effort 不回滚入队。单测锁幂等/best-effort;端到端随后端窗口。 name-server 各 profile(dev/local/staging)+base 指 mini-infra 100.64.0.8:9876;rocketmq-spring-boot-starter 在 starter-mq/websocket 均 optional=true 不传递,故 aigc-server 直依赖(版本走依赖管理)。 Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
parent
6d6bb52d9d
commit
4bc36ec6dd
@ -48,7 +48,7 @@ public interface ErrorCodeConstants {
|
||||
* 入口 QPS 突发闸拦截(切片一 阶段〇 C2):Sentinel FLOW_GRADE_QPS(规则热源 Nacos dataId sentinel-gen-flow-rules)
|
||||
* 判入队速率超限 → blockHandler 优雅转此业务码(经全局 GlobalExceptionHandler 转 CommonResult.error,HTTP 仍 200、
|
||||
* body.code 非 0),绝不 500、不静默丢。与门①②③(per-creator 业务配额)分层:本码 = 入口速率闸削洪峰;
|
||||
* 「在跑生成并发≤15」由 C3 RocketMQ consumeThreadMax 保证、不由本门保证(三层正交)。
|
||||
* C3 RocketMQ consumeThreadMax=15 只 cap dispatch 投递速率、真「在跑生成并发≤15」由有界 worker 池保证(follow-up,不在 C3),均不由本门保证(三层正交)。
|
||||
*/
|
||||
ErrorCode AIGC_ADMISSION_RATE_LIMITED = new ErrorCode(1_101_001_006, "生成服务繁忙(入队速率超限),请稍后重试");
|
||||
|
||||
|
||||
@ -59,7 +59,7 @@
|
||||
· spring-cloud-starter-alibaba-sentinel:@SentinelResource 切面 + 流控运行时(SentinelResourceAspect 自动装配);
|
||||
· sentinel-datasource-nacos:流控规则从 Nacos dataId sentinel-gen-flow-rules 热加载(改即生效、不重启)。
|
||||
版本由 huijing-dependencies 的 spring-cloud-alibaba-dependencies:2025.0.0.0 BOM 托管(实际 Sentinel 1.8.9),不写死。
|
||||
职责 = 生成入队入口速率闸(FLOW_GRADE_QPS 削洪峰),与 DB 三门(业务配额)/ C3 consumeThreadMax(在跑并发≤15)三层正交。 -->
|
||||
职责 = 生成入队入口速率闸(FLOW_GRADE_QPS 削洪峰),与 DB 三门(业务配额)/ C3 consumeThreadMax(dispatch 投递速率上限;真在跑并发≤15 属有界 worker 池=follow-up)三层正交。 -->
|
||||
<dependency>
|
||||
<groupId>com.alibaba.cloud</groupId>
|
||||
<artifactId>spring-cloud-starter-alibaba-sentinel</artifactId>
|
||||
@ -137,6 +137,13 @@
|
||||
<artifactId>huijing-spring-boot-starter-biz-tenant</artifactId>
|
||||
</dependency>
|
||||
|
||||
<!-- RocketMQ(切片一 阶段〇 C3:异步 gen 队列 生产者+有界并发派发消费者 consumeThreadMax=15):RocketMQTemplate / @RocketMQMessageListener。
|
||||
starter-mq / starter-websocket 均以 optional=true 声明本 starter(不传递),故本模块须显式直依赖;版本由 huijing 依赖管理统一控制。 -->
|
||||
<dependency>
|
||||
<groupId>org.apache.rocketmq</groupId>
|
||||
<artifactId>rocketmq-spring-boot-starter</artifactId>
|
||||
</dependency>
|
||||
|
||||
<!-- Web + 安全:@RestController / @PreAuthorize / 当前登录用户 -->
|
||||
<dependency>
|
||||
<groupId>com.wanxiang</groupId>
|
||||
|
||||
@ -12,15 +12,18 @@ import static com.wanxiang.huijing.framework.common.exception.util.ServiceExcept
|
||||
* 生成入队的 Sentinel 资源门 —— 入口 QPS 突发保护(FLOW_GRADE_QPS,防洪峰)。切片一 阶段〇 C2。
|
||||
*
|
||||
* <p><b>职责边界(经 Opus 评审纠正,别把空体 admit() 当缺陷)</b>:本门统计「每秒入队请求数」,不是「在跑生成并发」。
|
||||
* 空体入口的 entry/exit 是纳秒级,thread-grade 统计不到真正在跑的生成——那发生在 admit() 返回之后的 RocketMQ
|
||||
* 消费线程里。故「在跑生成并发≤15」由 C3 消费者 consumeThreadMax=15 保证,不由本门保证。
|
||||
* 空体入口的 entry/exit 是纳秒级,thread-grade 统计不到真正在跑的生成——真生成发生在 admit() 返回、入队、消费线程
|
||||
* 投递之后的 worker 侧(异步跑数十秒)。故在跑生成并发不由本门保证:C3 consumeThreadMax=15 只 cap dispatch 投递速率,
|
||||
* 真「在跑生成并发≤15」由有界 worker 池保证(follow-up,不在 C3)。
|
||||
*
|
||||
* <p><b>三层正交</b>:
|
||||
* <ul>
|
||||
* <li>入口速率(本门 · Sentinel FLOW_GRADE_QPS):每秒入队请求数上限,削洪峰;</li>
|
||||
* <li>在跑并发(C3 · RocketMQ consumeThreadMax=15):同时在跑的生成数结构性上限;</li>
|
||||
* <li>投递速率(C3 · RocketMQ consumeThreadMax=15):每秒投递 job 给 worker 的并发上限(非在跑生成并发);</li>
|
||||
* <li>业务配额(既有 DB 三门):per-creator×level 日额度 + per-creator 并发 + 全局背压,在 admit() 之后执行。</li>
|
||||
* </ul>
|
||||
* <p>真「在跑生成并发≤15」的硬上限由<b>有界 worker 池</b>保证(创始人 2026-06-26 决策:异步任务队列 + 有界 worker 池
|
||||
* N=4-6 起步压测调、≤15 硬上限),作 follow-up 工作单元实现、不在 C3。
|
||||
*
|
||||
* <p><b>越限行为</b>:Sentinel 抛 {@link BlockException} → {@link #onAdmissionBlocked} 优雅转成业务异常
|
||||
* {@code AIGC_ADMISSION_RATE_LIMITED}「生成服务繁忙…请稍后重试」(经全局 GlobalExceptionHandler 转 CommonResult.error,
|
||||
@ -43,7 +46,7 @@ public class GenAdmissionResource {
|
||||
* 通过则直接返回(调用方继续走 DB 三门 + 入队);被 block 则由 Sentinel 切面转调 {@link #onAdmissionBlocked}。
|
||||
*
|
||||
* <p><b>空体是设计,不是缺陷</b>:{@code FLOW_GRADE_QPS} 只统计「每秒通过本资源的次数」(与方法耗时无关),
|
||||
* 空体正合适——让 Sentinel 计入队 QPS 削洪峰。业务准入(配额)在调用方 DB 三门;在跑并发≤15 在 C3 consumeThreadMax,均不在此。
|
||||
* 空体正合适——让 Sentinel 计入队 QPS 削洪峰。业务准入(配额)在调用方 DB 三门;C3 consumeThreadMax=15 cap dispatch 投递速率、真在跑并发≤15 属有界 worker 池(follow-up),均不在此。
|
||||
*/
|
||||
@SentinelResource(value = RESOURCE, blockHandler = "onAdmissionBlocked")
|
||||
public void admit() {
|
||||
|
||||
@ -9,6 +9,7 @@ import com.wanxiang.huijing.game.module.aigc.service.executor.ExecutorLlmClient;
|
||||
import com.wanxiang.huijing.game.module.aigc.service.executor.GameConfigSchemaValidator;
|
||||
import com.wanxiang.huijing.game.module.aigc.service.executor.PromptResourceLoader;
|
||||
import com.wanxiang.huijing.game.module.aigc.service.executor.WorkerClassifyClient;
|
||||
import com.wanxiang.huijing.game.module.aigc.mq.GenTaskProducer;
|
||||
import com.wanxiang.huijing.game.module.aigc.saa.SaaGraphDispatcher;
|
||||
import com.wanxiang.huijing.game.module.aigc.service.executor.GenerationDispatcher;
|
||||
import com.wanxiang.huijing.game.module.aigc.service.executor.WorkerDispatchClient;
|
||||
@ -163,6 +164,8 @@ public class AigcExecutorConfiguration {
|
||||
* @param llmClient LLM 通道
|
||||
* @param workerDispatchClient 外置生成 worker 派发通道(P3 generic 路)
|
||||
* @param sourceProjectApiProvider 源项目反查 seam(A11 M1②:modify 路据 baseVersionId 反查 base 源注入 job.sourceProject;软取)
|
||||
* @param genTaskProducerProvider gen 队列生产者 seam(C3:@Scheduled 兜底 tick 对 stale queued 重发信号交有界消费者;软取,
|
||||
* 缺席→兜底退回 tick 直派,非破坏性。与本配置同 aigc.executor.enabled 条件,enabled 时恒在席)
|
||||
* @return 执行器(@Scheduled 注解由全局调度后处理器消费,fixedDelay 单飞)
|
||||
*/
|
||||
@Bean
|
||||
@ -174,15 +177,20 @@ public class AigcExecutorConfiguration {
|
||||
ExecutorLlmClient llmClient,
|
||||
WorkerDispatchClient workerDispatchClient,
|
||||
SaaGraphDispatcher saaGraphDispatcher,
|
||||
ObjectProvider<SourceProjectApi> sourceProjectApiProvider) {
|
||||
ObjectProvider<SourceProjectApi> sourceProjectApiProvider,
|
||||
ObjectProvider<GenTaskProducer> genTaskProducerProvider) {
|
||||
// A11 M1②:软取 SourceProjectApi(studio @Primary 就地解析;modify 路据 baseVersionId 反查 base 源注入 §6.1 job 的 sourceProject 键,
|
||||
// 供便宜档 HTTP worker 无 DB 被动消费——与 SaaGraphDispatcher 同一 seam)。缺席(单模块装配/studio 未在 classpath)→ modify 取源旁路(不注入 sourceProject 键),不阻断装配。
|
||||
SourceProjectApi sourceProjectApi = sourceProjectApiProvider.getIfAvailable();
|
||||
// SAA 迁移 #2:传入 http worker + saa 进程内派发器,由执行器 10 参构造按 aigc.executor.dispatcher 选用
|
||||
// C3:软取 GenTaskProducer(与本配置同 aigc.executor.enabled 条件装配,enabled 时恒在席)——传入使 @Scheduled 兜底 tick 对
|
||||
// stale queued「重发信号」交消费者重走认领+派发,不在 tick 直派(投递可靠性补偿,让 MQ 保持唯一投递路径;与并发 cap 无关)。
|
||||
// 缺席(防御分支)→ 传 null → 兜底 tick 退回「直接认领+派发」(韧性降级、不让任务卡死)。
|
||||
GenTaskProducer genTaskProducer = genTaskProducerProvider.getIfAvailable();
|
||||
// SAA 迁移 #2:传入 http worker + saa 进程内派发器,由执行器构造按 aigc.executor.dispatcher 选用
|
||||
// (默认 http=现行行为不变;saa=opt-in 进程内形态①)。
|
||||
return new AigcGenerateExecutor(properties, aigcTaskMapper, difyCallbackService,
|
||||
promptResourceLoader, schemaValidator, llmClient, workerDispatchClient, saaGraphDispatcher,
|
||||
sourceProjectApi, Clock.systemDefaultZone());
|
||||
sourceProjectApi, Clock.systemDefaultZone(), genTaskProducer);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@ -0,0 +1,49 @@
|
||||
package com.wanxiang.huijing.game.module.aigc.mq;
|
||||
|
||||
import com.wanxiang.huijing.game.module.aigc.service.executor.AigcGenerateExecutor;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
|
||||
import org.springframework.stereotype.Service;
|
||||
|
||||
/**
|
||||
* 生成任务「认领 + 派发」服务(切片一 阶段〇 C3)——MQ 消费者与既有执行器之间的薄接缝。
|
||||
*
|
||||
* <p><b>复用既有、不重写</b>:把「CAS 认领({@code claimQueuedTask}) + 生成流水派发({@code runGenerationPipeline}
|
||||
* → generic 走 {@code SaaGraphDispatcher}/http worker)」收敛为一次调用 {@link AigcGenerateExecutor#claimAndDispatchById}。
|
||||
* 与执行器 {@code @Scheduled} tick 的 {@code processCandidate} 同一执行体、同一系统身份 + 租户上下文,只是触发面
|
||||
* 从「批量扫描」换成「MQ 单 taskId 事件」。<b>不复制三表写链、不重写派发逻辑</b>。
|
||||
*
|
||||
* <p><b>幂等 + 派发失败落 failed</b>(细节见 {@link AigcGenerateExecutor#claimAndDispatchById}):
|
||||
* CAS 认领失败(被并发抢走/非 queued/已取消/任务不存在)→ 返回 false 幂等丢弃;认领成功但派发抛异常
|
||||
* → 内部落 failed(不留悬空 running)后返回 true。故本方法<b>不外抛</b>,消费者据返回值只做日志、不重投。
|
||||
*
|
||||
* <p><b>装配</b>:随 {@code aigc.executor.enabled=true} 装配(与执行器/生产者/消费者同一开关);执行器 Bean 由
|
||||
* {@code AigcExecutorConfiguration} 同条件注册,故此处硬注入 {@link AigcGenerateExecutor} 恒可解析。
|
||||
*
|
||||
* @author 绘境AI
|
||||
*/
|
||||
@Slf4j
|
||||
@Service
|
||||
@ConditionalOnProperty(prefix = "aigc.executor", name = "enabled", havingValue = "true")
|
||||
public class AigcGenDispatchService {
|
||||
|
||||
/** 既有生成执行器(AigcExecutorConfiguration @Bean 注册,同 aigc.executor.enabled 条件)——复用其认领+派发执行体。 */
|
||||
private final AigcGenerateExecutor generateExecutor;
|
||||
|
||||
public AigcGenDispatchService(AigcGenerateExecutor generateExecutor) {
|
||||
this.generateExecutor = generateExecutor;
|
||||
}
|
||||
|
||||
/**
|
||||
* 认领并派发一个生成任务(MQ 消费者收到 taskId 信号后调用)。
|
||||
*
|
||||
* @param taskId 任务 ID(消息体轻信号)
|
||||
* @return true=CAS 认领成功(已派发或已落终态);false=认领未命中=幂等丢弃(同 taskId 重复投递/被并发抢走/非 queued)
|
||||
*/
|
||||
public boolean claimAndDispatch(Long taskId) {
|
||||
// 委托既有执行器:CAS 认领 + 复用 runGenerationPipeline 派发;认领失败返 false,派发抛异常内部落 failed(不外抛)。
|
||||
boolean claimed = generateExecutor.claimAndDispatchById(taskId);
|
||||
log.debug("[aigc-mq] 认领+派发结论 taskId={} claimed={}", taskId, claimed);
|
||||
return claimed;
|
||||
}
|
||||
}
|
||||
@ -0,0 +1,74 @@
|
||||
package com.wanxiang.huijing.game.module.aigc.mq;
|
||||
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
|
||||
import org.apache.rocketmq.spring.core.RocketMQListener;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
/**
|
||||
* 生成任务 MQ 消费者(切片一 阶段〇 C3)——有界并发派发。
|
||||
*
|
||||
* <p><b>{@code consumeThreadMax=15} = dispatch(投递握手)并发速率上限,不是在跑生成并发上限</b>:消费线程收到
|
||||
* taskId 信号后 CAS 认领 → 组 §6.1 job 投递生成躯体(http worker POST 拿 ACK / saa 进程内交后台池),投递握手即
|
||||
* 返回释放(sub-second),<b>不在消费线程内长挂等生成结果</b>——生成是 worker 异步跑数十秒、回调驱动终态(见
|
||||
* {@link com.wanxiang.huijing.game.module.aigc.service.executor.GenerationDispatcher}「仅投递握手,不等生成结果」契约)。
|
||||
* 故这 15 线程 cap 的是「每秒能投递多少 job」,而非「同时真跑多少生成」(投递后任务即留 RUNNING,真在跑的生成挂在 worker 侧)。
|
||||
*
|
||||
* <p><b>真「在跑生成并发≤15」由有界 worker 池保证(follow-up,不在 C3)</b>:符合创始人 2026-06-26 决策——异步任务
|
||||
* 队列 + 有界 worker 池(N=4-6 起步、压测调,≤15 硬上限);该 in-flight 硬上限作独立 follow-up 工作单元实现,本 C3
|
||||
* 只落队列机制。三层正交:C2 Sentinel 管入口 QPS「每秒进多少」、DB 三门管业务配额「每人/全局多少在飞」、本消费线程
|
||||
* 池管投递速率「每秒投递多少 job」。多实例扩容后的集群级全局限流留放量后(Token Server)。
|
||||
*
|
||||
* <p><b>幂等(CAS 认领去重)</b>:收到 taskId 信号 → 经 {@link AigcGenDispatchService#claimAndDispatch} 用既有
|
||||
* CAS {@code claimQueuedTask}(queued→running) 认领,认领成功才派发;认领失败(已被别的消费线程/兜底 tick 拿走、
|
||||
* 非 queued、已取消、任务不存在)即幂等丢弃——同一 taskId 重复投递不会双跑。
|
||||
*
|
||||
* <p><b>不无限重投</b>:{@code claimAndDispatch} 内部已处置派发异常(认领成功但派发抛异常 → 落 failed,不留悬空
|
||||
* running),正常返回 true/false 不外抛;故 {@link #onMessage} 认领失败/毒消息只记日志、正常返回(ack),
|
||||
* 不让 RocketMQ 无限重投。消费模型 = 默认 CLUSTERING(每条消息集群内单实例消费,即工作队列语义)。
|
||||
*
|
||||
* <p><b>装配</b>:随 {@code aigc.executor.enabled=true} 装配(与执行器/生产者同一开关)——executor 关闭时不起消费者、
|
||||
* 不连 NameServer(避免无执行器时空转/无谓连接)。
|
||||
*
|
||||
* @author 绘境AI
|
||||
*/
|
||||
@Slf4j
|
||||
@Component
|
||||
@ConditionalOnProperty(prefix = "aigc.executor", name = "enabled", havingValue = "true")
|
||||
@RocketMQMessageListener(
|
||||
topic = GenTaskProducer.TOPIC,
|
||||
consumerGroup = "aigc_gen_consumer",
|
||||
consumeThreadMax = 15, // rocketmq-client 无界队列致 consumeThreadMax 不单独生效(线程数永不超 core),故与 consumeThreadNumber 同设 15 让投递并发上限真为 15
|
||||
consumeThreadNumber = 15 // min=max=15:15 个投递线程常驻(投递握手 sub-second、空闲线程零成本);这是 dispatch 投递并发、非在跑生成并发(后者=有界 worker 池 follow-up)
|
||||
)
|
||||
public class GenTaskConsumer implements RocketMQListener<String> {
|
||||
|
||||
/** 认领+派发接缝(复用既有执行器 CAS 认领 + runGenerationPipeline 派发,见 {@link AigcGenDispatchService})。 */
|
||||
private final AigcGenDispatchService dispatchService;
|
||||
|
||||
public GenTaskConsumer(AigcGenDispatchService dispatchService) {
|
||||
this.dispatchService = dispatchService;
|
||||
}
|
||||
|
||||
/**
|
||||
* 消费一条 gen 信号:解析 taskId → CAS 认领 + 派发(认领失败幂等丢弃)。
|
||||
*
|
||||
* @param taskIdStr 消息体(taskId 的字符串形态,见 {@link GenTaskProducer#publishNewTask})
|
||||
*/
|
||||
@Override
|
||||
public void onMessage(String taskIdStr) {
|
||||
Long taskId;
|
||||
try {
|
||||
taskId = Long.valueOf(taskIdStr);
|
||||
} catch (NumberFormatException e) {
|
||||
// 毒消息(taskId 非法):重投无益,直接丢弃(正常返回=ack),只告警留痕,防无限重投风暴。
|
||||
log.error("[aigc-mq] gen 信号 taskId 非法,丢弃不重投 body={}", taskIdStr, e);
|
||||
return;
|
||||
}
|
||||
// CAS 认领(既有 claimQueuedTask:仅 queued→running 才成功)→ 认领成功才跑,失败即幂等丢弃。
|
||||
// claimAndDispatch 内部已处置派发异常(落 failed 不留悬空 running),不外抛——故此处只记日志、不重投。
|
||||
boolean claimed = dispatchService.claimAndDispatch(taskId);
|
||||
log.info("[aigc-mq] 消费 gen 信号 taskId={} claimed={}", taskId, claimed);
|
||||
}
|
||||
}
|
||||
@ -0,0 +1,59 @@
|
||||
package com.wanxiang.huijing.game.module.aigc.mq;
|
||||
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.rocketmq.spring.core.RocketMQTemplate;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
/**
|
||||
* 生成任务 MQ 生产者(切片一 阶段〇 C3)——入队成功后发一条「有新任务」信号,把「触发」从执行器
|
||||
* {@code @Scheduled} DB 轮询改为事件驱动(RocketMQ)。
|
||||
*
|
||||
* <p><b>消息体只放 taskId(轻信号)</b>:任务真身/状态仍以 {@code game_aigc_task} 表为准——MQ 是投递触发、
|
||||
* 不是数据搬运。消费侧据 taskId 用既有 CAS {@code claimQueuedTask}(queued→running) 认领,认领成功才派发,
|
||||
* 同一 taskId 重复投递被 CAS 天然去重(不双跑)。
|
||||
*
|
||||
* <p><b>best-effort(命门)</b>:发 MQ 失败<b>不回滚入队</b>(任务已落库 queued)、也<b>不外抛</b>——只告警留痕。
|
||||
* 丢失的信号由执行器 {@code @Scheduled} 降级兜底轮询补偿捡起(stale queued 重发信号,见
|
||||
* {@link com.wanxiang.huijing.game.module.aigc.service.executor.AigcGenerateExecutor} 兜底补偿路)。
|
||||
* 故 MQ 抖动只影响「即时性」、不影响「不丢任务」。
|
||||
*
|
||||
* <p><b>装配条件</b>:随 {@code aigc.executor.enabled=true} 装配(与执行器/消费者同一开关)——executor 关闭时
|
||||
* 无消费者、无执行器,发信号无意义,故一并缺席;{@link com.wanxiang.huijing.game.module.aigc.service.task.AigcTaskServiceImpl}
|
||||
* 以 {@code ObjectProvider} 软注入本 Bean,缺席时入队照常落库、仅不发信号(现行骨架行为)。
|
||||
*
|
||||
* @author 绘境AI
|
||||
*/
|
||||
@Slf4j
|
||||
@Component
|
||||
@ConditionalOnProperty(prefix = "aigc.executor", name = "enabled", havingValue = "true")
|
||||
public class GenTaskProducer {
|
||||
|
||||
/** gen 队列 topic(与 B2 broker 在 mini-infra 建的 topic 对齐;消费者 {@link GenTaskConsumer} 同名订阅)。 */
|
||||
public static final String TOPIC = "AIGC_GEN_TOPIC";
|
||||
|
||||
/** RocketMQ 模板(rocketmq-spring 自动装配;name-server 见 application*.yaml → mini-infra 100.64.0.8:9876)。 */
|
||||
private final RocketMQTemplate rocketMQTemplate;
|
||||
|
||||
public GenTaskProducer(RocketMQTemplate rocketMQTemplate) {
|
||||
this.rocketMQTemplate = rocketMQTemplate;
|
||||
}
|
||||
|
||||
/**
|
||||
* 发「新任务」信号(同步发,拿发送结果落日志;失败只告警、不抛,由执行器兜底轮询补偿)。
|
||||
*
|
||||
* @param taskId 生成任务 ID(消息体轻信号;消费侧据此反查 game_aigc_task + CAS 认领)
|
||||
*/
|
||||
public void publishNewTask(Long taskId) {
|
||||
try {
|
||||
// 同步发:payload 只放 taskId 字符串(轻信号)。syncSend(String, Object) 拿 SendResult 落痕便于排障。
|
||||
var result = rocketMQTemplate.syncSend(TOPIC, String.valueOf(taskId));
|
||||
log.info("[aigc-mq] 发送 gen 任务信号 topic={} taskId={} msgId={} status={}",
|
||||
TOPIC, taskId, result != null ? result.getMsgId() : null,
|
||||
result != null ? result.getSendStatus() : null);
|
||||
} catch (Exception e) {
|
||||
// best-effort:MQ 抖动不连累入队(任务已落库 queued),也不外抛——由 @Scheduled 兜底轮询补偿捡起。
|
||||
log.warn("[aigc-mq] 发送 gen 信号失败(不回滚入队、不外抛,由执行器兜底轮询补偿)taskId={}", taskId, e);
|
||||
}
|
||||
}
|
||||
}
|
||||
@ -70,8 +70,12 @@ public class AigcExecutorProperties {
|
||||
/** LLM 模型名(沿 C1 spike / M-a 批跑同款通道) */
|
||||
private String llmModel = "MiniMax-M2.7";
|
||||
|
||||
/** tick 轮询间隔毫秒(@Scheduled fixedDelay:上一 tick 跑完才调度下一 tick,进程级单飞) */
|
||||
private Long pollIntervalMs = 10000L;
|
||||
/**
|
||||
* tick 轮询间隔毫秒(@Scheduled fixedDelay:上一 tick 跑完才调度下一 tick,进程级单飞)。
|
||||
* C3:MQ 事件驱动为主触发后,tick 降级为兜底补偿,<b>周期由 10s 拉大至 60s</b>(不再作主触发路径);
|
||||
* 同时作 C3 兜底路「stale queued」判据的一个周期基准(createTime 早于「一个 poll 周期前」才重发信号)。
|
||||
*/
|
||||
private Long pollIntervalMs = 60000L;
|
||||
|
||||
/** 每 tick 认领上限(灰度小流量;扫描按 id 升序 FIFO 公平) */
|
||||
private Integer scanBatchSize = 2;
|
||||
|
||||
@ -4,6 +4,7 @@ import com.wanxiang.huijing.game.module.aigc.controller.admin.task.vo.DifyCallba
|
||||
import com.wanxiang.huijing.game.module.aigc.dal.dataobject.task.AigcTaskDO;
|
||||
import com.wanxiang.huijing.game.module.aigc.dal.mysql.task.AigcTaskMapper;
|
||||
import com.wanxiang.huijing.game.module.aigc.enums.FailureReasonEnum;
|
||||
import com.wanxiang.huijing.game.module.aigc.mq.GenTaskProducer;
|
||||
import com.wanxiang.huijing.game.module.aigc.service.callback.DifyCallbackService;
|
||||
import com.wanxiang.huijing.game.module.studio.api.SourceProjectApi;
|
||||
import com.wanxiang.huijing.game.module.studio.dto.SourceProjectFetchRespDTO;
|
||||
@ -42,6 +43,13 @@ import static com.wanxiang.huijing.game.module.aigc.enums.ErrorCodeConstants.AIG
|
||||
* {@link DifyCallbackService#handleCallback}(复用 EXEC-001 §8.4 唯一写入路径,不绕过、不复制三表逻辑)。
|
||||
* 从此真实用户 submitGenerate 不再永久停 queued(M2 整体收口第②段)。
|
||||
*
|
||||
* 【C3 触发面变更(切片一 阶段〇:RocketMQ 异步 gen 队列)】主触发从「@Scheduled DB 轮询」改为「MQ 事件驱动」:
|
||||
* 入队即发 AIGC_GEN_TOPIC 信号 → 有界消费者(consumeThreadMax=15=dispatch 投递握手并发速率上限,非在跑生成并发;
|
||||
* 真在跑并发≤15 由有界 worker 池保证=follow-up)→ 经
|
||||
* {@link #claimAndDispatchById} 走本类同一 CAS 认领 + 生成流水(复用,不重写)。{@link #tick()} @Scheduled
|
||||
* <b>降级为兜底补偿</b>(周期拉大至 60s):① stale queued(信号丢)重发信号交消费者、② stale running 僵死收尸
|
||||
* failed。派发失败落 failed 不留悬空 running({@link #claimAndDispatchById} catch 路)。
|
||||
*
|
||||
* 【调度选型依据(§3,亲核 2026-06-10)】Spring @Scheduled 薄轮询,不依赖 XXL-Job:
|
||||
* staging xxl.job.enabled=false 且无 xxl-admin → HuijingXxlJobAutoConfiguration 整类不装配 →
|
||||
* @XxlJob 触发面在 staging 永不运行(trade SettlementJob 即此死法,B4 烟测实证)。
|
||||
@ -107,6 +115,15 @@ public class AigcGenerateExecutor {
|
||||
private final SourceProjectApi sourceProjectApi;
|
||||
/** 墙钟(预算门用;单测注入假时钟验证 M5) */
|
||||
private final Clock clock;
|
||||
/**
|
||||
* gen 队列生产者(C3 兜底补偿路专用,<b>可为 null</b>):MQ 为主触发后,{@code @Scheduled} tick 降级为兜底补偿,
|
||||
* 对捞到的 stale queued <b>重发信号</b>(交消费者重走认领+派发),<b>不在 tick 线程直接派发</b>——这是让 MQ 保持
|
||||
* 唯一投递路径的<b>投递可靠性补偿</b>(信号丢了补一条、回到同一条消费→派发链路,不另起第二条并行派发路径),
|
||||
* 与并发 cap 无关(真「在跑生成并发≤15」属有界 worker 池=follow-up,不由消费线程数或此重发决定)。
|
||||
* <p><b>null 时降级为直派</b>(producer 未注入的旧构造/单测态 / 生产者 Bean 因故缺席):tick 退回既有
|
||||
* 「扫描-认领-派发」行为,避免任务卡死(韧性降级,非主路径)。生产经 AigcExecutorConfiguration 注入非空。
|
||||
*/
|
||||
private final GenTaskProducer genTaskProducer;
|
||||
|
||||
/** 自禁用连续 tick 计数(自检恢复即清零;告警节流判据) */
|
||||
private int selfCheckFailTicks = 0;
|
||||
@ -155,19 +172,34 @@ public class AigcGenerateExecutor {
|
||||
}
|
||||
|
||||
/**
|
||||
* 10 参构造(A11 切片三 M1② 收口;经 AigcExecutorConfiguration @Bean 调用):在 9 参基础上接
|
||||
* {@link SourceProjectApi} 反查 seam——modify 路据 baseVersionId 反查 base 源注入 §6.1 job 的 sourceProject 键
|
||||
* (供便宜档 HTTP worker 无 DB 消费)。按 {@code aigc.executor.dispatcher} 选派发通道——{@code saa} 且 saaDispatcher
|
||||
* 非空 → 进程内 SaaGraphDispatcher(形态①);否则 → http {@link WorkerDispatchClient}(默认,现行)。
|
||||
*
|
||||
* @param saaDispatcher SAA 进程内派发器(可空:dispatcher=http 或 saa 未装配时为 null → 回落 http)
|
||||
* @param sourceProjectApi 源项目反查 seam(A11 modify 取 base 源;可 null=取源旁路,不注入 sourceProject)
|
||||
* 10 参构造(A11 切片三 M1② 收口;现行单测沿用此签名):委托 11 参构造,{@code genTaskProducer} 传 null →
|
||||
* {@code @Scheduled} 兜底 tick 退回「直接认领+派发」(C3 前行为,非破坏性)。生产经 11 参构造注入非空生产者,
|
||||
* 兜底 tick 改为「重发信号交消费者」(投递可靠性补偿,让 MQ 保持唯一投递路径;与并发 cap 无关)。
|
||||
*/
|
||||
public AigcGenerateExecutor(AigcExecutorProperties properties, AigcTaskMapper aigcTaskMapper,
|
||||
DifyCallbackService difyCallbackService, PromptResourceLoader promptResourceLoader,
|
||||
GameConfigSchemaValidator schemaValidator, ExecutorLlmClient llmClient,
|
||||
WorkerDispatchClient workerDispatchClient, GenerationDispatcher saaDispatcher,
|
||||
SourceProjectApi sourceProjectApi, Clock clock) {
|
||||
// 委托 11 参构造,genTaskProducer 传 null → 兜底 tick 直派(C3 前行为,非破坏性;现行单测沿用此签名)。
|
||||
this(properties, aigcTaskMapper, difyCallbackService, promptResourceLoader, schemaValidator,
|
||||
llmClient, workerDispatchClient, saaDispatcher, sourceProjectApi, clock, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* 11 参构造(切片一 阶段〇 C3 收口;经 AigcExecutorConfiguration @Bean 调用):在 10 参基础上接 gen 队列生产者
|
||||
* {@link GenTaskProducer}——MQ 为主触发后,{@code @Scheduled} tick 降级为兜底补偿,对 stale queued <b>重发信号</b>
|
||||
* (交消费者重走认领+派发),不在 tick 线程直接派发(投递可靠性补偿,让 MQ 保持唯一投递路径;与并发 cap 无关)。
|
||||
*
|
||||
* @param saaDispatcher SAA 进程内派发器(可空:dispatcher=http 或 saa 未装配时为 null → 回落 http)
|
||||
* @param sourceProjectApi 源项目反查 seam(A11 modify 取 base 源;可 null=取源旁路,不注入 sourceProject)
|
||||
* @param genTaskProducer gen 队列生产者(C3 兜底 tick 重发信号用;可 null=兜底退回直派,非破坏性/单测态)
|
||||
*/
|
||||
public AigcGenerateExecutor(AigcExecutorProperties properties, AigcTaskMapper aigcTaskMapper,
|
||||
DifyCallbackService difyCallbackService, PromptResourceLoader promptResourceLoader,
|
||||
GameConfigSchemaValidator schemaValidator, ExecutorLlmClient llmClient,
|
||||
WorkerDispatchClient workerDispatchClient, GenerationDispatcher saaDispatcher,
|
||||
SourceProjectApi sourceProjectApi, Clock clock, GenTaskProducer genTaskProducer) {
|
||||
// §5.2 阈值关系铁律:stale×60 > budget + 273,装配期校验(误配置宁可启动失败也不带病收割在飞任务)
|
||||
properties.validateThresholds();
|
||||
this.properties = properties;
|
||||
@ -184,11 +216,14 @@ public class AigcGenerateExecutor {
|
||||
log.info("[executor] 生成派发通道选定 dispatcher={}(saa装配={})= {}", properties.getDispatcher(),
|
||||
saaDispatcher != null, useSaa ? "进程内SaaGraphDispatcher(形态①)" : "http WorkerDispatchClient(现行)");
|
||||
this.clock = clock;
|
||||
// C3:gen 队列生产者(兜底 tick 重发信号用);非空=MQ 主触发 + 兜底重发信号(投递可靠性补偿),null=兜底退回直派(单测/降级态)。
|
||||
this.genTaskProducer = genTaskProducer;
|
||||
// 启动三行自检日志(§12-① 冒烟观察面:执行器已启用 / prompt 资源版本 / key 配置态;严禁打印密钥本身)
|
||||
log.info("[executor-selfcheck] aigc 生成执行器已启用(poll={}ms, batch={}, budget={}s, stale={}min, maxAge={}h, templates={}, sourceSeam={})",
|
||||
log.info("[executor-selfcheck] aigc 生成执行器已启用(poll={}ms, batch={}, budget={}s, stale={}min, maxAge={}h, templates={}, sourceSeam={}, 触发={})",
|
||||
properties.getPollIntervalMs(), properties.getScanBatchSize(), properties.getTaskBudgetSeconds(),
|
||||
properties.getStaleRunningMinutes(), properties.getMaxTaskAgeHours(), properties.getSupportedTemplates(),
|
||||
sourceProjectApi != null ? "在席" : "缺席(modify取源旁路)");
|
||||
sourceProjectApi != null ? "在席" : "缺席(modify取源旁路)",
|
||||
genTaskProducer != null ? "MQ主触发+tick兜底重发信号" : "tick直派(无MQ生产者/单测态)");
|
||||
// HJ-MC-TPL-EXEC-001 §6.4:多模板装载后无单一版本,按 templateId 逐模板列出版本(就绪时)
|
||||
log.info("[executor-selfcheck] prompt 资源版本 {}", promptResourceLoader.isReady()
|
||||
? promptResourceLoader.describePromptVersions() : ("未就绪:" + promptResourceLoader.getNotReadyReason()));
|
||||
@ -196,11 +231,13 @@ public class AigcGenerateExecutor {
|
||||
}
|
||||
|
||||
/**
|
||||
* 调度触发面(§5.1):fixedDelay 轮询(默认 10s,上一 tick 跑完才开下一 tick)+ 启动延迟 30s
|
||||
* (避开 app 启动期 Bean 装配/Flyway 迁移窗口)。XXL-Job 上线后本方法被 @XxlJob 触发面替换,
|
||||
* 执行体 {@link #runOneTick()} 原样复用(§3 演化路径)。
|
||||
* 调度触发面(§5.1;C3 降级为兜底补偿):fixedDelay 轮询(<b>默认 60s,周期已拉大</b>——MQ 事件驱动为主触发后,
|
||||
* 本 tick 不再作主触发路径,只当兜底补偿)+ 启动延迟 30s(避开 app 启动期 Bean 装配/Flyway 迁移窗口)。
|
||||
* <p>兜底补偿捞两类卡住任务({@link #runOneTick()}):① stale queued(入队了但 MQ 信号丢)→ 有生产者则
|
||||
* 重发信号交有界消费者认领、无生产者则直派兜底;② stale running 超阈值(进程崩溃/派发丢状态遗留)→ watchdog
|
||||
* 收尸 failed。XXL-Job 上线后本方法被 @XxlJob 触发面替换,执行体原样复用(§3 演化路径)。
|
||||
*/
|
||||
@Scheduled(fixedDelayString = "${aigc.executor.poll-interval-ms:10000}", initialDelay = 30_000)
|
||||
@Scheduled(fixedDelayString = "${aigc.executor.poll-interval-ms:60000}", initialDelay = 30_000)
|
||||
public void tick() {
|
||||
try {
|
||||
runOneTick();
|
||||
@ -225,8 +262,28 @@ public class AigcGenerateExecutor {
|
||||
int claimed = 0;
|
||||
int succeeded = 0;
|
||||
int failed = 0;
|
||||
int resignaled = 0; // C3 兜底补偿:对 stale queued 重发 MQ 信号的条数(有生产者的生产态走此路,不直派)
|
||||
// C3 stale queued 判据:createTime 早于「一个 poll 周期前」才算信号大概率已丢/滞留,才重发;
|
||||
// 新鲜 queued(刚入队、信号在途)不重发以减噪(重发即便发生也被消费侧 CAS 幂等去重,无害)。
|
||||
LocalDateTime staleQueuedBefore = LocalDateTime.now().minus(Duration.ofMillis(properties.getPollIntervalMs()));
|
||||
// ===== 步骤③:逐任务处理(tick 线程内串行;单任务异常不中断本 tick 其余任务)=====
|
||||
for (AigcTaskDO task : candidates) {
|
||||
// ===== C3 兜底补偿路(生产态:MQ 生产者在席)=====
|
||||
// 只对 stale queued「重发信号」交消费者重走认领+派发,不在 tick 线程直接派发——这是让 MQ 保持唯一投递
|
||||
// 路径的投递可靠性补偿(信号丢了补发一条、回到同一条消费→派发链路,不另起第二条并行派发路径);
|
||||
// 与并发 cap 无关(真「在跑生成并发≤15」属有界 worker 池=follow-up,不由消费线程数或此重发决定)。
|
||||
if (genTaskProducer != null) {
|
||||
if (task.getCreateTime() != null && task.getCreateTime().isAfter(staleQueuedBefore)) {
|
||||
continue; // 新鲜 queued(信号大概率在途),本轮不重发,减噪
|
||||
}
|
||||
// best-effort 重发(生产者内部吞异常不外抛);消费侧 CAS 幂等去重,重复信号无害。
|
||||
genTaskProducer.publishNewTask(task.getId());
|
||||
resignaled++;
|
||||
log.info("[executor-tick] 兜底补偿:stale queued 重发 gen 信号(交消费者认领,不在 tick 直派)taskId={}, traceId={}",
|
||||
task.getId(), task.getTraceId());
|
||||
continue;
|
||||
}
|
||||
// ===== 无 MQ 生产者(单测/降级态):退回既有「直接认领+派发」(行为逐字不变,兜底不让任务卡死)=====
|
||||
try {
|
||||
// 系统身份 + 所属租户上下文内处理整任务(含 CAS 认领与回调写链——三表行必须落对租户,§5.1-3b;
|
||||
// 系统身份保证审计列 updater 可填充,否则 UPDATE null 直出撞 NOT NULL,部署窗口#3 根因修复)
|
||||
@ -249,16 +306,78 @@ public class AigcGenerateExecutor {
|
||||
task.getId(), task.getTraceId(), e);
|
||||
}
|
||||
}
|
||||
// ===== 步骤④:stale watchdog(进程崩溃恢复:上次认领后死掉的任务收尸为 failed(timeout))=====
|
||||
// ===== 步骤④:stale watchdog(进程崩溃恢复 / 派发丢状态遗留:RUNNING 超阈值收尸为 failed(timeout))=====
|
||||
// 这即 C3 兜底补偿的第②类「stale running 僵死任务」——认领成功后进程崩溃/派发失败未落状态的漏网,超时收尸。
|
||||
int reaped = watchdogSweep();
|
||||
// tick 级汇总行(§5.5):有活动打 info;空转降为 debug 防 10s 级刷屏
|
||||
if (claimed + reaped > 0) {
|
||||
log.info("[executor-tick] tick 汇总:认领 {} / 成功 {} / 失败 {} / 收尸 {}", claimed, succeeded, failed, reaped);
|
||||
// tick 级汇总行(§5.5):有活动打 info;空转降为 debug 防刷屏
|
||||
if (claimed + reaped + resignaled > 0) {
|
||||
log.info("[executor-tick] tick 汇总(兜底补偿):认领 {} / 成功 {} / 失败 {} / 重发信号 {} / 收尸 {}",
|
||||
claimed, succeeded, failed, resignaled, reaped);
|
||||
} else {
|
||||
log.debug("[executor-tick] tick 汇总:认领 0 / 成功 0 / 失败 0 / 收尸 0(空转)");
|
||||
log.debug("[executor-tick] tick 汇总:认领 0 / 成功 0 / 失败 0 / 重发信号 0 / 收尸 0(空转)");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* MQ 事件驱动的单任务认领+派发(切片一 阶段〇 C3)——RocketMQ 消费者 {@link com.wanxiang.huijing.game.module.aigc.mq.GenTaskConsumer}
|
||||
* 收到 taskId 信号后经 {@link com.wanxiang.huijing.game.module.aigc.mq.AigcGenDispatchService} 调用本方法。与 @Scheduled tick 的
|
||||
* {@link #processCandidate} 同一执行体(CAS 认领 + {@link #runGenerationPipeline} 派发)、同一系统身份 + 租户上下文,
|
||||
* 只是触发面从「批量扫描」换成「单 taskId 事件」。
|
||||
*
|
||||
* <p><b>幂等(CAS 去重)</b>:既有 {@code claimQueuedTask} 仅 queued→running 才成功;认领成功才派发,认领失败
|
||||
* (被并发消费线程/兜底 tick 抢走、非 queued、已取消、任务不存在)即幂等丢弃返 false——同一 taskId 重复投递不双跑。
|
||||
*
|
||||
* <p><b>派发失败落 failed(不留悬空 running,Opus 评审)</b>:CAS 成功后 {@link #runGenerationPipeline} 抛异常
|
||||
* (派发/渲染/DB 抖动等)→ catch 里经 {@link #callbackFailed} 落 failed(llm_error 归因,与 dispatchGeneric 投递
|
||||
* 失败同桶),绝不留悬空 running 等 watchdog。若 pipeline 内部已自行落终态后再抛,二次回调被终态拒重入(M7)幂等吞掉,安全。
|
||||
*
|
||||
* <p><b>就绪门</b>:key/prompt/schema 未就绪({@link #isExecutorReady()})则不认领、返 false——任务留 queued,
|
||||
* 由 @Scheduled 兜底稍后重发信号重试(不烧一次认领把任务判失败)。
|
||||
*
|
||||
* @param taskId 任务 ID(消息体轻信号)
|
||||
* @return true=CAS 认领成功(已派发或已落终态);false=认领未命中/未就绪=幂等丢弃(不重投)
|
||||
*/
|
||||
public boolean claimAndDispatchById(Long taskId) {
|
||||
// 消费线程无租户上下文 → executeIgnore 跨租户按 id 取任务(与 tick 扫描同款绕租户过滤)。
|
||||
AigcTaskDO task = TenantUtils.executeIgnore(() -> aigcTaskMapper.selectById(taskId));
|
||||
if (task == null) {
|
||||
// 任务不存在(脏信号/已被物理清理):幂等丢弃,不重投。
|
||||
log.warn("[executor-mq] 任务不存在,幂等丢弃 taskId={}", taskId);
|
||||
return false;
|
||||
}
|
||||
// 就绪门(只读、不变更状态,多消费线程并发安全;不复用 selfCheckPass 以免竞争其节流计数):
|
||||
// 未就绪则不认领,任务留 queued 交兜底轮询稍后重发信号重试。
|
||||
if (!isExecutorReady()) {
|
||||
log.warn("[executor-mq] 执行器未就绪(api-key/prompt/schema),暂不认领(兜底轮询稍后重试)taskId={}, traceId={}",
|
||||
task.getId(), task.getTraceId());
|
||||
return false;
|
||||
}
|
||||
// 系统身份 + 所属租户上下文内:CAS 认领 → 复用生成流水派发;认领失败幂等丢弃,派发抛异常落 failed(不留悬空 running)。
|
||||
Boolean claimed = executeWithSystemIdentity(task, () -> {
|
||||
if (!aigcTaskMapper.claimQueuedTask(task.getId())) {
|
||||
// CAS 未命中:被并发消费线程/兜底 tick 抢走,或已取消/非 queued —— 幂等丢弃。
|
||||
log.info("[executor-mq] CAS 认领未命中(已被认领/非 queued),幂等丢弃 taskId={}, traceId={}",
|
||||
task.getId(), task.getTraceId());
|
||||
return Boolean.FALSE;
|
||||
}
|
||||
long queuedSeconds = task.getCreateTime() == null ? -1
|
||||
: Duration.between(task.getCreateTime(), LocalDateTime.now()).getSeconds();
|
||||
log.info("[executor-mq] 任务认领成功(MQ 触发)taskId={}, traceId={}, templateId={}, 排队时长={}s",
|
||||
task.getId(), task.getTraceId(), task.getTemplateId(), queuedSeconds);
|
||||
try {
|
||||
// 复用既有生成流水(与 tick processCandidate 同一 runGenerationPipeline:generic 走派发面 SaaGraphDispatcher/http worker)。
|
||||
runGenerationPipeline(task);
|
||||
} catch (Exception e) {
|
||||
// 认领成功但派发抛异常:落 failed(llm_error,与 dispatchGeneric 投递失败同桶),绝不留悬空 running(Opus 评审)。
|
||||
log.error("[executor-mq] 认领成功但派发异常,落 failed(不留悬空 running)taskId={}, traceId={}",
|
||||
task.getId(), task.getTraceId(), e);
|
||||
callbackFailed(task, FailureReasonEnum.LLM_ERROR);
|
||||
}
|
||||
return Boolean.TRUE;
|
||||
});
|
||||
return Boolean.TRUE.equals(claimed);
|
||||
}
|
||||
|
||||
/** 单测可观测口径:自禁用告警已发条数(验证 §8.3 节流正确性) */
|
||||
public int getSelfCheckAlertCount() {
|
||||
return selfCheckAlertCount;
|
||||
@ -294,6 +413,19 @@ public class AigcGenerateExecutor {
|
||||
return false;
|
||||
}
|
||||
|
||||
/**
|
||||
* 就绪只读检查(C3 MQ 路专用):key 非空 + prompt/schema classpath 资源就位。
|
||||
* 与 {@link #selfCheckPass()} 判据同源,但<b>不变更任何状态</b>(不碰 selfCheckFailTicks 节流计数)——
|
||||
* 供多消费线程(consumeThreadMax=15)并发安全调用(selfCheckPass 的节流计数是 tick 单线程专用,多线程并发会数据竞争)。
|
||||
*
|
||||
* @return true=就绪可认领
|
||||
*/
|
||||
private boolean isExecutorReady() {
|
||||
return StringUtils.hasText(properties.getApiKey())
|
||||
&& promptResourceLoader.isReady()
|
||||
&& schemaValidator.isReady();
|
||||
}
|
||||
|
||||
// ============================== 系统身份任务边界(部署窗口#3 根因修复) ==============================
|
||||
|
||||
/** 执行器系统用户编号:审计列 updater/creator 统一落 "0"(DefaultDBFieldHandler 以 getLoginUserId().toString() 填充) */
|
||||
|
||||
@ -9,6 +9,7 @@ import com.wanxiang.huijing.game.module.aigc.dal.mysql.task.AigcTaskMapper;
|
||||
import com.wanxiang.huijing.game.module.aigc.enums.AigcTaskSourceEnum;
|
||||
import com.wanxiang.huijing.game.module.aigc.enums.AigcTaskStatusEnum;
|
||||
import com.wanxiang.huijing.game.module.aigc.enums.FailureReasonEnum;
|
||||
import com.wanxiang.huijing.game.module.aigc.mq.GenTaskProducer;
|
||||
import com.wanxiang.huijing.game.module.aigc.service.executor.AigcTemplateConstants;
|
||||
import com.wanxiang.huijing.framework.common.pojo.PageResult;
|
||||
import com.wanxiang.huijing.module.infra.api.config.ConfigApi;
|
||||
@ -81,11 +82,20 @@ public class AigcTaskServiceImpl implements AigcTaskService {
|
||||
* 入口 QPS 突发闸(Sentinel FLOW_GRADE_QPS,规则热源 Nacos;切片一 阶段〇 C2)。
|
||||
* 与上方控制平面软注入件不同:GenAdmissionResource 是无条件 @Component(恒在席),故直接 @Resource 硬注入
|
||||
* (同 aigcTaskMapper/playerApi),不走 ObjectProvider。职责 = 入队入口速率闸削洪峰,与 DB 三门(业务配额)
|
||||
* / C3 consumeThreadMax(在跑并发≤15)三层正交,详见 {@link GenAdmissionResource} 类注释。
|
||||
* / C3 consumeThreadMax(dispatch 投递速率上限;真在跑并发≤15 属有界 worker 池=follow-up)三层正交,详见 {@link GenAdmissionResource} 类注释。
|
||||
*/
|
||||
@Resource
|
||||
private com.wanxiang.huijing.game.module.aigc.admission.GenAdmissionResource genAdmissionResource;
|
||||
|
||||
/**
|
||||
* gen 队列生产者软注入(切片一 阶段〇 C3):入队成功后发 AIGC_GEN_TOPIC 信号,触发有界消费者(consumeThreadMax=15)
|
||||
* 据 taskId CAS 认领拉起生成——把「触发」从执行器 @Scheduled DB 轮询改为 MQ 事件驱动。
|
||||
* 与执行器/消费者同 {@code aigc.executor.enabled} 条件装配 → 用 {@link ObjectProvider} 软注入:enabled=false
|
||||
* (无执行器/无消费者)时缺席,入队照常落库、仅不发信号(现行骨架行为,任务留 queued)。best-effort:发信号失败不回滚入队。
|
||||
*/
|
||||
@Resource
|
||||
private ObjectProvider<GenTaskProducer> genTaskProducerProvider;
|
||||
|
||||
/**
|
||||
* v0 会员档位占位值 = L1(执行版 §5.1/B3):AigcGenerateReqVO 无 level 字段、会员档数据模型未建,
|
||||
* 提交侧统一落 L1,避免为未建会员体系造孤儿字段。分档逻辑 v0 即可用(§8 手塞 level=2 读 L2 档单测覆盖)。
|
||||
@ -197,9 +207,11 @@ public class AigcTaskServiceImpl implements AigcTaskService {
|
||||
// D12 控制平面 + GP9(共享入队门,§6.1):降级门→配额并发→背压→GP9 safety→insert(enabled=false 全门旁路 = 现行行为)
|
||||
enqueueWithControlPlane(task);
|
||||
|
||||
// TODO 对接点:投递 RocketMQ 异步生成队列(T-AGC-07),消费者拉起 Dify 工作流(契约#6)。
|
||||
// 投递须在事务提交后执行(避免「消息已发、库未提交」),并保证投递失败可补偿(落消息表 / 重试)。
|
||||
// LLM 实名充值为人工闸门,本处不触发计费。当前骨架仅落库 queued + 返回 taskId,不同步等待结果。
|
||||
// C3:投递 RocketMQ 异步生成队列信号(落地 T-AGC-07 对接点)——入队落库后发 AIGC_GEN_TOPIC 信号,有界消费者
|
||||
// (consumeThreadMax=15)据 taskId CAS 认领拉起生成。submitGenerate 非 @Transactional,enqueueWithControlPlane
|
||||
// 的 insert 已自动提交(此处即「提交后」,规避「消息已发、库未提交」);best-effort:发信号失败不回滚入队,由执行器兜底轮询补偿。
|
||||
// LLM 实名充值为人工闸门,本处不触发计费。
|
||||
publishGenSignal(task);
|
||||
log.info("[submitGenerate] 生成任务已入队 taskId={} traceId={} userId={} templateId={} level={}",
|
||||
task.getId(), task.getTraceId(), userId, reqVO.getTemplateId(), task.getLevel());
|
||||
return task;
|
||||
@ -254,7 +266,8 @@ public class AigcTaskServiceImpl implements AigcTaskService {
|
||||
// 决策D:retry 也走控制平面(重跑全门 + GP9 safety)——违规/超配额请求无法靠 retry 绕过门禁
|
||||
enqueueWithControlPlane(retry);
|
||||
|
||||
// TODO 对接点:同 submitGenerate,投递 RocketMQ 异步生成队列拉起 Dify 工作流(契约#6)。
|
||||
// C3:同 submitGenerate,入队落库后发 gen 队列信号触发有界消费者认领(best-effort,失败由执行器兜底轮询补偿)。
|
||||
publishGenSignal(retry);
|
||||
log.info("[retryTask] 重试任务已入队 newTaskId={} retryOf={} traceId={} userId={} level={}",
|
||||
retry.getId(), origin.getId(), retry.getTraceId(), userId, retry.getLevel());
|
||||
return retry;
|
||||
@ -322,7 +335,7 @@ public class AigcTaskServiceImpl implements AigcTaskService {
|
||||
private void enqueueWithControlPlane(AigcTaskDO task) {
|
||||
// 入口 QPS 突发闸(Sentinel FLOW_GRADE_QPS,规则热源 Nacos dataId sentinel-gen-flow-rules)——先于下方全部 DB 门:
|
||||
// 统计「每秒入队请求数」削洪峰,越限由 blockHandler 优雅转 AIGC_ADMISSION_RATE_LIMITED(非 500、不静默丢)。
|
||||
// 与下方 DB 三门(per-creator 业务配额)分层正交;「在跑生成并发≤15」由 C3 RocketMQ consumeThreadMax 保证、不在此。
|
||||
// 与下方 DB 三门(per-creator 业务配额)分层正交;真「在跑生成并发≤15」由有界 worker 池保证(follow-up,C3 consumeThreadMax 只 cap dispatch 投递速率)、不在此。
|
||||
// 置于 props 判空之前 = 独立于控制平面 enabled 开关(基础设施级速率保护恒生效);无 Nacos 规则时为纳秒级 no-op,
|
||||
// 其回滚杠杆是 Nacos 流控规则本身(调大 count / 删 dataId),与控制平面 enabled=false 各自独立。
|
||||
genAdmissionResource.admit();
|
||||
@ -379,6 +392,26 @@ public class AigcTaskServiceImpl implements AigcTaskService {
|
||||
creatorUserId, level, task.getId(), dailyCount, quota.getDaily());
|
||||
}
|
||||
|
||||
/**
|
||||
* 入队后发 gen 队列信号(切片一 阶段〇 C3;submitGenerate/retryTask 共用)——把「触发」从执行器 @Scheduled DB 轮询
|
||||
* 改为 RocketMQ 事件驱动。<b>双重兜底</b>:① 软注入——生产者缺席({@code aigc.executor.enabled=false},无执行器/
|
||||
* 无消费者)时跳过、任务留 queued(现行骨架行为);② best-effort——发信号失败由 {@link GenTaskProducer#publishNewTask}
|
||||
* 内部吞异常并告警留痕、不回滚入队,执行器 @Scheduled 兜底轮询对 stale queued 重发信号补偿。
|
||||
* <p>须在入队落库后调用({@code task.getId()} 已由 insert 回填)。
|
||||
*
|
||||
* @param task 已落库的 queued 任务(id 已回填)
|
||||
*/
|
||||
private void publishGenSignal(AigcTaskDO task) {
|
||||
GenTaskProducer producer = genTaskProducerProvider.getIfAvailable();
|
||||
if (producer == null) {
|
||||
// 执行器/消费者未装配(aigc.executor.enabled=false):无消费方,发信号无意义,跳过;任务留 queued(现行骨架行为)。
|
||||
log.debug("[aigc-mq] gen 队列生产者缺席(executor 未启用),跳过发信号 taskId={}", task.getId());
|
||||
return;
|
||||
}
|
||||
// best-effort:publishNewTask 内部已吞异常(发失败只告警、不外抛、不回滚入队),执行器兜底轮询补偿丢失信号。
|
||||
producer.publishNewTask(task.getId());
|
||||
}
|
||||
|
||||
/**
|
||||
* 降级门读取(执行版 §6.1):读 infra ConfigApi 键 aigc.generate.paused。
|
||||
* fail-open on switch read:ConfigApi 缺席 / 读不到值 / 非 "true" / 异常 → 一律视为「未暂停」,
|
||||
|
||||
@ -0,0 +1,53 @@
|
||||
package com.wanxiang.huijing.game.module.aigc.mq;
|
||||
|
||||
import com.wanxiang.huijing.framework.test.core.ut.BaseMockitoUnitTest;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.mockito.InjectMocks;
|
||||
import org.mockito.Mock;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
|
||||
import static org.mockito.ArgumentMatchers.anyLong;
|
||||
import static org.mockito.Mockito.never;
|
||||
import static org.mockito.Mockito.verify;
|
||||
import static org.mockito.Mockito.when;
|
||||
|
||||
/**
|
||||
* {@link GenTaskConsumer} 单元测试(纯 Mockito,不触真执行器/DB)——把守 C3 消费者「幂等丢弃 + 不重投」契约。
|
||||
*
|
||||
* <p>覆盖:① {@code claimAndDispatch} 返 false(已被并发消费线程/兜底 tick 认领、非 queued)→ {@link GenTaskConsumer#onMessage}
|
||||
* 正常返回、<b>不抛</b>(幂等丢弃,天然去重不双跑);② 返 true(认领成功)→ onMessage 不抛、taskId 正确解析下派;
|
||||
* ③ 毒消息(taskId 非数字)→ onMessage 不抛且<b>不调</b> dispatch(丢弃不重投,防重投风暴)。
|
||||
*
|
||||
* @author 绘境AI
|
||||
*/
|
||||
class GenTaskConsumerTest extends BaseMockitoUnitTest {
|
||||
|
||||
@Mock
|
||||
private AigcGenDispatchService dispatchService;
|
||||
|
||||
@InjectMocks
|
||||
private GenTaskConsumer consumer;
|
||||
|
||||
/** 幂等丢弃:claimAndDispatch 返 false → onMessage 正常返回、不抛(同 taskId 重复投递不双跑)。 */
|
||||
@Test
|
||||
void testOnMessage_claimMiss_idempotentDiscardNoThrow() {
|
||||
when(dispatchService.claimAndDispatch(789L)).thenReturn(false);
|
||||
assertDoesNotThrow(() -> consumer.onMessage("789"));
|
||||
verify(dispatchService).claimAndDispatch(789L);
|
||||
}
|
||||
|
||||
/** 认领成功:claimAndDispatch 返 true → onMessage 不抛,taskId 正确解析下派。 */
|
||||
@Test
|
||||
void testOnMessage_claimed_dispatchInvokedNoThrow() {
|
||||
when(dispatchService.claimAndDispatch(100L)).thenReturn(true);
|
||||
assertDoesNotThrow(() -> consumer.onMessage("100"));
|
||||
verify(dispatchService).claimAndDispatch(100L);
|
||||
}
|
||||
|
||||
/** 毒消息:taskId 非数字 → onMessage 不抛且不调 dispatch(丢弃不重投,防无限重投风暴)。 */
|
||||
@Test
|
||||
void testOnMessage_malformedTaskId_discardNoDispatchNoThrow() {
|
||||
assertDoesNotThrow(() -> consumer.onMessage("not-a-number"));
|
||||
verify(dispatchService, never()).claimAndDispatch(anyLong());
|
||||
}
|
||||
}
|
||||
@ -0,0 +1,51 @@
|
||||
package com.wanxiang.huijing.game.module.aigc.mq;
|
||||
|
||||
import com.wanxiang.huijing.framework.test.core.ut.BaseMockitoUnitTest;
|
||||
import org.apache.rocketmq.spring.core.RocketMQTemplate;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.mockito.InjectMocks;
|
||||
import org.mockito.Mock;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
|
||||
import static org.mockito.ArgumentMatchers.anyString;
|
||||
import static org.mockito.ArgumentMatchers.eq;
|
||||
import static org.mockito.Mockito.doThrow;
|
||||
import static org.mockito.Mockito.verify;
|
||||
|
||||
/**
|
||||
* {@link GenTaskProducer} 单元测试(纯 Mockito,不触真 MQ)——把守 C3 生产者 best-effort 契约。
|
||||
*
|
||||
* <p>覆盖:① 正常发 AIGC_GEN_TOPIC 信号、payload = taskId 字符串;② {@code syncSend} 抛异常时
|
||||
* {@link GenTaskProducer#publishNewTask} <b>不外抛</b>(best-effort:不回滚入队、由执行器兜底轮询补偿)。
|
||||
*
|
||||
* <p>注:{@code RocketMQTemplate.syncSend} 有 (String,Object) 与 (String,Message) 两个 2 参重载;本测用
|
||||
* String payload 匹配器({@code eq("123")}/{@code anyString()})强制解析到 (String,Object) 重载,与生产代码
|
||||
* {@code syncSend(TOPIC, String.valueOf(taskId))} 同一重载。
|
||||
*
|
||||
* @author 绘境AI
|
||||
*/
|
||||
class GenTaskProducerTest extends BaseMockitoUnitTest {
|
||||
|
||||
@Mock
|
||||
private RocketMQTemplate rocketMQTemplate;
|
||||
|
||||
@InjectMocks
|
||||
private GenTaskProducer producer;
|
||||
|
||||
/** 正常路:发到正确 topic,payload = taskId 字符串。 */
|
||||
@Test
|
||||
void testPublishNewTask_sendsSignalToTopicWithTaskId() {
|
||||
assertDoesNotThrow(() -> producer.publishNewTask(123L));
|
||||
// 验证发到 AIGC_GEN_TOPIC + payload="123"(未 stub 返回值=默认 null,publishNewTask 的 result!=null 守卫吞掉,不 NPE)
|
||||
verify(rocketMQTemplate).syncSend(eq(GenTaskProducer.TOPIC), eq("123"));
|
||||
}
|
||||
|
||||
/** best-effort(核心断言):syncSend 抛异常 → publishNewTask 吞异常不外抛(否则会连累入队事务/接口)。 */
|
||||
@Test
|
||||
void testPublishNewTask_bestEffort_swallowsSendException() {
|
||||
doThrow(new RuntimeException("MQ down")).when(rocketMQTemplate).syncSend(anyString(), anyString());
|
||||
// 关键:MQ 抖动不外抛(任务已落库 queued,丢失信号由执行器 @Scheduled 兜底轮询补偿)
|
||||
assertDoesNotThrow(() -> producer.publishNewTask(456L));
|
||||
verify(rocketMQTemplate).syncSend(eq(GenTaskProducer.TOPIC), eq("456"));
|
||||
}
|
||||
}
|
||||
@ -4,6 +4,7 @@ import com.wanxiang.huijing.game.module.aigc.controller.admin.task.vo.DifyCallba
|
||||
import com.wanxiang.huijing.game.module.aigc.dal.dataobject.task.AigcTaskDO;
|
||||
import com.wanxiang.huijing.game.module.aigc.dal.mysql.task.AigcTaskMapper;
|
||||
import com.wanxiang.huijing.game.module.aigc.framework.executor.config.AigcExecutorConfiguration;
|
||||
import com.wanxiang.huijing.game.module.aigc.mq.GenTaskProducer;
|
||||
import com.wanxiang.huijing.game.module.aigc.service.callback.DifyCallbackService;
|
||||
import com.wanxiang.huijing.game.module.studio.api.SourceProjectApi;
|
||||
import com.wanxiang.huijing.game.module.studio.dto.SourceProjectFetchRespDTO;
|
||||
@ -39,6 +40,7 @@ import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
import static org.mockito.ArgumentMatchers.any;
|
||||
import static org.mockito.ArgumentMatchers.anyInt;
|
||||
import static org.mockito.ArgumentMatchers.anyLong;
|
||||
import static org.mockito.ArgumentMatchers.anyString;
|
||||
import static org.mockito.ArgumentMatchers.eq;
|
||||
import static org.mockito.Mockito.mock;
|
||||
@ -832,4 +834,77 @@ class AigcGenerateExecutorTest extends BaseMockitoUnitTest {
|
||||
verifyNoInteractions(sourceProjectApi);
|
||||
}
|
||||
|
||||
// ============================== 用例27-29:C3 MQ 事件驱动认领+派发 & 兜底重发信号 ==============================
|
||||
|
||||
/**
|
||||
* C3-A 幂等({@link AigcGenerateExecutor#claimAndDispatchById}):CAS 认领未命中(被并发消费线程/兜底 tick 抢走、
|
||||
* 非 queued)→ 返 false,且不派发、不回调(同 taskId 重复投递不双跑=天然去重)。
|
||||
*/
|
||||
@Test
|
||||
void testClaimAndDispatchById_casMiss_idempotentNoDispatch() {
|
||||
AigcTaskDO task = queuedTask(901L, "aigc-mq-miss");
|
||||
when(aigcTaskMapper.selectById(901L)).thenReturn(task);
|
||||
when(aigcTaskMapper.claimQueuedTask(901L)).thenReturn(false); // CAS 未命中(已被认领/非 queued)
|
||||
|
||||
boolean claimed = newExecutor(readyProps(), Clock.systemDefaultZone()).claimAndDispatchById(901L);
|
||||
|
||||
assertFalse(claimed); // 幂等丢弃
|
||||
verifyNoInteractions(difyCallbackService); // 未派发、未回调(无终态写入)
|
||||
}
|
||||
|
||||
/**
|
||||
* C3-B 派发失败落 failed 不留悬空 running(Opus 评审命门):CAS 认领成功后派发通道抛异常 →
|
||||
* {@link AigcGenerateExecutor#claimAndDispatchById} catch 里经唯一写入路径落 failed(llm_error),绝不留悬空 running。
|
||||
*/
|
||||
@Test
|
||||
void testClaimAndDispatchById_dispatchThrows_marksFailedNoOrphanRunning() {
|
||||
AigcExecutorProperties props = genericProps();
|
||||
props.setDispatcher("saa"); // 选进程内派发通道 = 下方注入的抛异常 dispatcher
|
||||
ExecutorLlmClient llmClient = new ExecutorLlmClient(props, sender, millis -> { });
|
||||
WorkerDispatchClient wdc = new WorkerDispatchClient(props, dispatchSender);
|
||||
GenerationDispatcher boom = mock(GenerationDispatcher.class);
|
||||
when(boom.dispatch(any())).thenThrow(new RuntimeException("dispatch boom")); // 派发抛异常(非返 false 的前置失败)
|
||||
// 11 参构造:saaDispatcher=boom + dispatcher=saa → generationDispatcher=boom;producer=null(走 MQ 单任务路,非 tick)。
|
||||
AigcGenerateExecutor exec = new AigcGenerateExecutor(props, aigcTaskMapper, difyCallbackService,
|
||||
LOADER, VALIDATOR, llmClient, wdc, boom, null, Clock.systemDefaultZone(), null);
|
||||
AigcTaskDO task = genericQueuedTask(902L, "aigc-mq-throw");
|
||||
when(aigcTaskMapper.selectById(902L)).thenReturn(task);
|
||||
when(aigcTaskMapper.claimQueuedTask(902L)).thenReturn(true); // CAS 认领成功
|
||||
|
||||
boolean claimed = exec.claimAndDispatchById(902L);
|
||||
|
||||
assertTrue(claimed); // 已认领(虽派发失败),不重投
|
||||
// 关键:落 failed(llm_error) 走唯一写入路径 handleCallback,绝不留悬空 running 等 watchdog
|
||||
ArgumentCaptor<DifyCallbackReqVO> cap = ArgumentCaptor.forClass(DifyCallbackReqVO.class);
|
||||
verify(difyCallbackService).handleCallback(cap.capture());
|
||||
assertEquals("failed", cap.getValue().getStatus());
|
||||
assertEquals("llm_error", cap.getValue().getFailureReason());
|
||||
}
|
||||
|
||||
/**
|
||||
* C3-C 兜底补偿:@Scheduled tick 有 MQ 生产者(生产态)时对 stale queued「重发信号」交消费者重走认领+派发,
|
||||
* <b>不在 tick 线程直接派发/认领</b>——投递可靠性补偿(让 MQ 保持唯一投递路径),与并发 cap 无关
|
||||
* (真「在跑生成并发≤15」属有界 worker 池=follow-up,不由消费线程数或此重发决定)。
|
||||
*/
|
||||
@Test
|
||||
void testRunOneTick_withProducer_resignalsStaleQueuedNoDirectDispatch() {
|
||||
GenTaskProducer producer = mock(GenTaskProducer.class);
|
||||
AigcExecutorProperties props = readyProps();
|
||||
ExecutorLlmClient llmClient = new ExecutorLlmClient(props, sender, millis -> { });
|
||||
WorkerDispatchClient wdc = new WorkerDispatchClient(props, dispatchSender);
|
||||
// 11 参构造:producer 非空 → 兜底 tick 走「重发信号」路。
|
||||
AigcGenerateExecutor exec = new AigcGenerateExecutor(props, aigcTaskMapper, difyCallbackService,
|
||||
LOADER, VALIDATOR, llmClient, wdc, null, null, Clock.systemDefaultZone(), producer);
|
||||
AigcTaskDO stale = queuedTask(903L, "aigc-mq-resig");
|
||||
stale.setCreateTime(LocalDateTime.now().minusSeconds(180)); // 早于一个 poll 周期(默认 60s)前 → stale → 重发
|
||||
when(aigcTaskMapper.selectClaimableQueued(anyInt(), any(LocalDateTime.class))).thenReturn(List.of(stale));
|
||||
when(aigcTaskMapper.selectStaleRunning(anyInt(), any(LocalDateTime.class))).thenReturn(List.of());
|
||||
|
||||
exec.runOneTick();
|
||||
|
||||
verify(producer).publishNewTask(903L); // 重发信号交有界消费者
|
||||
verify(aigcTaskMapper, never()).claimQueuedTask(anyLong()); // 兜底不在 tick 直接认领
|
||||
verifyNoInteractions(difyCallbackService); // 兜底不直派、不回调
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@ -5,6 +5,7 @@ import com.wanxiang.huijing.game.module.aigc.controller.app.task.vo.TemplateResp
|
||||
import com.wanxiang.huijing.game.module.aigc.dal.dataobject.task.AigcTaskDO;
|
||||
import com.wanxiang.huijing.game.module.aigc.dal.mysql.task.AigcTaskMapper;
|
||||
import com.wanxiang.huijing.game.module.aigc.enums.AigcTaskStatusEnum;
|
||||
import com.wanxiang.huijing.game.module.aigc.mq.GenTaskProducer;
|
||||
import com.wanxiang.huijing.game.module.aigc.service.executor.AigcTemplateConstants;
|
||||
import com.wanxiang.huijing.framework.common.exception.ServiceException;
|
||||
import com.wanxiang.huijing.framework.common.pojo.CommonResult;
|
||||
@ -58,6 +59,10 @@ class AigcTaskServiceImplTest extends BaseMockitoUnitTest {
|
||||
// C2 入口 QPS 闸(硬注入 @Component):@InjectMocks 注入本 mock,void admit() 默认 no-op → 现行提交/重试基线行为不变。
|
||||
@Mock
|
||||
private com.wanxiang.huijing.game.module.aigc.admission.GenAdmissionResource genAdmissionResource;
|
||||
// C3 gen 队列生产者软注入:@InjectMocks 注入本 mock,getIfAvailable() 默认返 null(@Mock 默认)→ publishGenSignal 跳过发信号
|
||||
// → 现行提交/重试基线断言不受影响(发信号是 additive 触发面,不改状态机/落库结果)。字段名与被测一致,Mockito 按名注入。
|
||||
@Mock
|
||||
private ObjectProvider<GenTaskProducer> genTaskProducerProvider;
|
||||
|
||||
// ============================== submitGenerate ==============================
|
||||
|
||||
|
||||
@ -76,7 +76,7 @@ xxl:
|
||||
|
||||
# rocketmq 配置项,对应 RocketMQProperties 配置类
|
||||
rocketmq:
|
||||
name-server: 127.0.0.1:9876 # RocketMQ Namesrv
|
||||
name-server: 100.64.0.8:9876 # RocketMQ Namesrv(mini-infra 自托管,切片一 阶段〇 C3;替旧 127.0.0.1)
|
||||
|
||||
spring:
|
||||
# RabbitMQ 配置项,对应 RabbitProperties 配置类
|
||||
|
||||
@ -95,7 +95,7 @@ xxl:
|
||||
|
||||
# rocketmq 配置项,对应 RocketMQProperties 配置类
|
||||
rocketmq:
|
||||
name-server: 127.0.0.1:9876 # RocketMQ Namesrv
|
||||
name-server: 100.64.0.8:9876 # RocketMQ Namesrv(mini-infra 自托管,切片一 阶段〇 C3;替旧 127.0.0.1)
|
||||
|
||||
spring:
|
||||
# RabbitMQ 配置项,对应 RabbitProperties 配置类
|
||||
|
||||
@ -96,7 +96,7 @@ xxl:
|
||||
|
||||
--- #################### 消息队列 ####################
|
||||
rocketmq:
|
||||
name-server: 127.0.0.1:9876
|
||||
name-server: 100.64.0.8:9876 # mini-infra 自托管 RocketMQ(切片一 阶段〇 C3:异步 gen 队列;替旧 127.0.0.1)
|
||||
spring:
|
||||
rabbitmq:
|
||||
host: 127.0.0.1
|
||||
|
||||
@ -70,7 +70,7 @@ spring:
|
||||
namespace: public
|
||||
file-extension: yaml # 配置文件格式
|
||||
# Sentinel 入口 QPS 突发保护(切片一 阶段〇 C2):流控规则从 Nacos dataId sentinel-gen-flow-rules 热加载(改即生效、不重启)。
|
||||
# 职责 = 生成入队入口速率闸(FLOW_GRADE_QPS 削洪峰),与 DB 三门(业务配额)/ C3 consumeThreadMax(在跑并发≤15)三层正交。
|
||||
# 职责 = 生成入队入口速率闸(FLOW_GRADE_QPS 削洪峰),与 DB 三门(业务配额)/ C3 consumeThreadMax(dispatch 投递速率上限;真在跑并发≤15 属有界 worker 池=follow-up)三层正交。
|
||||
sentinel:
|
||||
transport:
|
||||
dashboard: ${SENTINEL_DASHBOARD:} # 留空 = 不接 Sentinel 控制台(mini-infra RAM 紧、不部面板);规则直存 Nacos
|
||||
@ -183,6 +183,9 @@ aj:
|
||||
|
||||
# rocketmq 配置项,对应 RocketMQProperties 配置类
|
||||
rocketmq:
|
||||
# 基址默认指向 mini-infra 自托管 RocketMQ(切片一 阶段〇 C3:异步 gen 队列 NameServer;B2 已部署 topic AIGC_GEN_TOPIC)。
|
||||
# 各 profile(dev/local/staging)同值显式声明,内网阶段一律直连 mini-infra(无本地 RocketMQ,见 docs/内网凭据与端点.md)。
|
||||
name-server: 100.64.0.8:9876 # RocketMQ Namesrv(mini-infra,替旧 127.0.0.1)
|
||||
# Producer 配置项
|
||||
producer:
|
||||
group: ${spring.application.name}_PRODUCER # 生产者分组
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user