[bugfix] Fix duplicate collector job timeouts#4169
Merged
tomsun28 merged 2 commits intoJul 4, 2026
Merged
Conversation
Contributor
There was a problem hiding this comment.
Pull request overview
Fixes a collector memory/timeout accumulation issue by ensuring that when a job with an existing ID is re-added, the previous HashedWheelTimer timeout is cancelled instead of being left scheduled in the wheel.
Changes:
- Cancel the previous
Timeoutwhen replacing entries incurrentCyclicTaskMap/currentTempTaskMap. - Add regression tests verifying timeout replacement and cancellation behavior for cyclic and temporary jobs.
Reviewed changes
Copilot reviewed 2 out of 2 changed files in this pull request and generated 1 comment.
| File | Description |
|---|---|
| hertzbeat-collector/hertzbeat-collector-common/src/main/java/org/apache/hertzbeat/collector/timer/TimerDispatcher.java | Cancels the previous scheduled timeout when a job ID is replaced to avoid accumulating pending wheel timeouts. |
| hertzbeat-collector/hertzbeat-collector-common/src/test/java/org/apache/hertzbeat/collector/timer/TimerDispatcherTest.java | Adds regression tests that validate the old timeout is cancelled and the new timeout remains active for both job types. |
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
Comment on lines
+261
to
+271
| @SuppressWarnings("unchecked") | ||
| private Map<Long, Timeout> currentTempTaskMap() { | ||
| try { | ||
| Field field = TimerDispatcher.class.getDeclaredField("currentTempTaskMap"); | ||
| field.setAccessible(true); | ||
| return (Map<Long, Timeout>) field.get(timerDispatcher); | ||
| } catch (ReflectiveOperationException e) { | ||
| throw new AssertionError(e); | ||
| } | ||
| } | ||
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
What's changed?
Fixes #4153.
When the collector receives the same job id again,
TimerDispatcherreplaced the timeout stored in the task map but did not cancel the previousHashedWheelTimertimeout. The replaced timeout could stay active until it fired, so repeated job updates may accumulate unnecessaryHashedWheelTimeoutinstances.This change cancels the previous timeout when a cyclic or temporary job is replaced by a new timeout for the same job id. It keeps the latest timeout in the map and only cancels the old one.
Added regression coverage for both cyclic and temporary jobs to verify that:
Validation:
git diff --check -- hertzbeat-collector/hertzbeat-collector-common/src/main/java/org/apache/hertzbeat/collector/timer/TimerDispatcher.java hertzbeat-collector/hertzbeat-collector-common/src/test/java/org/apache/hertzbeat/collector/timer/TimerDispatcherTest.java./mvnw -pl hertzbeat-collector/hertzbeat-collector-common -am -Dsurefire.failIfNoSpecifiedTests=false testChecklist
Add or update API
No API is added or updated in this PR.