修复打包失败
All checks were successful
Build and Push to Target Registry / 构建并推送镜像到目标仓库 (push) Successful in 12m8s
All checks were successful
Build and Push to Target Registry / 构建并推送镜像到目标仓库 (push) Successful in 12m8s
This commit is contained in:
@@ -1,44 +1,44 @@
|
||||
package org.dromara.sis.rocketmq.consumer;
|
||||
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.rocketmq.common.message.MessageExt;
|
||||
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
|
||||
import org.apache.rocketmq.spring.core.RocketMQListener;
|
||||
import org.dromara.sis.rocketmq.RocketMqConstants;
|
||||
import org.dromara.sis.rocketmq.producer.ProducerService;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
/**
|
||||
* @author lsm
|
||||
* @apiNote MeterRecordConsumer
|
||||
* @since 2025/8/25
|
||||
*/
|
||||
@Slf4j
|
||||
@Component
|
||||
@RequiredArgsConstructor
|
||||
@RocketMQMessageListener(
|
||||
topic = RocketMqConstants.TOPIC,
|
||||
consumerGroup = RocketMqConstants.METER_GROUP,
|
||||
selectorExpression = RocketMqConstants.METER_RECORD,
|
||||
nameServer = "${rocketmq.cluster1.name-server}"
|
||||
)
|
||||
public class MeterRecordConsumer implements RocketMQListener<MessageExt> {
|
||||
|
||||
private final ProducerService producerService;
|
||||
|
||||
@Override
|
||||
public void onMessage(MessageExt ext) {
|
||||
try {
|
||||
if (ext.getBody() == null) {
|
||||
log.info("仪表上报消息数据,不转发!");
|
||||
} else {
|
||||
producerService.defaultSend(RocketMqConstants.TOPIC, RocketMqConstants.METER_RECORD, new String(ext.getBody()));
|
||||
log.info("转发仪表上报数据处理成功");
|
||||
}
|
||||
} catch (Exception e) {
|
||||
log.error("转发仪表上报数据处理失败,", e);
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
//package org.dromara.sis.rocketmq.consumer;
|
||||
//
|
||||
//import lombok.RequiredArgsConstructor;
|
||||
//import lombok.extern.slf4j.Slf4j;
|
||||
//import org.apache.rocketmq.common.message.MessageExt;
|
||||
//import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
|
||||
//import org.apache.rocketmq.spring.core.RocketMQListener;
|
||||
//import org.dromara.sis.rocketmq.RocketMqConstants;
|
||||
//import org.dromara.sis.rocketmq.producer.ProducerService;
|
||||
//import org.springframework.stereotype.Component;
|
||||
//
|
||||
///**
|
||||
// * @author lsm
|
||||
// * @apiNote MeterRecordConsumer
|
||||
// * @since 2025/8/25
|
||||
// */
|
||||
//@Slf4j
|
||||
//@Component
|
||||
//@RequiredArgsConstructor
|
||||
//@RocketMQMessageListener(
|
||||
// topic = RocketMqConstants.TOPIC,
|
||||
// consumerGroup = RocketMqConstants.METER_GROUP,
|
||||
// selectorExpression = RocketMqConstants.METER_RECORD,
|
||||
// nameServer = "${rocketmq.cluster1.name-server}"
|
||||
//)
|
||||
//public class MeterRecordConsumer implements RocketMQListener<MessageExt> {
|
||||
//
|
||||
// private final ProducerService producerService;
|
||||
//
|
||||
// @Override
|
||||
// public void onMessage(MessageExt ext) {
|
||||
// try {
|
||||
// if (ext.getBody() == null) {
|
||||
// log.info("仪表上报消息数据,不转发!");
|
||||
// } else {
|
||||
// producerService.defaultSend(RocketMqConstants.TOPIC, RocketMqConstants.METER_RECORD, new String(ext.getBody()));
|
||||
// log.info("转发仪表上报数据处理成功");
|
||||
// }
|
||||
// } catch (Exception e) {
|
||||
// log.error("转发仪表上报数据处理失败,", e);
|
||||
// }
|
||||
//
|
||||
// }
|
||||
//}
|
||||
|
Reference in New Issue
Block a user