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,136 @@ public void decCorePoolSize() {

}

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

/**
* Map the filtered-list ackIndex to the original msgs position.
* This is a helper method that can be tested independently.
*
* @param originalMsgs The original message list (contains duplicates)
* @param filteredMsgs The filtered message list (duplicates removed)
* @param filteredAckIndex The ackIndex from the listener (based on filteredMsgs)
* @return The mapped ackIndex for the original list, or -1 if mapping fails
*/
static int mapAckIndex(List<MessageExt> originalMsgs, List<MessageExt> filteredMsgs, int filteredAckIndex) {
if (originalMsgs == null || filteredMsgs == null || filteredAckIndex < 0) {
return -1;
}

// No duplicates, no mapping needed
if (filteredMsgs.size() == originalMsgs.size()) {
return Math.min(filteredAckIndex, originalMsgs.size() - 1);
}

// Clamp filteredAckIndex to valid range
if (filteredAckIndex >= filteredMsgs.size()) {
filteredAckIndex = filteredMsgs.size() - 1;
}

// When all filtered messages succeed, set ackIndex to cover all original messages
// This ensures trailing duplicates are not treated as failed
if (filteredAckIndex == filteredMsgs.size() - 1) {
return originalMsgs.size() - 1;
}

// Partial success: find the position in original msgs that corresponds to filteredMsgs[filteredAckIndex]
MessageExt lastAckedMsg = filteredMsgs.get(filteredAckIndex);
for (int i = 0; i < originalMsgs.size(); i++) {
if (originalMsgs.get(i) == lastAckedMsg) {
return i;
}
}

return -1; // Should not happen if lists are consistent
}

/**
* 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 @@ -387,14 +514,19 @@ 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);
final boolean hasDuplicates = filteredMsgs.size() < msgs.size();

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 +535,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 +578,66 @@ public void run() {
if (null == status) {
log.warn("consumeMessage return null, Group: {} Msgs: {} MQ: {}",
ConsumeMessageConcurrentlyService.this.consumerGroup,
msgs,
filteredMsgs,
messageQueue);
status = ConsumeConcurrentlyStatus.RECONSUME_LATER;
}

// Handle ackIndex semantics when duplicates were filtered
// ackIndex is based on filteredMsgs, but processConsumeResult uses original msgs
// We need to:
// 1. Only mark messages up to ackIndex in filteredMsgs as processed
// 2. Map the filtered-list ackIndex to the original msgs position
if (status == ConsumeConcurrentlyStatus.CONSUME_SUCCESS) {
int filteredAckIndex = context.getAckIndex();

if (!filteredMsgs.isEmpty() && filteredAckIndex >= 0) {
// Clamp ackIndex to valid range for filteredMsgs
if (filteredAckIndex >= filteredMsgs.size()) {
filteredAckIndex = filteredMsgs.size() - 1;
}

// Only mark successfully consumed messages (up to filteredAckIndex)
List<MessageExt> successfullyConsumed = new ArrayList<>(filteredAckIndex + 1);
for (int i = 0; i <= filteredAckIndex; i++) {
successfullyConsumed.add(filteredMsgs.get(i));
}
markMessagesAsProcessed(successfullyConsumed);
}

// Map filtered-list ackIndex to original msgs position
if (hasDuplicates && !filteredMsgs.isEmpty() && filteredAckIndex >= 0) {
// When all filtered messages succeed, set ackIndex to cover all original messages
// This ensures trailing duplicates are not treated as failed (which would cause unnecessary retry)
if (filteredAckIndex == filteredMsgs.size() - 1) {
// Full success: all filtered messages consumed successfully
// Set ackIndex to cover the entire original batch (including duplicates that were skipped)
context.setAckIndex(msgs.size() - 1);
log.debug("Duplicate messages filtered, full success. Set ackIndex to {} (covers all original messages)",
msgs.size() - 1);
} else {
// Partial success: find the position in original msgs that corresponds to filteredMsgs[filteredAckIndex]
MessageExt lastAckedMsg = filteredMsgs.get(filteredAckIndex);
int originalAckIndex = -1;
for (int i = 0; i < msgs.size(); i++) {
if (msgs.get(i) == lastAckedMsg) {
originalAckIndex = i;
break;
}
}

if (originalAckIndex >= 0) {
context.setAckIndex(originalAckIndex);
log.debug("Duplicate messages filtered, mapped filteredAckIndex {} to originalAckIndex {}",
filteredAckIndex, originalAckIndex);
}
}
}
// If no duplicates, ackIndex already refers to original msgs, no mapping needed
}
// For RECONSUME_LATER, don't mark any messages as processed
// processConsumeResult will set ackIndex = -1, sending all messages back

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