nahidupa opened a new pull request, #18322:
URL: https://github.com/apache/iceberg/pull/18322
## Problem
The coordinator's source partition count was fixed at construction. If the
source assignment expands while that coordinator remains active, reports from
the old partitions can admit a full commit before the added partition reports.
The resulting snapshot and completion event can carry a completeness timestamp
that does not cover the current assignment.
Refreshing only once per cycle is insufficient: an unavailable or unstable
description can preserve old eligibility, assignments can change before full
admission, and a late metadata result can overwrite a callback's invalidation.
## Changes
- Read a stable, nonempty source-group description through the existing
owned Admin client before each commit cycle.
- Pin an immutable member-to-topic-partition assignment snapshot for that
cycle and refresh its expected partition count.
- Recheck the complete assignment immediately before admitting a full
commit, including equal-size partition replacements and ownership changes.
- Invalidate through `CoordinatorThread` from both assignment and revocation
callbacks, including callbacks that retain the coordinator.
- Bracket metadata reads with an atomic revision so a successful late
response cannot restore invalidated eligibility.
- Keep buffered responses and the existing timeout-partial path. An unknown
or changed cycle can commit partially without a full completeness watermark,
but cannot regain full eligibility until a later verified cycle.
No extra Admin client, source-consumer access from the coordinator thread,
event schema, configuration, public API, or election/fencing protocol is
introduced.
## Partial Commits And Limits
When assignment verification fails or the cycle is invalidated, timeout can
still commit buffered files. Such a partial snapshot omits
`kafka.connect.valid-through-ts`, and completion events carry a null timestamp.
A later cycle with freshly verified assignments can complete fully.
This verifies the assignments observed through full-commit admission. It is
not a distributed generation fence: assignments can change after admission, and
already-published table snapshots cannot be retracted. Multi-table commits are
not atomic. Snapshot-expiration recovery and bounded startup replay are outside
this fix.
This standalone branch intentionally retains main's existing readiness-count
mechanism. The distinct expected-partition coverage fix belongs to #17925; this
PR does not duplicate that unmerged change.
## Validation
Fresh branch from upstream `main` at
`48330b8dacab6662242d252b39c8444190979bb2`.
Standalone head: `216658fc1723875772a849d390f36605ca0610ab`.
- **17** new parameterized regression cases cover expansion,
unknown/unstable/empty metadata after a successful cycle, failed full
verification, equal-size topology changes, notifications during metadata reads,
and retained-coordinator added-only/full opens and non-leader closes.
- Assertions inspect real Iceberg snapshots and added files, exact control
checkpoints, commit IDs, and full versus null partial timestamps, including
later verified recovery.
- The old-partitions-ready/new-partition-missing regression failed on
unmodified main with a premature full snapshot and passed after the fix.
- Disposable negative controls failed behavior assertions: frozen
constructor count **1/1**, omitted callback effect **3/3**, weakened revision
bracketing **2/2**, omitted pre-full topology comparison **3/3**. Exact source
restoration was verified and all 17 positives rerun. Compilation failures were
not counted as behavioral evidence.
- Standalone full connector gate: **160 tests, 21 suites**, zero failures,
errors, or skips. Explicit `spotlessCheck`, Checkstyle main/test, class
uniqueness, and whitespace checks passed.
```sh
./gradlew -DsparkVersions= -DflinkVersions= -DkafkaVersions=3 \
:iceberg-kafka-connect:iceberg-kafka-connect:spotlessCheck \
:iceberg-kafka-connect:iceberg-kafka-connect:check \
-x integrationTest --rerun-tasks --console=plain
```
Tests use Corretto 21.0.3, real in-memory Iceberg tables, Kafka mock
clients, and controlled callback interleavings. They are not live-broker
expansion/rebalance tests, parallel-thread stress tests, or local JDK 17
validation. GitHub CI for the published head is separate.
## Compatibility And Related Work
#17925 remains contributor-approved but unmerged and was left unchanged. In
an unpublished integration branch, its expected identity set was derived from
this fix's verified cycle snapshot, preserving identity coverage and
expected-member overlap. Matching Admin fixtures were added without weakening
assertions. That combination passed **169 tests, 21 suites** at
`ae6b264892887b442bd35abc1c8d00dd37642546`.
The unpublished combination with #18006 and #18012 passed **190 tests, 24
suites** at `25771fca3aa119c40a52423a67f0346d20f1aa6f`, after wiring the
recovery/lookup metadata fixtures and the consumer factory overload. Those
integration adaptations will be needed when the independent changes meet; the
combined branch is not part of this PR.
#17713 is merged and its channel replay guard is preserved. #17450 is closed
unmerged and is not a dependency. The three bug fixes remain on separate
branches and PRs.
---
**AI Disclosure**
- Model: [unknown - human to fill in]
- Platform/Tool: GitHub Copilot.
- Human Oversight: unreviewed.
- Prompt Summary: Implement assignment freshness on a fresh independent
branch, retain timeout partial commits while blocking unverified full
watermarks, apply approved readiness-review testing lessons, and validate
standalone plus unpublished combinations without modifying the approved PR.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]