Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -273,6 +273,31 @@ public class DefaultMQPushConsumer extends ClientConfig implements MQPushConsume
*/
private int popBatchNums = 32;

/**
* Enable message deduplication based on message keys.
* When enabled, duplicate messages will be detected and skipped before consumption.
* Default: false (disabled)
*
* <p><strong>Note:</strong> Deduplication is currently only supported for concurrent consumption mode.
* Orderly consumption and POP mode do not support deduplication and will log a warning if enabled.</p>
*/
private boolean enableMessageDeduplication = false;

/**
* Maximum size of deduplication cache (number of message keys).
* Range: [1000, 100000]
* Default: 10000
*/
private int deduplicationCacheSize = 10000;

/**
* Deduplication cache entry expire time in milliseconds.
* Older entries will be removed during cleanup.
* Range: [10000, 3600000] (10s - 1h)
* Default: 60000 (1 minute)
*/
private long deduplicationCacheExpireTime = 60000;

/**
* Maximum time to await message consuming when shutdown consumer, 0 indicates no await.
*/
Expand Down Expand Up @@ -987,6 +1012,30 @@ public void setPopBatchNums(int popBatchNums) {
this.popBatchNums = popBatchNums;
}

public boolean isEnableMessageDeduplication() {
return enableMessageDeduplication;
}

public void setEnableMessageDeduplication(boolean enableMessageDeduplication) {
this.enableMessageDeduplication = enableMessageDeduplication;
}

public int getDeduplicationCacheSize() {
return deduplicationCacheSize;
}

public void setDeduplicationCacheSize(int deduplicationCacheSize) {
this.deduplicationCacheSize = deduplicationCacheSize;
}

public long getDeduplicationCacheExpireTime() {
return deduplicationCacheExpireTime;
}

public void setDeduplicationCacheExpireTime(long deduplicationCacheExpireTime) {
this.deduplicationCacheExpireTime = deduplicationCacheExpireTime;
}

public boolean isClientRebalance() {
return clientRebalance;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,11 @@
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.Executors;
import java.util.concurrent.LinkedBlockingQueue;
Expand Down Expand Up @@ -122,11 +124,95 @@ public void decCorePoolSize() {

}

@Override
public int getCorePoolSize() {
return this.consumeExecutor.getCorePoolSize();
}

/**
* Filter duplicate messages from the list.
* Creates a new list with non-duplicate messages for consumption.
*
* <p>A message is treated as a duplicate when its deduplication key is either:
* <ul>
* <li>already present in the global processed cache (cross-batch dedup), or</li>
* <li>already seen earlier in this same batch (intra-batch dedup).</li>
* </ul>
* The global cache is only populated after successful consumption, so without the batch-local
* check two identical messages in one batch would both reach the listener.
*
* @param msgs Original message list (will not be modified)
* @return New list containing only non-duplicate messages
*/
private List<MessageExt> filterDuplicateMessages(List<MessageExt> msgs) {
return filterDuplicateMessages(msgs,
this.defaultMQPushConsumerImpl.getMessageDeduplicator(), this.consumerGroup);
}

/**
* Filter duplicate messages from the list, deduplicating both against the global cache and
* within the batch itself. Exposed as package-private for unit testing.
*
* @param msgs Original message list (will not be modified)
* @param deduplicator The deduplicator holding already-processed keys, or null to disable
* @param consumerGroup Consumer group, for logging only
* @return New list containing only non-duplicate messages, or the original list when
* deduplication is disabled
*/
static List<MessageExt> filterDuplicateMessages(List<MessageExt> msgs,
MessageDeduplicator deduplicator, String consumerGroup) {
if (deduplicator == null || msgs == null || msgs.isEmpty()) {
return msgs;
}

List<MessageExt> filteredMsgs = new ArrayList<>(msgs.size());
Set<String> seenInBatch = new HashSet<>(msgs.size() * 2);
int duplicateCount = 0;

for (MessageExt msg : msgs) {
String deduplicationKey = MessageDeduplicator.getDeduplicationKey(msg);

// A null/empty key means we cannot deduplicate this message; keep it.
boolean duplicate = deduplicationKey != null && !deduplicationKey.isEmpty()
&& (deduplicator.isDuplicate(deduplicationKey) || !seenInBatch.add(deduplicationKey));

if (duplicate) {
// Duplicate message detected
duplicateCount++;
log.warn("Duplicate message detected. msgId={}, keys={}, topic={}, queueId={}, queueOffset={}",
msg.getMsgId(), msg.getKeys(), msg.getTopic(), msg.getQueueId(), msg.getQueueOffset());
} else {
// Non-duplicate message, add to filtered list
filteredMsgs.add(msg);
}
}

if (duplicateCount > 0) {
log.info("Found {} duplicate messages in batch of {} messages for consumer group {}. Filtered list size: {}",
duplicateCount, msgs.size(), consumerGroup, filteredMsgs.size());
}

return filteredMsgs;
}

/**
* Mark successfully consumed messages as processed in deduplication cache.
*
* @param msgs Messages to mark as processed
*/
private void markMessagesAsProcessed(List<MessageExt> msgs) {
MessageDeduplicator deduplicator = this.defaultMQPushConsumerImpl.getMessageDeduplicator();
if (deduplicator == null || msgs == null || msgs.isEmpty()) {
return;
}

for (MessageExt msg : msgs) {
String deduplicationKey = MessageDeduplicator.getDeduplicationKey(msg);
if (deduplicationKey != null) {
deduplicator.markProcessed(deduplicationKey);
}
}
}

@Override
public ConsumeMessageDirectlyResult consumeMessageDirectly(MessageExt msg, String brokerName) {
ConsumeMessageDirectlyResult result = new ConsumeMessageDirectlyResult();
Expand Down Expand Up @@ -246,39 +332,46 @@ public void processConsumeResult(
) {
int ackIndex = context.getAckIndex();

if (consumeRequest.getMsgs().isEmpty())
// The messages the listener actually consumed (duplicates filtered out), or the
// original batch when deduplication is disabled. ackIndex is relative to this list,
// so the send-back loop and consume stats must operate on it too. Skipped duplicates
// are neither consumed nor failed and must never be sent back.
final List<MessageExt> processedMsgs = consumeRequest.getProcessedMsgs();
final List<MessageExt> originalMsgs = consumeRequest.getMsgs();

if (originalMsgs.isEmpty())
return;

switch (status) {
case CONSUME_SUCCESS:
if (ackIndex >= consumeRequest.getMsgs().size()) {
ackIndex = consumeRequest.getMsgs().size() - 1;
if (ackIndex >= processedMsgs.size()) {
ackIndex = processedMsgs.size() - 1;
}
int ok = ackIndex + 1;
int failed = consumeRequest.getMsgs().size() - ok;
int failed = processedMsgs.size() - ok;
this.getConsumerStatsManager().incConsumeOKTPS(consumerGroup, consumeRequest.getMessageQueue().getTopic(), ok);
this.getConsumerStatsManager().incConsumeFailedTPS(consumerGroup, consumeRequest.getMessageQueue().getTopic(), failed);
break;
case RECONSUME_LATER:
ackIndex = -1;
this.getConsumerStatsManager().incConsumeFailedTPS(consumerGroup, consumeRequest.getMessageQueue().getTopic(),
consumeRequest.getMsgs().size());
processedMsgs.size());
break;
default:
break;
}

switch (this.defaultMQPushConsumer.getMessageModel()) {
case BROADCASTING:
for (int i = ackIndex + 1; i < consumeRequest.getMsgs().size(); i++) {
MessageExt msg = consumeRequest.getMsgs().get(i);
for (int i = ackIndex + 1; i < processedMsgs.size(); i++) {
MessageExt msg = processedMsgs.get(i);
log.warn("BROADCASTING, the message consume failed, drop it, {}", msg.toString());
}
break;
case CLUSTERING:
List<MessageExt> msgBackFailed = new ArrayList<>(consumeRequest.getMsgs().size());
for (int i = ackIndex + 1; i < consumeRequest.getMsgs().size(); i++) {
MessageExt msg = consumeRequest.getMsgs().get(i);
List<MessageExt> msgBackFailed = new ArrayList<>(processedMsgs.size());
for (int i = ackIndex + 1; i < processedMsgs.size(); i++) {
MessageExt msg = processedMsgs.get(i);
// Maybe message is expired and cleaned, just ignore it.
if (!consumeRequest.getProcessQueue().containsMessage(msg)) {
log.info("Message is not found in its process queue; skip send-back-procedure, topic={}, "
Expand All @@ -294,7 +387,7 @@ public void processConsumeResult(
}

if (!msgBackFailed.isEmpty()) {
consumeRequest.getMsgs().removeAll(msgBackFailed);
originalMsgs.removeAll(msgBackFailed);

this.submitConsumeRequestLater(msgBackFailed, consumeRequest.getProcessQueue(), consumeRequest.getMessageQueue());
}
Expand All @@ -303,7 +396,21 @@ public void processConsumeResult(
break;
}

long offset = consumeRequest.getProcessQueue().removeMessage(consumeRequest.getMsgs());
// Only populate the dedup cache once the consume result is actually committed, i.e. the
// processQueue has not been dropped. Marking earlier (before the drop check) would let
// uncommitted messages suppress their own redelivery, breaking at-least-once.
if (status == ConsumeConcurrentlyStatus.CONSUME_SUCCESS
&& ackIndex >= 0 && !consumeRequest.getProcessQueue().isDropped()) {
List<MessageExt> successfullyConsumed = new ArrayList<>(ackIndex + 1);
for (int i = 0; i <= ackIndex; i++) {
successfullyConsumed.add(processedMsgs.get(i));
}
markMessagesAsProcessed(successfullyConsumed);
}

// removeMessage drains by queueOffset, so it must receive the full original batch
// (including skipped duplicates) to keep msgTreeMap and the committed offset consistent.
long offset = consumeRequest.getProcessQueue().removeMessage(originalMsgs);
if (offset >= 0 && !consumeRequest.getProcessQueue().isDropped()) {
this.defaultMQPushConsumerImpl.getOffsetStore().updateOffset(consumeRequest.getMessageQueue(), offset, true);
}
Expand Down Expand Up @@ -359,17 +466,31 @@ class ConsumeRequest implements Runnable {
private final List<MessageExt> msgs;
private final ProcessQueue processQueue;
private final MessageQueue messageQueue;
/**
* The messages the listener actually consumed: the deduped list when deduplication is
* enabled, otherwise the original {@link #msgs}. Set in {@link #run()} after filtering.
* Defaults to {@code msgs} so the value is always non-null for processConsumeResult.
*/
private List<MessageExt> processedMsgs;

public ConsumeRequest(List<MessageExt> msgs, ProcessQueue processQueue, MessageQueue messageQueue) {
this.msgs = msgs;
this.processQueue = processQueue;
this.messageQueue = messageQueue;
this.processedMsgs = msgs;
}

public List<MessageExt> getMsgs() {
return msgs;
}

/**
* @return the messages the listener consumed (duplicates filtered); never null.
*/
public List<MessageExt> getProcessedMsgs() {
return processedMsgs;
}

public ProcessQueue getProcessQueue() {
return processQueue;
}
Expand All @@ -387,14 +508,23 @@ public void run() {
defaultMQPushConsumerImpl.tryResetPopRetryTopic(msgs, consumerGroup);
defaultMQPushConsumerImpl.resetRetryAndNamespace(msgs, defaultMQPushConsumer.getConsumerGroup());

// Filter duplicate messages if deduplication is enabled
// Create a new list with non-duplicate messages for consumption
List<MessageExt> filteredMsgs = filterDuplicateMessages(msgs);
// processedMsgs is the list the listener actually consumes (duplicates removed).
// It defaults to the original batch when dedup is off or every message is a duplicate,
// and is later read by processConsumeResult for send-back/stats. ackIndex is relative
// to this list.
this.processedMsgs = filteredMsgs.isEmpty() ? msgs : filteredMsgs;

ConsumeMessageContext consumeMessageContext = null;
if (ConsumeMessageConcurrentlyService.this.defaultMQPushConsumerImpl.hasHook()) {
consumeMessageContext = new ConsumeMessageContext();
consumeMessageContext.setNamespace(defaultMQPushConsumer.getNamespace());
consumeMessageContext.setConsumerGroup(defaultMQPushConsumer.getConsumerGroup());
consumeMessageContext.setProps(new HashMap<>());
consumeMessageContext.setMq(messageQueue);
consumeMessageContext.setMsgList(msgs);
consumeMessageContext.setMsgList(filteredMsgs.isEmpty() ? msgs : filteredMsgs);
consumeMessageContext.setSuccess(false);
ConsumeMessageConcurrentlyService.this.defaultMQPushConsumerImpl.executeHookBefore(consumeMessageContext);
}
Expand All @@ -403,17 +533,24 @@ public void run() {
boolean hasException = false;
ConsumeReturnType returnType = ConsumeReturnType.SUCCESS;
try {
if (msgs != null && !msgs.isEmpty()) {
for (MessageExt msg : msgs) {
// Prepare messages for consumption
List<MessageExt> msgsToConsume = filteredMsgs.isEmpty() ? Collections.emptyList() : filteredMsgs;

if (!msgsToConsume.isEmpty()) {
for (MessageExt msg : msgsToConsume) {
MessageAccessor.setConsumeStartTimeStamp(msg, String.valueOf(System.currentTimeMillis()));
}
status = listener.consumeMessage(Collections.unmodifiableList(msgsToConsume), context);
} else {
// All messages were duplicates, mark as success without calling listener
status = ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
log.info("All messages in batch were duplicates for queue {}. Skipping listener invocation.", messageQueue);
}
status = listener.consumeMessage(Collections.unmodifiableList(msgs), context);
} catch (Throwable e) {
log.warn("consumeMessage exception: {} Group: {} Msgs: {} MQ: {}",
UtilAll.exceptionSimpleDesc(e),
ConsumeMessageConcurrentlyService.this.consumerGroup,
msgs,
filteredMsgs,
messageQueue, e);
hasException = true;
}
Expand All @@ -439,11 +576,17 @@ public void run() {
if (null == status) {
log.warn("consumeMessage return null, Group: {} Msgs: {} MQ: {}",
ConsumeMessageConcurrentlyService.this.consumerGroup,
msgs,
filteredMsgs,
messageQueue);
status = ConsumeConcurrentlyStatus.RECONSUME_LATER;
}

// No ackIndex remapping here: the listener's ackIndex is already relative to
// the filtered list (processedMsgs) it consumed. processConsumeResult uses
// processedMsgs for send-back and stats, so the index needs no translation.
// Dedup-cache marking is deferred to processConsumeResult, gated on the
// processQueue not being dropped (at-least-once safety).

if (ConsumeMessageConcurrentlyService.this.defaultMQPushConsumerImpl.hasHook()) {
consumeMessageContext.setStatus(status.toString());
consumeMessageContext.setSuccess(ConsumeConcurrentlyStatus.CONSUME_SUCCESS == status);
Expand Down
Loading
Loading