Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -803,6 +803,10 @@ public void truncateDirtyFiles(long offsetToTruncate) throws RocksDBException {

this.recoverTopicQueueTable();

if (this.timerMessageStore != null) {
this.timerMessageStore.onCommitLogDispatchTruncate(offsetToTruncate);
}

if (!messageStoreConfig.isEnableBuildConsumeQueueConcurrently()) {
this.reputMessageService = new ReputMessageService();
} else {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -784,6 +784,9 @@ public boolean enqueue(int queueId) {
MessageExt msgExt = getMessageByCommitOffset(offsetPy, sizePy);
if (null == msgExt) {
perfCounterTicks.getCounter("enqueue_get_miss");
} else if (!isValidTimerMessage(msgExt)) {
LOGGER.warn("Skip invalid timer message during enqueue, offsetPy:{}, sizePy:{}, topic:{}",
offsetPy, sizePy, msgExt.getTopic());
} else {
lastEnqueueButExpiredTime = System.currentTimeMillis();
lastEnqueueButExpiredStoreTime = msgExt.getStoreTimestamp();
Expand Down Expand Up @@ -1123,6 +1126,56 @@ public int dequeue() throws Exception {
return 1;
}

public void onCommitLogDispatchTruncate(long offsetToTruncate) {
discardTruncatedRequests(enqueuePutQueue, offsetToTruncate);
discardTruncatedRequests(dequeuePutQueue, offsetToTruncate);

List<List<TimerRequest>> retainedLists = new ArrayList<>();
List<TimerRequest> requests;
while ((requests = dequeueGetQueue.poll()) != null) {
List<TimerRequest> retained = new ArrayList<>(requests.size());
for (TimerRequest request : requests) {
if (request.getOffsetPy() >= offsetToTruncate) {
request.idempotentRelease(false);
} else {
retained.add(request);
}
}
if (!retained.isEmpty()) {
retainedLists.add(retained);
}
}
for (List<TimerRequest> retained : retainedLists) {
dequeueGetQueue.offer(retained);
}

ConsumeQueueInterface cq = messageStore.getConsumeQueue(TIMER_TOPIC, 0);
if (cq != null && currQueueOffset > cq.getMaxOffsetInQueue()) {
LOGGER.warn("Timer currQueueOffset:{} is larger than maxOffsetInQueue:{} after CommitLog truncation to {}",
currQueueOffset, cq.getMaxOffsetInQueue(), offsetToTruncate);
currQueueOffset = cq.getMaxOffsetInQueue();
}
if (commitQueueOffset > currQueueOffset) {
commitQueueOffset = currQueueOffset;
}
prepareTimerCheckPoint();
}

private void discardTruncatedRequests(BlockingQueue<TimerRequest> queue, long offsetToTruncate) {
List<TimerRequest> retained = new ArrayList<>();
TimerRequest request;
while ((request = queue.poll()) != null) {
if (request.getOffsetPy() >= offsetToTruncate) {
request.idempotentRelease(false);
} else {
retained.add(request);
}
}
for (TimerRequest requestToRetain : retained) {
queue.offer(requestToRetain);
}
}

private List<List<TimerRequest>> splitIntoLists(List<TimerRequest> origin) {
//this method assume that the origin is not null;
List<List<TimerRequest>> lists = new LinkedList<>();
Expand Down Expand Up @@ -1157,6 +1210,13 @@ private List<List<TimerRequest>> splitIntoLists(List<TimerRequest> origin) {
}

private MessageExt getMessageByCommitOffset(long offsetPy, int sizePy) {
long minPhyOffset = messageStore.getMinPhyOffset();
long maxPhyOffset = messageStore.getMaxPhyOffset();
if (sizePy <= 0 || offsetPy < minPhyOffset || offsetPy > maxPhyOffset - sizePy) {
LOGGER.warn("Skip timer message outside CommitLog range, offsetPy:{}, sizePy:{}, minPhyOffset:{}, maxPhyOffset:{}",
offsetPy, sizePy, minPhyOffset, maxPhyOffset);
return null;
}
for (int i = 0; i < 3; i++) {
MessageExt msgExt = StoreUtil.getMessage(offsetPy, sizePy, messageStore, bufferLocal.get());
if (null == msgExt) {
Expand All @@ -1168,6 +1228,23 @@ private MessageExt getMessageByCommitOffset(long offsetPy, int sizePy) {
return null;
}

boolean isValidTimerMessage(MessageExt msgExt) {
if (msgExt == null || !TIMER_TOPIC.equals(msgExt.getTopic())) {
return false;
}
String realTopic = msgExt.getProperty(MessageConst.PROPERTY_REAL_TOPIC);
String realQueueId = msgExt.getProperty(MessageConst.PROPERTY_REAL_QUEUE_ID);
if (realTopic == null || realTopic.isEmpty() || realQueueId == null || realQueueId.isEmpty()) {
return false;
}
try {
Integer.parseInt(realQueueId);
return true;
} catch (NumberFormatException ignored) {
return false;
}
}

public MessageExtBrokerInner convert(MessageExt messageExt, long enqueueTime, boolean needRoll) {
if (enqueueTime != -1) {
MessageAccessor.putProperty(messageExt, TIMER_ENQUEUE_MS, enqueueTime + "");
Expand Down Expand Up @@ -1717,7 +1794,12 @@ public void run() {
long start = System.currentTimeMillis();
MessageExt msgExt = getMessageByCommitOffset(tr.getOffsetPy(), tr.getSizePy());
if (null != msgExt) {
if (needDelete(tr.getMagic()) && !needRoll(tr.getMagic())) {
if (!isValidTimerMessage(msgExt)) {
LOGGER.warn("Skip invalid timer message during dequeue, offsetPy:{}, sizePy:{}, topic:{}",
tr.getOffsetPy(), tr.getSizePy(), msgExt.getTopic());
tr.idempotentRelease(false);
doRes = true;
} else if (needDelete(tr.getMagic()) && !needRoll(tr.getMagic())) {
if (msgExt.getProperty(MessageConst.PROPERTY_TIMER_DEL_UNIQKEY) != null && tr.getDeleteList() != null) {
//Execute metric plus one for messages that fail to be deleted
addMetric(msgExt, 1);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -265,6 +265,48 @@ public void testRetryUntilSuccess() throws Exception {
verify(mockMessageStore, times(6)).putMessage(any(MessageExtBrokerInner.class));
}

@Test
public void testInvalidTimerMessageIsRejected() throws Exception {
TimerMessageStore timerMessageStore = createTimerMessageStore(null, true);

MessageExt invalidMessage = new MessageExt();
invalidMessage.setTopic("ordinary-topic");
assertFalse(timerMessageStore.isValidTimerMessage(invalidMessage));

MessageExt missingQueueId = new MessageExt();
missingQueueId.setTopic(TimerMessageStore.TIMER_TOPIC);
MessageAccessor.putProperty(missingQueueId, MessageConst.PROPERTY_REAL_TOPIC, "ordinary-topic");
assertFalse(timerMessageStore.isValidTimerMessage(missingQueueId));

MessageExt validMessage = new MessageExt();
validMessage.setTopic(TimerMessageStore.TIMER_TOPIC);
MessageAccessor.putProperty(validMessage, MessageConst.PROPERTY_REAL_TOPIC, "ordinary-topic");
MessageAccessor.putProperty(validMessage, MessageConst.PROPERTY_REAL_QUEUE_ID, "0");
assertTrue(timerMessageStore.isValidTimerMessage(validMessage));
}

@Test
public void testOnCommitLogDispatchTruncateClearsStaleRequests() throws Exception {
TimerMessageStore timerMessageStore = createTimerMessageStore(null, false);
timerMessageStore.load();

TimerRequest staleEnqueueRequest = new TimerRequest(100, 20, 0, 0, 0);
TimerRequest retainedEnqueueRequest = new TimerRequest(50, 20, 0, 0, 0);
timerMessageStore.enqueuePutQueue.offer(staleEnqueueRequest);
timerMessageStore.enqueuePutQueue.offer(retainedEnqueueRequest);

TimerRequest staleDequeueRequest = new TimerRequest(120, 20, 0, 0, 0);
timerMessageStore.dequeuePutQueue.offer(staleDequeueRequest);

timerMessageStore.onCommitLogDispatchTruncate(100);

assertFalse(staleEnqueueRequest.isSucc());
assertFalse(staleDequeueRequest.isSucc());
assertEquals(retainedEnqueueRequest, timerMessageStore.enqueuePutQueue.poll());
assertTrue(timerMessageStore.enqueuePutQueue.isEmpty());
assertTrue(timerMessageStore.dequeuePutQueue.isEmpty());
}

@Test
public void testTimerFlowControl() throws Exception {
String topic = "TimerTest_testTimerFlowControl";
Expand Down