Sylwester Lachiewicz created FLINK-40362:
--------------------------------------------

             Summary: DynamicKafkaSource idle reader never finishes when 
no-more-splits precedes the metadata update
                 Key: FLINK-40362
                 URL: https://issues.apache.org/jira/browse/FLINK-40362
             Project: Flink
          Issue Type: Bug
          Components: Connectors / Kafka
            Reporter: Sylwester Lachiewicz


An idle reader (a subtask with no assigned splits) that receives the 
no-more-splits signal before the {{MetadataUpdateEvent}} never forwards that 
signal to the sub-readers created by the update.

{{DynamicKafkaSourceReader.handleSourceEvents}} guards the re-notification with 
{{!pendingSplits.isEmpty()}}, which is false in exactly this case, so 
{{notifyNoMoreSplits()}} is never re-delivered. The reader then returns 
{{NOTHING_AVAILABLE}} indefinitely and a bounded job never finishes.

Reproduced deterministically at unit level: 
{{DynamicKafkaSourceReaderTest#testIdleReaderFinishesWhenNoMoreSplitsArrivesBeforeMetadata}}
 calls {{notifyNoMoreSplits()}} before {{handleSourceEvents(metadata)}} and 
then polls for {{END_OF_INPUT}}. On main it fails with "expected: END_OF_INPUT 
but was: NOTHING_AVAILABLE"; with the fix it passes.

Fix: re-deliver the no-more-splits signal after processing a metadata update 
regardless of whether pending splits exist. {{notifyNoMoreSplits()}} is 
idempotent and only fires when the signal has already arrived, so this adds no 
new cross-component traffic.

This was originally suspected of causing {{DynamicKafkaSourceITTest}} 
{{testIdleReader}} hangs in CI, but that has not been reproduced locally (11/11 
passes on unfixed main), so the ITCase link is unconfirmed.

Related: FLINK-40094 and FLINK-39979 touch the same reader idleness area.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to