with one click
mq-middleware
指导如何使用 toLink-Service 的 MQ 消息中台进行消息收发、定义新消息类型以及处理多厂商适配逻辑。
Install with Codex or Claude Copy this prompt, paste it into Codex, Claude, or another assistant, and let it review the skill page and install it for you.
Menu
指导如何使用 toLink-Service 的 MQ 消息中台进行消息收发、定义新消息类型以及处理多厂商适配逻辑。
Install with Codex or Claude Copy this prompt, paste it into Codex, Claude, or another assistant, and let it review the skill page and install it for you.
Based on SOC occupation classification
| 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 的消息清单与字段说明当用户要求基于某个需求、功能、改造、技术方案或项目实践生成博客文章时使用;必须结合用户给出的完整需求、当前项目真实实现逻辑、业务场景和代码/文档上下文做深入分析,输出通俗易懂且专业的 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 请求,断言响应,最终在对话中返回测试结果汇总。