원클릭으로
mq-middleware
指导如何使用 toLink-Service 的 MQ 消息中台进行消息收发、定义新消息类型以及处理多厂商适配逻辑。
Codex 또는 Claude로 설치 이 Prompt를 복사해 Codex, Claude 또는 다른 어시스턴트에 붙여 넣으면 Skill 페이지를 검토하고 설치를 진행할 수 있습니다.
메뉴
指导如何使用 toLink-Service 的 MQ 消息中台进行消息收发、定义新消息类型以及处理多厂商适配逻辑。
Codex 또는 Claude로 설치 이 Prompt를 복사해 Codex, Claude 또는 다른 어시스턴트에 붙여 넣으면 Skill 페이지를 검토하고 설치를 진행할 수 있습니다.
SOC 직업 분류 기준
当用户要求基于某个需求、功能、改造、技术方案或项目实践生成博客文章时使用;必须结合用户给出的完整需求、当前项目真实实现逻辑、业务场景和代码/文档上下文做深入分析,输出通俗易懂且专业的 Markdown 博客到 `.specs/blog/《博客名称》.md`。
当修改 AGENTS.md/CLAUDE.md、docs/api、docs/internals、docs/ops,或代码变更影响这些文档记录的 API、MySQL schema、MQ 契约、Redis 缓存、OSS、错误码、模块架构、配置时,检查并同步更新对应文档,保证项目文档自动维护。
MySQL 建表与字段规范(面向 Java 管理端业务:用户、LLM 配置、数据集、知识文件、解析任务)。统一命名、索引、字段类型、时间戳、引擎字符集与注释要求,便于研发与 DBA 评审落地。
SpringDoc OpenAPI 3 中文注解生成工作流。为 Spring Boot Controller 和 DTO 生成符合企业级规范的中文 Swagger 注解(@Tag、@Operation、@Parameter、@Schema)。
brief.md 和 acceptance.feature 已冻结后,生成 .specs/<需求名>/technical_design.md;必须基于真实 Java 代码、组件文档和契约。
为 toLink-Service 的 HTTP 接口构建并执行全面的 curl 黑盒测试。分析待测接口与边界条件,必要时直连数据库或经接口造数,对本地已启动服务发起 curl 请求,断言响应,最终在对话中返回测试结果汇总。
| name | mq-middleware |
| description | 指导如何使用 toLink-Service 的 MQ 消息中台进行消息收发、定义新消息类型以及处理多厂商适配逻辑。 |
| when_to_use | 当用户要求接入 Kafka/RabbitMQ、发送或订阅消息、新增 MQ 消息类型、实现消息消费者或修改现有消息模型时激活。触发示例:'发送一条 MQ 消息'、'写个消费者'、'新增消息类型'、'改 MQ 消息字段' |
MQ 模块位于 link-components/toLink-components-mq,通过接口 + AutoConfiguration 实现多厂商(Kafka/RabbitMQ)切换。
强制同步要求:凡涉及 MQ 模块操作(新增/修改/删除消息模型、消费者、Topic、厂商适配逻辑、配置项),完成代码修改后必须:
docs/api/mq_contracts.md 中的消息清单与字段说明。| 接口 | 位置 | 职责 |
|---|---|---|
AbstractMQ | link-components/.../mq/AbstractMQ.java | 业务消息契约:getMQName()、getMQType()、getMessage() |
MQSend | link-components/.../mq/MQSend.java | 统一发送入口:send(AbstractMQ), send(AbstractMQ, int delay) |
MQMsgReceiver | link-components/.../mq/MQMsgReceiver.java | 框架侧原始消息接收契约:receive(String msg) |
MQSendType | link-components/.../mq/constant/MQSendType.java | 投递语义枚举:QUEUE(点对点)、BROADCAST(广播) |
| 场景 | 位置 |
|---|---|
| Java↔Python 双端共享(组件级) | link-components/toLink-components-mq/.../mq/model/ |
| Java 内部或服务级消息 | link-service/src/main/java/com/qingluo/link/service/mq/ |
| 消息模型 | Topic/Queue | 位置 | 方向 | 说明 |
|---|---|---|---|---|
DocumentParseTaskMQ | tolink.rag.parse_task | link-service/.../mq/ | Java→Python | 文档解析任务投递 |
DocumentParseResultMQ | tolink.rag.parse_result | link-components/.../model/ | Python→Java | 解析终态结果回传 |
CacheCompensationMQ | tolink.cache.evict | link-service/.../mq/ | 补偿生产者→Java | 缓存补偿删除 |
用 Spring 注入 MQSend,不直接实例化 Kafka / RabbitMQ vendor:
@Service
@RequiredArgsConstructor
public class YourService {
private final MQSend mqSend;
public void doSomething() {
DocumentParseTaskMQ message = new DocumentParseTaskMQ(payload);
mqSend.send(message);
}
// 延迟消息(仅 RabbitMQ 有效)
public void doDelayed() {
mqSend.send(message, 30); // 30 秒后投递
}
}
link-components/.../mq/model/link-service/.../mq/import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.annotation.JSONField;
import com.qingluo.link.components.mq.AbstractMQ;
import com.qingluo.link.components.mq.constant.MQSendType;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.springframework.util.StringUtils;
/**
* 简要描述消息用途及方向,例如:Java 向 Python 投递的 XXX 任务消息。
*/
public class YourEventMQ implements AbstractMQ {
public static final String MQ_NAME = "tolink.xxx.your_event";
private MsgPayload msgPayload;
public YourEventMQ() { this.msgPayload = new MsgPayload(); }
public YourEventMQ(MsgPayload msgPayload) { this.msgPayload = msgPayload; }
/** 仅接收方需要 parseMsg;纯发送方可省略。 */
public static MsgPayload parseMsg(String msg) {
MsgPayload payload = JSON.parseObject(msg).toJavaObject(MsgPayload.class);
validate(payload);
return payload;
}
@Override public String getMQName() { return MQ_NAME; }
@Override public MQSendType getMQType() { return MQSendType.QUEUE; } // 广播用 BROADCAST
@Override public String getMessage() { validate(msgPayload); return JSON.toJSONString(msgPayload); }
public interface MQReceiver {
void receive(MsgPayload payload);
}
@Data
@NoArgsConstructor
@AllArgsConstructor
public static class MsgPayload {
@JSONField(name = "biz_id")
private String bizId;
// 字段名使用 snake_case(与 Python 端对齐),Java 字段使用 camelCase
}
private static void validate(MsgPayload payload) {
if (payload == null) {
throw new IllegalArgumentException("your_event payload is missing");
}
if (!StringUtils.hasText(payload.getBizId())) {
throw new IllegalArgumentException("your_event biz_id is missing");
}
}
}
MQ_NAME(Topic/Queue):tolink.<domain>.<event_name>,使用稳定小写下划线,例如 tolink.rag.parse_task。@JSONField(name = "snake_case"),Java 变量用 camelCase。消费者实现 XxxMQ.MQReceiver 接口,由 KafkaMQTopologyScanner / RabbitMQTopologyScanner 自动发现并绑定:
@Component
@RequiredArgsConstructor
public class YourEventConsumer implements YourEventMQ.MQReceiver {
private final YourService yourService;
@Override
public void receive(YourEventMQ.MsgPayload payload) {
yourService.handle(payload);
}
}
若需要直接操作原始字符串(如特殊错误处理),实现 MQMsgReceiver:
@Component
public class RawConsumer implements MQMsgReceiver {
@Override
public void receive(String msg) { ... }
}
配置前缀:tolink.mq(MQProperties)
| 配置项 | 说明 | 默认值 |
|---|---|---|
tolink.mq.vender / tolink.mq.vendor | 厂商选择:rabbitMQ 或 kafka | — |
tolink.mq.scanBasePackages | 扫描消息模型的包路径 | com.qingluo |
tolink.mq.kafkaAutoCreateTopics | 启动时是否自动创建 Kafka Topics | true |
tolink.mq.kafkaTopicPartitions | Kafka Topic 默认分区数 | 1 |
tolink.mq.kafkaTopicReplicas | Kafka Topic 默认副本数 | 1 |
tolink.mq.delayedExchangeName | RabbitMQ 延迟交换机名称 | delayExchange |
tolink.mq.fanoutExchangeNamePrefix | RabbitMQ 广播交换机前缀 | fanout_exchange_ |
新增或修改消息模型后必须:
docs/api/mq_contracts.md 的消息清单与字段说明