SAlexandru opened a new pull request, #29389:
URL: https://github.com/apache/flink/pull/29389
Backport of #29218 to release-1.20.
The production change is identical: the per-producer state map in
FutureCompletingBlockingQueue and releaseProducer() in the SplitFetcherManager
shutdown hook.
Test adaptations, because release-1.20 doesn't have FLINK-35924 (no
records-processed latch / shutdown synchronization batch in SplitFetcher):
- testFetcherShutdownReleasesWakeupStateAcrossLifecycles creates the
producer state with an explicit wakeUpPuttingThread instead of relying on the
synchronization batch, and uses CommonTestUtils.waitUntilCondition
(TestUtils.waitUntil came with FLINK-35924).
- testFailedFetcherIsNotReapedAndIsCleanedUpOnClose is omitted: without
the latch, a failed fetcher closes its reader and runs its shutdown hook
immediately, so it cannot be stranded.
- The new tests are public, since the 1.20 test classes are JUnitĀ 4.
--
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]