mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-01 05:46:12 +08:00
Fixed Transactional message will be lost under extreme condition
This commit is contained in:
committed by
von gosling
parent
f27dc75511
commit
5d3560dc5c
+18
-25
@@ -198,7 +198,7 @@ public class TransactionalMessageServiceImpl implements TransactionalMessageServ
|
||||
if (null != checkImmunityTimeStr) {
|
||||
checkImmunityTime = getImmunityTime(checkImmunityTimeStr, transactionTimeout);
|
||||
if (valueOfCurrentMinusBorn < checkImmunityTime) {
|
||||
if (checkPrepareQueueOffset(removeMap, doneOpOffset, msgExt, checkImmunityTime)) {
|
||||
if (checkPrepareQueueOffset(removeMap, doneOpOffset, msgExt)) {
|
||||
newOffset = i + 1;
|
||||
i++;
|
||||
continue;
|
||||
@@ -315,33 +315,26 @@ public class TransactionalMessageServiceImpl implements TransactionalMessageServ
|
||||
* @param removeMap Op message map to determine whether a half message was responded by producer.
|
||||
* @param doneOpOffset Op Message which has been checked.
|
||||
* @param msgExt Half message
|
||||
* @param checkImmunityTime User defined time to avoid being detected early.
|
||||
* @return Return true if put success, otherwise return false.
|
||||
*/
|
||||
private boolean checkPrepareQueueOffset(HashMap<Long, Long> removeMap, List<Long> doneOpOffset, MessageExt msgExt,
|
||||
long checkImmunityTime) {
|
||||
if (System.currentTimeMillis() - msgExt.getBornTimestamp() < checkImmunityTime) {
|
||||
String prepareQueueOffsetStr = msgExt.getUserProperty(MessageConst.PROPERTY_TRANSACTION_PREPARED_QUEUE_OFFSET);
|
||||
if (null == prepareQueueOffsetStr) {
|
||||
return putImmunityMsgBackToHalfQueue(msgExt);
|
||||
} else {
|
||||
long prepareQueueOffset = getLong(prepareQueueOffsetStr);
|
||||
if (-1 == prepareQueueOffset) {
|
||||
return false;
|
||||
} else {
|
||||
if (removeMap.containsKey(prepareQueueOffset)) {
|
||||
long tmpOpOffset = removeMap.remove(prepareQueueOffset);
|
||||
doneOpOffset.add(tmpOpOffset);
|
||||
return true;
|
||||
} else {
|
||||
return putImmunityMsgBackToHalfQueue(msgExt);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
private boolean checkPrepareQueueOffset(HashMap<Long, Long> removeMap, List<Long> doneOpOffset,
|
||||
MessageExt msgExt) {
|
||||
String prepareQueueOffsetStr = msgExt.getUserProperty(MessageConst.PROPERTY_TRANSACTION_PREPARED_QUEUE_OFFSET);
|
||||
if (null == prepareQueueOffsetStr) {
|
||||
return putImmunityMsgBackToHalfQueue(msgExt);
|
||||
} else {
|
||||
return true;
|
||||
long prepareQueueOffset = getLong(prepareQueueOffsetStr);
|
||||
if (-1 == prepareQueueOffset) {
|
||||
return false;
|
||||
} else {
|
||||
if (removeMap.containsKey(prepareQueueOffset)) {
|
||||
long tmpOpOffset = removeMap.remove(prepareQueueOffset);
|
||||
doneOpOffset.add(tmpOpOffset);
|
||||
return true;
|
||||
} else {
|
||||
return putImmunityMsgBackToHalfQueue(msgExt);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user