messageManager

This commit is contained in:
redkale
2023-12-24 08:40:54 +08:00
parent b52090d1f2
commit 240b2b226b
3 changed files with 33 additions and 4 deletions

View File

@@ -0,0 +1,29 @@
/*
*
*/
package org.redkale.mq;
import java.util.List;
import org.redkale.inject.Resourcable;
/**
* MQ消息管理器
*
* <p>
* 详情见: https://redkale.org
*
* @author zhangjx
*
* @since 2.8.0
*/
public interface MessageManager extends Resourcable {
//
public boolean createTopic(String... topics);
//删除topic如果不存在则跳过
public boolean deleteTopic(String... topics);
//查询所有topic
public abstract List<String> queryTopic();
}

View File

@@ -24,10 +24,10 @@ import org.redkale.convert.Convert;
import org.redkale.convert.ConvertFactory;
import org.redkale.convert.ConvertType;
import org.redkale.convert.json.JsonConvert;
import org.redkale.inject.Resourcable;
import org.redkale.inject.ResourceEvent;
import org.redkale.mq.MessageConext;
import org.redkale.mq.MessageConsumer;
import org.redkale.mq.MessageManager;
import org.redkale.mq.MessageProducer;
import org.redkale.mq.ResourceConsumer;
import org.redkale.mq.ResourceProducer;
@@ -47,7 +47,7 @@ import org.redkale.util.*;
*
* @since 2.1.0
*/
public abstract class MessageAgent implements Resourcable {
public abstract class MessageAgent implements MessageManager {
protected final Logger logger = Logger.getLogger(this.getClass().getSimpleName());

View File

@@ -3,8 +3,6 @@
*/
package org.redkale.mq.spi;
import org.redkale.mq.spi.MessageAgent;
import org.redkale.mq.spi.MessageAgentProvider;
import java.lang.reflect.Field;
import java.util.ArrayList;
import java.util.HashSet;
@@ -26,6 +24,7 @@ import org.redkale.inject.ResourceEvent;
import org.redkale.inject.ResourceFactory;
import org.redkale.inject.ResourceTypeLoader;
import org.redkale.mq.MessageConsumer;
import org.redkale.mq.MessageManager;
import org.redkale.mq.MessageProducer;
import org.redkale.mq.ResourceConsumer;
import org.redkale.mq.ResourceProducer;
@@ -202,6 +201,7 @@ public class MessageModuleEngine extends ModuleEngine {
for (MessageAgent agent : this.messageAgents) {
this.resourceFactory.inject(agent);
agent.init(agent.getConfig());
this.resourceFactory.register(agent.getName(), MessageManager.class, agent);
this.resourceFactory.register(agent.getName(), MessageAgent.class, agent);
}
logger.info("MessageAgent init in " + (System.currentTimeMillis() - s) + " ms");