Skip to content

Kafka support manual partition pause/resume with pause-if-no-requests=true - #3538

Open
ozangunalp wants to merge 1 commit into
smallrye:mainfrom
ozangunalp:kafka_manual_partition_pause_resume
Open

ozangunalp wants to merge 1 commit into
smallrye:mainfrom
ozangunalp:kafka_manual_partition_pause_resume

Conversation

@ozangunalp

Copy link
Copy Markdown
Collaborator

Fixes #3295

@ozangunalp
ozangunalp force-pushed the kafka_manual_partition_pause_resume branch from 7c4e746 to 96a4098 Compare September 25, 2026 07:01
@codecov

codecov Bot commented Sep 25, 2026

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 81.81818% with 6 lines in your changes missing coverage. Please review.
✅ Project coverage is 77.93%. Comparing base (a96442f) to head (96a4098).
⚠️ Report is 1257 commits behind head on main.

Files with missing lines Patch % Lines
...ve/messaging/kafka/impl/ReactiveKafkaConsumer.java 81.25% 1 Missing and 5 partials ⚠️
Additional details and impacted files

Impacted file tree graph

@@             Coverage Diff              @@
##               main    #3538      +/-   ##
============================================
+ Coverage     77.47%   77.93%   +0.46%     
- Complexity     3778     5641    +1863     
============================================
  Files           306      487     +181     
  Lines         12673    18842    +6169     
  Branches       1648     2319     +671     
============================================
+ Hits           9818    14685    +4867     
- Misses         2116     3009     +893     
- Partials        739     1148     +409     
Files with missing lines Coverage Δ
...ye/reactive/messaging/kafka/i18n/KafkaLogging.java 100.00% <ø> (ø)
...ctive/messaging/kafka/impl/RebalanceListeners.java 92.30% <100.00%> (+6.34%) ⬆️
...ve/messaging/kafka/impl/ReactiveKafkaConsumer.java 80.33% <81.25%> (-1.61%) ⬇️

... and 292 files with indirect coverage changes

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

@ozangunalp
ozangunalp force-pushed the kafka_manual_partition_pause_resume branch from 96a4098 to 4695f0e Compare September 29, 2026 14:29
@cescoffier cescoffier changed the title Kafka support manual partition pause/resume with pause-if-requests=true Kafka support manual partition pause/resume with pause-if-no-requests=true Sep 30, 2026
}

void removeManuallyPausedPartitions(Collection<TopicPartition> partitions) {
if (manuallyPausedPartitions.removeAll(partitions)) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is this correct? removeAll returns true if it removes at least one partition from the collection. I feel like you want to remove all of them. If partitions {tp0, tp1, tp2} are revoked but only tp0 was manually paused, the log says "Cleaning up paused partitions [tp0, tp1, tp2] due to revocation". Slightly misleading, no?.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nice catch, we should log the intersection if there are any.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Instead I added both revoked and paused partitions to the log.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't believe it's correct.
removeAll mutates manuallyPausedPartitions before the log call. The log message template is: "Cleaning up manually paused partitions %s due to revocation of %s".

The first %s receives the remaining set, not the cleaned-up set.

I think you need:

void removeManuallyPausedPartitions(Collection<TopicPartition> partitions) {
    Set<TopicPartition> removed = new HashSet<>(manuallyPausedPartitions);
    removed.retainAll(partitions);
    if (!removed.isEmpty()) {
        manuallyPausedPartitions.removeAll(removed);
        log.cleaningUpPausedPartitions(removed, partitions);
    }
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

or update the log messsage: "Remaining manually paused partitions %s after revocation of %s"

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, let's change the log message.

@ozangunalp
ozangunalp force-pushed the kafka_manual_partition_pause_resume branch from 4695f0e to ace0a93 Compare September 30, 2026 08:53
@ozangunalp
ozangunalp requested a review from cescoffier October 6, 2026 14:11
@ozangunalp
ozangunalp force-pushed the kafka_manual_partition_pause_resume branch from ace0a93 to ee98dc5 Compare October 7, 2026 07:29
@ozangunalp
ozangunalp force-pushed the kafka_manual_partition_pause_resume branch from ee98dc5 to 3f3a63f Compare October 8, 2026 09:52

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Kafka manual partition pausing with pause-if-no-requests backpressure support

2 participants