Prechádzať zdrojové kódy

feature:记录消费的offset

xiari 1 rok pred
rodič
commit
e23142e890

+ 2 - 2
ruoyi-server/ruoyi-server-mqdata/src/main/java/org/dromara/server/mq/consumer/KafkaCloudConsumer.java

@@ -63,7 +63,7 @@ public class KafkaCloudConsumer {
         //判断 offset
         Boolean canConsume = xfOffsetService.judgeCanConsume(topic, groupId, offset, partition);
         if(!canConsume){
-            log.info("[kafka消息处理]-[消息:{}-[已消费,不能重复消费,offset: {},topic:{},groupId:{},]", record.value(), offset,topic,groupId);
+            log.info("[kafka消息处理]-[消息:{}-[已消费,不能重复消费,offset: {},topic:{},groupId:{},partition:{}]", record.value(), offset,topic,groupId,partition);
             return;
         }
 
@@ -98,7 +98,7 @@ public class KafkaCloudConsumer {
         //判断 offset
         Boolean canConsume = xfOffsetService.judgeCanConsume(topic, groupId, offset,partition);
         if(!canConsume){
-            log.info("[kafka消息处理]-[消息:{}-[已消费,不能重复消费,offset: {},topic:{},groupId:{},]", record.value(), offset,topic,groupId);
+            log.info("[kafka消息处理]-[消息:{}-[已消费,不能重复消费,offset: {},topic:{},groupId:{},partition:{}]", record.value(), offset,topic,groupId,partition);
             return;
         }
 

+ 1 - 1
ruoyi-server/ruoyi-server-mqdata/src/main/java/org/dromara/server/mq/consumer/KafkaLocalConsumer.java

@@ -56,7 +56,7 @@ public class KafkaLocalConsumer {
         //判断 offset
         Boolean canConsume = xfOffsetService.judgeCanConsume(topic, groupId, offset, partition);
         if(!canConsume){
-            log.info("[kafka消息处理]-[消息:{}-[已消费,不能重复消费,offset: {},topic:{},groupId:{},]", record.value(), offset,topic,groupId);
+            log.info("[kafka消息处理]-[消息:{}-[已消费,不能重复消费,offset: {},topic:{},groupId:{},partition:{}]", record.value(), offset,topic,groupId,partition);
             return;
         }