Skip to content

Commit

Permalink
[fix][client] Fix negative acknowledgement by messageId (apache#23060)
Browse files Browse the repository at this point in the history
(cherry picked from commit d4bbf10)
(cherry picked from commit 02f3ecc)
  • Loading branch information
izumo27 authored and nikhil-ctds committed Jul 30, 2024
1 parent 734e3da commit dd022fa
Show file tree
Hide file tree
Showing 2 changed files with 9 additions and 6 deletions.
Original file line number Diff line number Diff line change
Expand Up @@ -134,7 +134,7 @@ public void testNegativeAcks(boolean batching, boolean usePartitions, Subscripti
Set<String> sentMessages = new HashSet<>();

final int N = 10;
for (int i = 0; i < N; i++) {
for (int i = 0; i < N * 2; i++) {
String value = "test-" + i;
producer.sendAsync(value);
sentMessages.add(value);
Expand All @@ -145,14 +145,19 @@ public void testNegativeAcks(boolean batching, boolean usePartitions, Subscripti
Message<String> msg = consumer.receive();
consumer.negativeAcknowledge(msg);
}
for (int i = 0; i < N; i++) {
Message<String> msg = consumer.receive();
consumer.negativeAcknowledge(msg.getMessageId());
}


assertTrue(consumer instanceof ConsumerBase<String>);
assertEquals(((ConsumerBase<String>) consumer).getUnAckedMessageTracker().size(), 0);

Set<String> receivedMessages = new HashSet<>();

// All the messages should be received again
for (int i = 0; i < N; i++) {
for (int i = 0; i < N * 2; i++) {
Message<String> msg = consumer.receive();
receivedMessages.add(msg.getValue());
consumer.acknowledge(msg);
Expand Down Expand Up @@ -310,9 +315,7 @@ public void testNegativeAcksDeleteFromUnackedTracker() throws Exception {
assertEquals(unAckedMessageTracker.size(), 0);
negativeAcksTracker.close();
// negative batch message id
unAckedMessageTracker.add(batchMessageId);
unAckedMessageTracker.add(batchMessageId2);
unAckedMessageTracker.add(batchMessageId3);
unAckedMessageTracker.add(messageId);
consumer.negativeAcknowledge(batchMessageId);
consumer.negativeAcknowledge(batchMessageId2);
consumer.negativeAcknowledge(batchMessageId3);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -762,7 +762,7 @@ public void negativeAcknowledge(MessageId messageId) {
negativeAcksTracker.add(messageId);

// Ensure the message is not redelivered for ack-timeout, since we did receive an "ack"
unAckedMessageTracker.remove(messageId);
unAckedMessageTracker.remove(MessageIdAdvUtils.discardBatch(messageId));
}

@Override
Expand Down

0 comments on commit dd022fa

Please sign in to comment.