Savonitar opened a new pull request, #322:
URL: https://github.com/apache/flink-connector-kafka/pull/322
<!--
*Thank you very much for contributing to the Apache Flink Kafka connector -
we are happy that you want to help us improve Flink. To help the community
review your contribution in the best possible way, please go through the
checklist below, which will get the contribution into a shape in which it can
be best reviewed.*
*Please understand that we do not do this to make contributions to Flink a
hassle. In order to uphold a high standard of quality for code contributions,
while at the same time managing a large number of contributions, we need
contributors to prepare the contributions well, and give reviewers enough
contextual information for the review. Please also understand that
contributions that do not follow this guide will take longer to review and thus
typically be picked up with lower priority by the community.*
## Contribution Checklist
- Make sure that the pull request corresponds to a [JIRA
issue](https://issues.apache.org/jira/projects/FLINK/issues). Exceptions are
made for typos in JavaDoc or documentation files, which need no JIRA issue.
- Name the pull request in the form "[FLINK-XXXX] [component] Title of the
pull request", where *FLINK-XXXX* should be replaced by the actual issue
number. Skip *component* if you are unsure about which is the best component.
Typo fixes that have no associated JIRA issue should be named following
this pattern: `[hotfix] [docs] Fix typo in event time introduction` or
`[hotfix] [javadocs] Expand JavaDoc for PuncuatedWatermarkGenerator`.
- Fill out the template below to describe the changes contributed by the
pull request. That will give reviewers the context they need to do the review.
- Make sure that the change passes the automated tests, i.e., `./mvnw
clean verify` passes. GitHub Actions runs the same build for every push and
pull request against the Flink versions and JDKs listed in
`.github/workflows/push_pr.yml`.
- Each pull request should address only one issue, not mix up code from
multiple issues.
- Each commit in the pull request has a meaningful commit message
(including the JIRA id)
- Once all items of the checklist are addressed, remove the above text and
this checklist, leaving only the filled out template below.
**(The sections below can be removed for hotfixes of typos)**
-->
## What is the purpose of the change
Adds deterministic integration coverage for exactly-once KafkaSink recovery
with pooled transactional IDs. The test verifies that restoring an earlier
checkpoint after a transaction ID has been reused preserves committed records
and successfully commits replayed records without loss or duplication.
It is a usefull test that will prevent further regressions in this area.
## Brief change log
- Add KafkaSinkRecoveryITCase with controlled record emission and manually
triggered checkpoints.
- Complete checkpoint C1, keep C2 and C3 incomplete, and verify through
Kafka that C1’s transactional ID has been reused.
- Cancel the job and restore C1 with the same transactional ID prefix and
operator identities.
- Assert the committed output after recovery and verify that the final
records appear exactly once.
## Verifying this change
This change added tests and can be verified as follows:
- Added
KafkaSinkRecoveryITCase#recoversPooledTransactionsWithoutLosingNewRecords.
- Passed three consecutive runs using the standard Maven/Testcontainers
lifecycle with JDK 17 and Flink 2.2.1:./mvnw test -pl flink-connector-kafka
-Dtest=KafkaSinkRecoveryITCase
- Spotless and Checkstyle checks passed.
## Does this pull request potentially affect one of the following parts:
- Dependencies (does it add or upgrade a dependency, including
`kafka.version` or `flink.version`): (no)
- The public API, i.e., is any changed class annotated with
`@Public(Evolving)` or `@Experimental`, or are the Table options or the PyFlink
wrappers changed: (no)
- Checkpointed state, its serializers, or exactly-once delivery (splits,
enumerator state, writer state, committables, transactions): (yes, adds
additional test coverage for exactly-once delivery)
- The per-record code paths (split reader, record emitter, writer,
serialization schemas; performance sensitive): (no)
## Documentation
- Does this pull request introduce a new feature? (no)
- If yes, how is the feature documented? (not applicable)
- If the docs changed, are both `docs/content` and `docs/content.zh`
updated? (not applicable)
---
##### Was generative AI tooling used to co-author this PR?
<!--
If generative AI tooling has been used in the process of authoring this PR,
please
change the checkbox below to `[X]` and replace the placeholder in the
"Generated-by"
line with the tool name and version. Otherwise remove the "Generated-by"
line.
See the ASF Generative Tooling Guidance for details:
https://www.apache.org/legal/generative-tooling.html
You are responsible for the quality and correctness of every change in this
PR
regardless of the tooling used. Low-effort AI-generated PRs will be closed.
See
AGENTS.md for the full guidance.
-->
- [x] Yes (please specify the tool below)
Generated-by: Claude Fable 5.1
--
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]