mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
#ROCKETMQ-314# msg send back must sync change process queue msg size .
This commit is contained in:
@@ -17,6 +17,7 @@
|
|||||||
package org.apache.rocketmq.client.impl.consumer;
|
package org.apache.rocketmq.client.impl.consumer;
|
||||||
|
|
||||||
import java.util.ArrayList;
|
import java.util.ArrayList;
|
||||||
|
import java.util.Collections;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.TreeMap;
|
import java.util.TreeMap;
|
||||||
@@ -101,7 +102,7 @@ public class ProcessQueue {
|
|||||||
try {
|
try {
|
||||||
if (!msgTreeMap.isEmpty() && msg.getQueueOffset() == msgTreeMap.firstKey()) {
|
if (!msgTreeMap.isEmpty() && msg.getQueueOffset() == msgTreeMap.firstKey()) {
|
||||||
try {
|
try {
|
||||||
msgTreeMap.remove(msgTreeMap.firstKey());
|
removeMessage(Collections.singletonList(msg));
|
||||||
} catch (Exception e) {
|
} catch (Exception e) {
|
||||||
log.error("send expired msg exception", e);
|
log.error("send expired msg exception", e);
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user