| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent a292bf2 commit fcd31ea
1 file changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -31,6 +31,7 @@ | |||
| 31 | 31 | import java.util.ArrayList; | |
| 32 | 32 | import java.util.HashMap; | |
| 33 | 33 | import java.util.Iterator; | |
| 34 | + import java.util.LinkedHashMap; | ||
| 34 | 35 | import java.util.List; | |
| 35 | 36 | import java.util.Map; | |
| 36 | 37 | import java.util.Map.Entry; | |
@@ -92,8 +93,8 @@ class MessageDispatcher { | |||
| 92 | 93 | private final LinkedBlockingQueue<AckRequestData> pendingAcks = new LinkedBlockingQueue<>(); | |
| 93 | 94 | private final LinkedBlockingQueue<AckRequestData> pendingNacks = new LinkedBlockingQueue<>(); | |
| 94 | 95 | private final LinkedBlockingQueue<AckRequestData> pendingReceipts = new LinkedBlockingQueue<>(); | |
| 95 | - private final ConcurrentMap<String, ReceiptCompleteData> outstandingReceipts = | ||
| 96 | - new ConcurrentHashMap<String, ReceiptCompleteData>(); | ||
| 96 | + private final LinkedHashMap<String, ReceiptCompleteData> outstandingReceipts = | ||
| 97 | + new LinkedHashMap<String, ReceiptCompleteData>(); | ||
| 97 | 98 | private final AtomicInteger messageDeadlineSeconds = new AtomicInteger(); | |
| 98 | 99 | private final AtomicBoolean extendDeadline = new AtomicBoolean(true); | |
| 99 | 100 | private final Lock jobLock; | |
@@ -397,7 +398,9 @@ void processReceivedMessages(List<ReceivedMessage> messages) { | |||
| 397 | 398 | if (this.exactlyOnceDeliveryEnabled.get()) { | |
| 398 | 399 | // For exactly once deliveries we don't add to outstanding batch because we first | |
| 399 | 400 | // process the receipt modack. If that is successful then we process the message. | |
| 400 | - outstandingReceipts.put(message.getAckId(), new ReceiptCompleteData(outstandingMessage)); | ||
| 401 | + synchronized (outstandingReceipts) { | ||
| 402 | + outstandingReceipts.put(message.getAckId(), new ReceiptCompleteData(outstandingMessage)); | ||
| 403 | + } | ||
| 401 | 404 | } else if (pendingMessages.putIfAbsent(message.getAckId(), ackHandler) != null) { | |
| 402 | 405 | // putIfAbsent puts ackHandler if ackID isn't previously mapped, then return the | |
| 403 | 406 | // previously-mapped element. | |
@@ -417,33 +420,36 @@ void processReceivedMessages(List<ReceivedMessage> messages) { | |||
| 417 | 420 | } | |
| 418 | 421 | ||
| 419 | 422 | void notifyAckSuccess(AckRequestData ackRequestData) { | |
| 420 | - | ||
| 421 | - if (outstandingReceipts.containsKey(ackRequestData.getAckId())) { | ||
| 422 | - outstandingReceipts.get(ackRequestData.getAckId()).notifyReceiptComplete(); | ||
| 423 | - List<OutstandingMessage> outstandingBatch = new ArrayList<>(); | ||
| 424 | - | ||
| 425 | - for (Iterator<Entry<String, ReceiptCompleteData>> it = | ||
| 426 | - outstandingReceipts.entrySet().iterator(); | ||
| 427 | - it.hasNext(); ) { | ||
| 428 | - Map.Entry<String, ReceiptCompleteData> receipt = it.next(); | ||
| 429 | - // If receipt is complete then add to outstandingBatch to process the batch | ||
| 430 | - if (receipt.getValue().isReceiptComplete()) { | ||
| 431 | - it.remove(); | ||
| 432 | - if (pendingMessages.putIfAbsent( | ||
| 433 | - receipt.getKey(), receipt.getValue().getOutstandingMessage().ackHandler) | ||
| 434 | - == null) { | ||
| 435 | - outstandingBatch.add(receipt.getValue().getOutstandingMessage()); | ||
| 423 | + synchronized (outstandingReceipts) { | ||
| 424 | + if (outstandingReceipts.containsKey(ackRequestData.getAckId())) { | ||
| 425 | + outstandingReceipts.get(ackRequestData.getAckId()).notifyReceiptComplete(); | ||
| 426 | + List<OutstandingMessage> outstandingBatch = new ArrayList<>(); | ||
| 427 | + | ||
| 428 | + for (Iterator<Entry<String, ReceiptCompleteData>> it = | ||
| 429 | + outstandingReceipts.entrySet().iterator(); | ||
| 430 | + it.hasNext(); ) { | ||
| 431 | + Map.Entry<String, ReceiptCompleteData> receipt = it.next(); | ||
| 432 | + // If receipt is complete then add to outstandingBatch to process the batch | ||
| 433 | + if (receipt.getValue().isReceiptComplete()) { | ||
| 434 | + it.remove(); | ||
| 435 | + if (pendingMessages.putIfAbsent( | ||
| 436 | + receipt.getKey(), receipt.getValue().getOutstandingMessage().ackHandler) | ||
| 437 | + == null) { | ||
| 438 | + outstandingBatch.add(receipt.getValue().getOutstandingMessage()); | ||
| 439 | + } | ||
| 440 | + } else { | ||
| 441 | + break; | ||
| 436 | 442 | } | |
| 437 | - } else { | ||
| 438 | - break; | ||
| 439 | 443 | } | |
| 444 | + processBatch(outstandingBatch); | ||
| 440 | 445 | } | |
| 441 | - processBatch(outstandingBatch); | ||
| 442 | 446 | } | |
| 443 | 447 | } | |
| 444 | 448 | ||
| 445 | 449 | void notifyAckFailed(AckRequestData ackRequestData) { | |
| 446 | - outstandingReceipts.remove(ackRequestData.getAckId()); | ||
| 450 | + synchronized (outstandingReceipts) { | ||
| 451 | + outstandingReceipts.remove(ackRequestData.getAckId()); | ||
| 452 | + } | ||
| 447 | 453 | } | |
| 448 | 454 | ||
| 449 | 455 | private void processBatch(List<OutstandingMessage> batch) { | |
| Back | FazBrowse Home | New Git URL |
0 commit comments