|
@@ -144,24 +144,24 @@ public class RocketSessionConsumer implements
|
|
|
// return ConsumeOrderlyStatus.SUCCESS;//成功
|
|
|
// }
|
|
|
|
|
|
- @Service
|
|
|
- @RocketMQMessageListener(consumerGroup = "${mq.config.sessionConsumerWebGroup}", topic = "${mq.config.sessionTopic}", selectorType = SelectorType.TAG, selectorExpression = "${mq.config.sessionTopicWebTag}")
|
|
|
- public class sessionConsumerWeb implements RocketMQListener<Message>, RocketMQPushConsumerLifecycleListener {
|
|
|
-
|
|
|
- @Override
|
|
|
- public void onMessage(Message message) {
|
|
|
- //实现RocketMQPushConsumerLifecycleListener监听器之后,此方法不调用
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public void prepareStart(DefaultMQPushConsumer defaultMQPushConsumer) {
|
|
|
- defaultMQPushConsumer.setConsumeMessageBatchMaxSize(SystemConstant.CONSUME_MESSAGE_BATCH_MAX_SIZE);//每次拉取10条
|
|
|
- defaultMQPushConsumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
|
|
|
- defaultMQPushConsumer.setMaxReconsumeTimes(SystemConstant.MAXRECONSUMETIMES);//最大重试次数
|
|
|
-// defaultMQPushConsumer.setMessageModel(MessageModel.BROADCASTING);
|
|
|
- defaultMQPushConsumer.registerMessageListener(RocketSessionConsumer.this::consumeMessage);
|
|
|
- }
|
|
|
- }
|
|
|
+// @Service
|
|
|
+// @RocketMQMessageListener(consumerGroup = "${mq.config.sessionConsumerWebGroup}", topic = "${mq.config.sessionTopic}", selectorType = SelectorType.TAG, selectorExpression = "${mq.config.sessionTopicWebTag}")
|
|
|
+// public class sessionConsumerWeb implements RocketMQListener<Message>, RocketMQPushConsumerLifecycleListener {
|
|
|
+//
|
|
|
+// @Override
|
|
|
+// public void onMessage(Message message) {
|
|
|
+// //实现RocketMQPushConsumerLifecycleListener监听器之后,此方法不调用
|
|
|
+// }
|
|
|
+//
|
|
|
+// @Override
|
|
|
+// public void prepareStart(DefaultMQPushConsumer defaultMQPushConsumer) {
|
|
|
+// defaultMQPushConsumer.setConsumeMessageBatchMaxSize(SystemConstant.CONSUME_MESSAGE_BATCH_MAX_SIZE);//每次拉取10条
|
|
|
+// defaultMQPushConsumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
|
|
|
+// defaultMQPushConsumer.setMaxReconsumeTimes(SystemConstant.MAXRECONSUMETIMES);//最大重试次数
|
|
|
+//// defaultMQPushConsumer.setMessageModel(MessageModel.BROADCASTING);
|
|
|
+// defaultMQPushConsumer.registerMessageListener(RocketSessionConsumer.this::consumeMessage);
|
|
|
+// }
|
|
|
+// }
|
|
|
|
|
|
@Service
|
|
|
@RocketMQMessageListener(consumerGroup = "${mq.config.sessionConsumerPcGroup}", topic = "${mq.config.sessionTopic}", selectorType = SelectorType.TAG, selectorExpression = "${mq.config.sessionTopicPcTag}")
|