Skip to content

Ensure lifted committers carry a reference to root committer so they're not duplicated - #1531

Merged
vlovgr merged 2 commits into
typelevel:mainfrom
alejandrohdezma:fix/committer-identity
Sep 30, 2026
Merged

vlovgr merged 2 commits into
typelevel:mainfrom
alejandrohdezma:fix/committer-identity

Conversation

@alejandrohdezma

Copy link
Copy Markdown
Contributor

New mapK function introduced in #1522 creates a new committer on every call and its Eq instance is based on equals so every mapK generates a new committer equal to nothing else. The problem is CommittableOffsetBatch groups by committer, so a lifted stream like the one we create in trace4cats-kafka fell back to one comit per record, in parallel.

This solves it by adding a "pointer" to the source committer, so that can be used for hashCode and equals.

…re not duplicated

New `mapK` function introduced in typelevel#1522 creates a new committer on every call and its `Eq` instance is based on `equals` so every mapK generates a new committer equal to nothing else. The problem is CommittableOffsetBatch groups by committer, so a lifted stream like the one we create in trace4cats-kafka fell back to one comit per record, in parallel.

This solves it by adding a "pointer" to the source committer, so that can be used for `hashCode` and `equals`.
The class is sealed with a package-private constructor so no one outside this repo can implement it.
@mergify

mergify Bot commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

This pull request does not currently match the merge queue conditions, so it cannot be queued from here. The box comes back if it matches again.

@vlovgr

vlovgr commented Sep 30, 2026

Copy link
Copy Markdown
Member

Thanks @alejandrohdezma!

@vlovgr
vlovgr merged commit 61c763e into typelevel:main Sep 30, 2026
8 checks passed
@alejandrohdezma
alejandrohdezma deleted the fix/committer-identity branch September 30, 2026 10:33
@vlovgr

vlovgr commented Sep 30, 2026

Copy link
Copy Markdown
Member

@alejandrohdezma v4.1.2 is on it's way to Maven Central with this and #1530.

@alejandrohdezma

Copy link
Copy Markdown
Contributor Author

Thank you @vlovgr!

scala-steward pushed a commit to scala-steward/trace4cats-kafka that referenced this pull request Oct 1, 2026
fs2-kafka 4.1.2 (typelevel/fs2-kafka#1531) makes a committer produced by
`KafkaCommitter#mapK` equal to the committer it was derived from, so
`CommittableOffsetBatch` merges lifted offsets again. `injectK` no longer
needs to cache one lifted committer per consumer, and the
`fs2.kafka.LiftedCommittableOffset` shim that rebuilt offsets around the
cached committer goes with it.

4.1.2 also fixes the deadlock on `KafkaConsumer#unsubscribe` once partitions
are assigned (typelevel/fs2-kafka#1530).
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.

2 participants