Skip to content

Partition removal state management and concurrency improvements - #190

Merged
Antony Stubbs (astubbs) merged 2 commits into
confluentinc:masterfrom
astubbs:partition-removal-state
Feb 25, 2022
Merged

Antony Stubbs (astubbs) merged 2 commits into
confluentinc:masterfrom
astubbs:partition-removal-state

Conversation

@astubbs

@astubbs Antony Stubbs (astubbs) commented Feb 11, 2022 •

Copy link
Copy Markdown
Contributor

Relates:

@astubbs

Copy link
Copy Markdown
Contributor Author

Niels Oertel (@nioertel) you can try running against this snapshot - should clear up the issues. In the team time I want to try to capture the previous failures in some tests and clean it up.

@astubbs Antony Stubbs (astubbs) left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

draft complete

@nioertel

Niels Oertel (nioertel) commented Feb 22, 2022 •

Copy link
Copy Markdown
Contributor

Hi Anthony, we performed some performance testing and observed a problematic behaviour.
Setup:

  • Topic with 24 partitions
  • 4 java processes processing data from the topic (i.e. Kafka distributes the 24 partitions equally to the 4 processes)
  • The 4 java processes are starting up pretty much in parallel
  • There is a small backlog on the topic, so the processes will start working immediately once they're connected to Kafka

Right after startup we see messages like this in the logs:

[WARN ] 2022-02-21 17:51:51.841 [pc-broker-poll]  i.c.p.s.PartitionMonitor - New assignment of partition test.topic-11 which already exists in partition state. Could be a state bug.

I'm not sure if this message has any side effects. It seems all partitions are still handled. But I want to cross-check if this message may cause any issues.

Note that we are using a version of the parallel consumer that is merged from the following branches:

@astubbs

Antony Stubbs (astubbs) commented Feb 22, 2022 •

Copy link
Copy Markdown
Contributor Author

Safe to ignore. I thought I'd removed it actually already. It's hit braise because of the new way state is managed. It's there but marked as previously removed. I'll reassess the path and logging and modify it, but yea safe to ignore. Sorry about that.

@astubbs
Antony Stubbs (astubbs) force-pushed the partition-removal-state branch 4 times, most recently from 3bdc29e to b848d46 Compare February 23, 2022 20:06
@astubbs

Copy link
Copy Markdown
Contributor Author

Btw Niels Oertel (@nioertel) did you really mean #93?

@astubbs
Antony Stubbs (astubbs) marked this pull request as draft February 23, 2022 20:21
@nioertel

Copy link
Copy Markdown
Contributor

Btw Niels Oertel (@nioertel) did you really mean #93?

No, I edited my comment. I meant #193 ;-)

@astubbs
Antony Stubbs (astubbs) force-pushed the partition-removal-state branch 2 times, most recently from 6c751cd to 2112069 Compare February 24, 2022 21:51

@astubbs Antony Stubbs (astubbs) left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

lgtm

Antony Stubbs (astubbs) referenced this pull request in astubbs/parallel-consumer Feb 25, 2022
save point - try a no op object for removed partitions

save point - try a no op object for removed partitions - implemented

license

step: npe

concurrency improvements

step - check in is assigned isn't instance of revoked

step: remove state access - otherwise need to return no-op collections

refactorings / improvements

more fixes and refactorings

fixes

review

cleanup

step: fix log message

step: fix: When checking stale epoch, also check if partition still assigned

step: improve impact of no-op removed state, so it wont' be retrieved for analysis

log settings fix

docs

tidyup

Minor test refactoring for clarity

START: Test for PR190 - make number of instances dynamic

step

step: seperate test

refactorings / improvements

test capture - currently backed off though

the rough fixes

step

step: error well caught

step: error well caught

step: error well caught

step: error well caught

step: Increase message count and go back to starting pc1 only after a portion of messages are sent - otherwise in old test, pc2 doesn't get a chance to consume any.

step: Increase message count and go back to starting pc1 only after a portion of messages are sent - otherwise in old test, pc2 doesn't get a chance to consume any.

step: Increase message count and go back to starting pc1 only after a portion of messages are sent - otherwise in old test, pc2 doesn't get a chance to consume any.

step: Increase message count and go back to starting pc1 only after a portion of messages are sent - otherwise in old test, pc2 doesn't get a chance to consume any.

step: Increase message count and go back to starting pc1 only after a portion of messages are sent - otherwise in old test, pc2 doesn't get a chance to consume any.

step: Increase message count and go back to starting pc1 only after a portion of messages are sent - otherwise in old test, pc2 doesn't get a chance to consume any.

step: tweek test settings for slower machines

turn down test logging

turn off unrelated tests

apply doc changes

apply workmanager

apply ShardManager

apply MultiInstanceRebalanceTest

apply PartitionMonitor

step: review

refactor: WorkContainer / Shard removal cleanup

refactor: WorkContainer / Shard removal cleanup

step: apply license

step: license now working

step: refactor isAssigned to isRemoved

step: turn tests back on

step: disable remaining flakey test
Antony Stubbs (astubbs) referenced this pull request in astubbs/parallel-consumer Feb 25, 2022
- Improve impact of no-op removed state, so it won't be retrieved for analysis
- Change from removing revoked PartitionState instances to replacing with a PartitionRemovedState subclass
-- Advantage of this is there will never be NPE from stale references of work still in flight
-- Also means we can track history of assignment
- Several refactorings that have been waiting to happen related to the change
- Improved cohesion
Antony Stubbs (astubbs) referenced this pull request in astubbs/parallel-consumer Feb 25, 2022
- Improve impact of no-op removed state, so it won't be retrieved for analysis
- Change from removing revoked PartitionState instances to replacing with a PartitionRemovedState subclass
-- Advantage of this is there will never be NPE from stale references of work still in flight
-- Also means we can track history of assignment
- Several refactorings that have been waiting to happen related to the change
- Improved cohesion
@astubbs
Antony Stubbs (astubbs) marked this pull request as ready for review February 25, 2022 11:23
Addresses issues reported:

- ConcurrentModificationException in ShardManager during Rebalancing #188
- NullPointerException in PartitionMonitor during Rebalancing #189

Notes:

- Improve impact of no-op removed state, so it won't be retrieved for analysis
- Change from removing revoked PartitionState instances to replacing with a PartitionRemovedState subclass
-- Advantage of this is there will never be NPE from stale references of work still in flight
-- Also means we can track history of assignment
- Several refactorings that have been waiting to happen related to the change
- Improved cohesion

Relates / Future improvements on this:

Refactor: Consider a shared nothing architecture, to reduce thread complexity #200
@astubbs
Antony Stubbs (astubbs) merged commit 3988c68 into confluentinc:master Feb 25, 2022
@astubbs
Antony Stubbs (astubbs) deleted the partition-removal-state branch February 25, 2022 11:48
Antony Stubbs (astubbs) added a commit that referenced this pull request Mar 23, 2022
Addresses issues reported:

- ConcurrentModificationException in ShardManager during Rebalancing #188
- NullPointerException in PartitionMonitor during Rebalancing #189

Notes:

- Improve impact of no-op removed state, so it won't be retrieved for analysis
- Change from removing revoked PartitionState instances to replacing with a PartitionRemovedState subclass
-- Advantage of this is there will never be NPE from stale references of work still in flight
-- Also means we can track history of assignment
- Several refactorings that have been waiting to happen related to the change
- Improved cohesion

Relates / Future improvements on this:

Refactor: Consider a shared nothing architecture, to reduce thread complexity #200
Antony Stubbs (astubbs) added a commit to astubbs/parallel-consumer that referenced this pull request Aug 7, 2026
The issue-reference gate caught three bare numbers below #1000 on added lines,
and it is right to: the fork's numbering sits inside upstream's range, so a bare
number is a coin flip.

- "PR #49" (twice, in the audit and the backlog) is a fork PR: #49.
- "(#190)" was quoted verbatim from a squashed commit subject, which is exactly
  how an upstream number leaks in looking innocent. Reworded so the quote stops
  at the subject and the number is cited as confluentinc#190 outside it.

bin/check-issue-refs.sh now passes on this branch.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01SizQDD2hVUjD7EhESe9Hkb
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.

NullPointerException in PartitionMonitor during Rebalancing ConcurrentModificationException in ShardManager during Rebalancing

2 participants