Skip to content

KAFKA-20503: Add integration tests for transaction buffering - #22958

Open
nicktelford wants to merge 1 commit into
apache:trunkfrom
nicktelford:KIP-892/eos-txn-buffering-tests
Open

KAFKA-20503: Add integration tests for transaction buffering#22958
nicktelford wants to merge 1 commit into
apache:trunkfrom
nicktelford:KIP-892/eos-txn-buffering-tests

Conversation

@nicktelford

Copy link
Copy Markdown
Contributor

We need to verify that, when transactional state stores are enabled and processing.guarantee is exactly_once_v2, records are buffered in the transaction buffer (invisible to a READ_COMMITTED reader) until the Streams commit cycle completes, and then committed to the store as expected.

This reuses EosIntegrationTest, which already owns the EOS invariant machinery and a transactionalStateStores parameter. A new shouldBufferStateStoreWritesUntilCommitUnderEos test drives three bursts of writes across a commit boundary and asserts the READ_COMMITTED/READ_UNCOMMITTED store views at each step, via a new isolation-aware verifyStateStore/queryStateStore and a waitForStateStore helper that polls for the buffer flush (context.commit() only requests a commit; the actual flush happens asynchronously afterwards).

The transactional dimension is orthogonal to group protocol and processing-threads, so — consistent with this file's existing sparse-matrix convention — the new test and the existing transactional parameterization are each exercised via a single representative combination rather than the full matrix, to avoid unnecessary integration-test runtime.

We need to verify that, when transactional state stores are enabled
and processing.guarantee is exactly_once_v2, records are buffered in
the transaction buffer (invisible to a READ_COMMITTED reader) until
the Streams commit cycle completes, and then committed to the store
as expected.

Reuses EosIntegrationTest, which already owns the EOS invariant
machinery: a new shouldBufferStateStoreWritesUntilCommitUnderEos test
drives three bursts of writes across a commit boundary and asserts
the READ_COMMITTED/READ_UNCOMMITTED store views at each step via a
new isolation-aware verifyStateStore/queryStateStore, waiting on the
async buffer flush with a new waitForStateStore helper (context.commit()
only requests a commit; the buffer flush happens afterwards).

The transactional dimension is orthogonal to group protocol and
processing-threads, so the new test and existing sparse transactional
parameterization are each exercised via a single representative
combination rather than the full matrix, to avoid unnecessary
integration-test runtime.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
@nicktelford

Copy link
Copy Markdown
Contributor Author

@bbejeck

@github-actions github-actions Bot added streams tests Test fixes (including flaky tests) triage PRs from the community labels Jul 27, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

streams tests Test fixes (including flaky tests) triage PRs from the community

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant