Kafka Connect: Rework commit-coordinator leader election and harden the coordinator - #17450
Open
kumarpritam863 wants to merge 2 commits into
Open
Kafka Connect: Rework commit-coordinator leader election and harden the coordinator#17450kumarpritam863 wants to merge 2 commits into
kumarpritam863 wants to merge 2 commits into
Conversation
added 2 commits
July 31, 2026 17:09
…he coordinator Reworks how the sink elects its single commit Coordinator and hardens the coordinator's fencing, recovery, and shutdown paths. No control-topic wire-format change; exactly-once semantics are preserved. Leader election (level read, no Admin call) - Before: the coordinator was chosen on every rebalance by calling Admin.describeConsumerGroups(connectGroupId) and picking the task that owned the globally-lowest (topic, partition) from a transient, mid-rebalance member snapshot. This needed a DESCRIBE ACL, added an Admin round-trip to the rebalance path, and was racy during cooperative rebalancing. - After: the leader is the task whose assignment contains partition 0 of the lexicographically-smallest topic in its own consumer subscription() -- connector-wide, identical across tasks, and valid for both a topics list and a topics.regex. No Admin call. Leadership is read as a level on the task thread: open()/close() only flag that a reconcile is needed, and save() reconciles once (start the coordinator if this task owns the leader partition, stop it otherwise). Keeping start/stop off the rebalance callback avoids blocking RPCs there and stops an eager rebalance from needlessly restarting a still-leading coordinator (which would discard in-flight commit state). - The coordinator's commit-readiness partition count is derived from consumer.partitionsFor(...) over the subscribed topics instead of a member-assignment snapshot. - Adds a Committer.configure(Catalog, IcebergSinkConfig, SinkTaskContext) lifecycle hook (invoked from IcebergSinkTask.start) for one-time setup. Coordinator hardening - Zombie fencing via a fixed transactional.id. The coordinator producer id is now connectGroupId-connectorName-coord -- identical across a connector's tasks and stable across restarts -- so a newly elected coordinator's initTransactions() epoch-fences a prior (zombie) coordinator's control-plane writes. Worker ids are unchanged. This fences the brief two-coordinator overlap that the level-read handoff can create (the losing task stops on its next save()). - Fenced != fatal. A coordinator that terminates because it was fenced (ProducerFenced / InvalidProducerEpoch / UnknownProducerId, matched across the cause chain) is cleared without failing the task; any other termination still fails the task. - Recovery reads from earliest. The coordinator's -coord consumer group defaults to auto.offset.reset=earliest, so a fresh or expired group re-reads uncommitted control events instead of skipping to the log end. Replay is idempotent via the snapshot offset floor + distinctByKey(location) dedup + the offsets compare-and-swap. - Idempotency filter in Channel.consumeAvailable: control-topic records at or below the already-consumed per-partition offset are skipped, so a re-delivered or rewound record is never re-buffered or re-counted. - Transient Kafka commit errors from commitConsumerOffsets (Iceberg CommitFailedException, Kafka CommitFailedException, RebalanceInProgressException, RetriableException) are retried rather than failing the task; the table commit already succeeded. - Bounded, interrupt-safe shutdown: stopCoordinator clears state first, then bounded-joins the thread; terminate() failures are best-effort and interrupts are preserved. - KafkaClientFactory.createProducer closes the producer if initTransactions() throws. Compatibility - No config or control-topic wire-format changes. - Exactly-once unchanged -- anchored on the Iceberg offsets compare-and-swap + file-location dedup; the fixed transactional.id adds control-plane fencing on top. - Known edge (OCC-safe): with topics.regex, a lexicographically-smaller topic appearing later can briefly run two coordinators if the losing task receives no rebalance callback. Testing - Unit tests for the election key (leaderPartition), the retryable-commit classification, and coordinator fenced-vs-fatal termination; integration suite (13 tests) passes. - Raised the integration-test commit-wait from 30s to 60s to reduce CI-load flakiness.
Contributor
Author
|
@laskoviymishka @danicafine @bryanck can you please review. |
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
Reworks how the Iceberg Kafka Connect sink elects its single commit
Coordinatorandhardens the coordinator's fencing, recovery, and shutdown paths. No control-topic
wire-format change; exactly-once semantics are preserved.
Leader election — level read, no Admin call
Admin.describeConsumerGroups(connectGroupId)and picking the task that owned theglobally-lowest
(topic, partition)from a transient, mid-rebalance member snapshot.This required a
DESCRIBEACL, added an Admin round-trip to the rebalance path, and wasracy during cooperative rebalancing.
assignment()contains partition 0 of thelexicographically-smallest topic in its own
subscription()— connector-wide,identical across tasks, and valid for both a
topicslist and atopics.regex. No Admincall.
open/closeonly flag that areconcile is needed;
save()reconciles once — start the coordinator if this task ownsthe leader partition, stop it otherwise (both idempotent). Keeping start/stop off the
rebalance callback avoids blocking RPCs there and stops an eager rebalance from
needlessly restarting a still-leading coordinator (which would discard in-flight commit
state).
consumer.partitionsFor(...)over the subscribed topics instead of a member snapshot.Committer.configure(Catalog, IcebergSinkConfig, SinkTaskContext)lifecycle hook(invoked from
IcebergSinkTask.start) for one-time setup.Coordinator hardening
transactional.id. The coordinator producer id is now<connectGroupId>-<connectorName>-coord— identical across a connector's tasks andstable across restarts — so a newly elected coordinator's
initTransactions()epoch-fences a prior (zombie) coordinator's control-plane writes. Worker ids are
unchanged. This fences the brief two-coordinator overlap the level-read handoff can
create (the losing task stops on its next
save()).(
ProducerFenced/InvalidProducerEpoch/UnknownProducerId, matched across thecause chain) is cleared without failing the task; any other termination still fails the
task.
earliest. The-coordconsumer group defaults toauto.offset.reset=earliest, so a fresh or expired group re-reads uncommitted controlevents instead of skipping to the log end. Replay is idempotent via the snapshot offset
floor +
distinctByKey(location)dedup + the offsets compare-and-swap.Channel.consumeAvailable. Control-topic records at or belowthe already-consumed per-partition offset are skipped, so a re-delivered or rewound
record is never re-buffered or re-counted.
commitConsumerOffsets(Iceberg
CommitFailedException, KafkaCommitFailedException,RebalanceInProgressException,RetriableException) are retried rather than failing thetask — the table commit already succeeded.
stopCoordinatorclears state first, thenbounded-joins the thread;
terminate()failures are best-effort and interrupts arepreserved.
KafkaClientFactory.createProducercloses theproducer if
initTransactions()throws.Compatibility
file-location dedup; the fixed
transactional.idadds control-plane fencing on top.Testing
leaderPartition), the retryable-commit classification,and coordinator fenced-vs-fatal termination.