Skip to content

KAFKA-20684 [4/N]: Migrate Connect to RebalanceListener - #23105

Open
adikou wants to merge 5 commits into
apache:trunkfrom
adikou:akousik/KAFKA-20684-connect
Open

adikou wants to merge 5 commits into
apache:trunkfrom
adikou:akousik/KAFKA-20684-connect

Conversation

@adikou

@adikou adikou commented Aug 6, 2026 •

Copy link
Copy Markdown
Contributor

WorkerSinkTask's HandleRebalance becomes a RebalanceListener, registered
via setRebalanceListener. Tests capture the listener from
setRebalanceListener rather than from subscribe. Callback bodies keep
using the enclosing consumer field.

Reviewers: Andrew Schofield aschofield@confluent.io, Chia-Ping Tsai chia7712@gmail.com

WorkerSinkTask's HandleRebalance becomes a RebalanceListener, registered via
setRebalanceListener. Tests capture the listener from setRebalanceListener rather
than from subscribe. Callback bodies keep using the enclosing consumer field.
@github-actions

Copy link
Copy Markdown

A label of 'needs-attention' was automatically added to this PR in order to raise the
attention of the committers. Once this issue has been triaged, the triage label
should be removed to prevent this automation from happening again.

@muralibasani

Copy link
Copy Markdown
Contributor

@adikou Thanks for the PR.
Can you pls look into the conflicts ?

@AndrewJSchofield
AndrewJSchofield self-requested a review October 6, 2026 06:57
@adikou
adikou force-pushed the akousik/KAFKA-20684-connect branch from 0a94faf to 8d2416a Compare October 6, 2026 18:27
private class HandleRebalance implements RebalanceListener {
@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
public void onPartitionsAssigned(Collection<TopicPartition> partitions, RebalanceConsumer rebalanceConsumer) {

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.

rebalanceConsumer argument is unused here and in below overridden methods. Is it deliberate?

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

I would say so. If we had Java 22, I think it would be appropriate to use the _ for an unused argument.


doAnswer((Answer<ConsumerRecords<byte[], byte[]>>) invocation -> {
rebalanceListener.getValue().onPartitionsRevoked(INITIAL_ASSIGNMENT);
rebalanceListener.getValue().onPartitionsRevoked(INITIAL_ASSIGNMENT, null);

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.

Do you see any implications of passing null here, and all the below calls ?
As at this moment, that argument is ignored. But later when it's implemented, probably this needs a change.

@github-actions github-actions Bot removed the triage PRs from the community label Oct 8, 2026
@chia7712

chia7712 commented Oct 8, 2026

Copy link
Copy Markdown
Member

@adikou Would you please rebase code?

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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants