mirror of
https://github.com/apache/rocketmq.git
synced 2026-08-30 18:10:44 +08:00
Signed-off-by: Aman Gautam <amangautam2128@gmail.com>
This commit is contained in:
@@ -216,6 +216,13 @@ public class PopBufferMergeService extends ServiceThread {
|
||||
}
|
||||
}
|
||||
|
||||
private boolean isSubscriptionGroupNotExist(PopCheckPointWrapper pointWrapper) {
|
||||
String group = pointWrapper.getCk().getCId();
|
||||
return brokerController.getSubscriptionGroupManager()
|
||||
.findSubscriptionGroupConfig(group) == null;
|
||||
}
|
||||
|
||||
|
||||
private void scan() {
|
||||
long startTime = System.currentTimeMillis();
|
||||
AtomicInteger count = new AtomicInteger(0);
|
||||
@@ -225,6 +232,19 @@ public class PopBufferMergeService extends ServiceThread {
|
||||
Map.Entry<String, PopCheckPointWrapper> entry = iterator.next();
|
||||
PopCheckPointWrapper pointWrapper = entry.getValue();
|
||||
|
||||
// Skip invalid POP records when consumer group does not exist
|
||||
if (isSubscriptionGroupNotExist(pointWrapper)) {
|
||||
POP_LOGGER.warn(
|
||||
"[PopBuffer] skip pop record because consumer group not exist, group={}, ck={}",
|
||||
pointWrapper.getCk().getCId(),
|
||||
pointWrapper
|
||||
);
|
||||
iterator.remove();
|
||||
counter.decrementAndGet();
|
||||
continue;
|
||||
}
|
||||
|
||||
|
||||
// just process offset(already stored at pull thread), or buffer ck(not stored and ack finish)
|
||||
if (pointWrapper.isJustOffset() && pointWrapper.isCkStored() || isCkDone(pointWrapper)
|
||||
|| isCkDoneForFinish(pointWrapper) && pointWrapper.isCkStored()) {
|
||||
|
||||
Reference in New Issue
Block a user