Efrat19 commented on PR #290: URL: https://github.com/apache/flink-connector-kafka/pull/290#issuecomment-5329176562
> Thanks for tracking down [FLINK-36434](https://issues.apache.org/jira/browse/FLINK-36434) — lazy creation on the fetcher thread is the right fix, and confining it to `KafkaPartitionSplitReader` is the right call for this ticket (I looked at pushing the fix into `SplitFetcherManager#createSplitFetcher`/`SplitFetcher` instead, but `getSplitReader()` is public and several callers, including `KafkaSourceFetcherManager#commitOffsets`, capture the reader reference before the fetcher thread ever runs — that's a much bigger, cross-connector change and belongs in its own ticket as you noted). > > One simplification on the locking, though: I don't think `pendingWakeup`/`consumerLock` are reachable in practice, so we can drop them (and the lock) entirely. > > `SplitReader#wakeUp()`'s contract is narrow: "wake up the split reader in case the fetcher thread is blocking in `fetch()`". In the whole of [flink-connector-base](https://issues.apache.org/jira/browse/FLINK-connector-base) + [flink-connector-kafka](https://issues.apache.org/jira/browse/FLINK-connector-kafka), `splitReader.wakeUp()` has exactly one call site: `FetchTask#wakeUp()`, which only fires when `SplitFetcher`'s `runningTask == fetchTask`. `fetchTask` is only ever selected once `assignedSplits` is non-empty, and `assignedSplits` is populated in exactly one place — `AddSplitsTask#run()` — which unconditionally calls `handleSplitsChanges()` (i.e. `ensureConsumer()`) immediately after, on the same fetcher thread, before that thread ever loops back to pick a next task. So by construction, `wakeUp()` can never reach `KafkaPartitionSplitReader` before `ensureConsumer()` has already run once. > > That means a plain `volatile` field is enough — no `AtomicReference`, no `synchronized`, no pending-wakeup bookkeeping: > > ```java > /** Created lazily on first use by the fetcher thread — see #ensureConsumer(). */ > private volatile KafkaConsumer<byte[], byte[]> consumer; > > private KafkaConsumer<byte[], byte[]> ensureConsumer() { > KafkaConsumer<byte[], byte[]> c = consumer; > if (c == null) { > c = createConsumer(consumerProps); > maybeRegisterKafkaConsumerMetrics(consumerProps, kafkaSourceReaderMetrics, c); > kafkaSourceReaderMetrics.registerNumBytesIn(c); > consumer = c; > } > return c; > } > > @Override > public void wakeUp() { > // wakeUp() only ever reaches a running fetch task (see SplitReader#wakeUp javadoc), > // which the base fetcher can only schedule after handleSplitsChanges() has run at > // least once, so consumer is always non-null here in practice. > KafkaConsumer<byte[], byte[]> c = consumer; > if (c != null) { > c.wakeup(); > } > } > > @Override > public void close() throws Exception { > KafkaConsumer<byte[], byte[]> c = consumer; > if (c != null) { > c.close(); > } > } > ``` > > `ensureConsumer()`'s creation branch needs no guard because it's single-writer (fetcher thread only calls it); `wakeUp()`'s cross-thread read just needs the `volatile` for visibility. Worth noting this leans on an implementation invariant of `SplitFetcher`/`AddSplitsTask` rather than a documented contract beyond `wakeUp()`'s own javadoc scoping — stable code, but that's why I'd keep the comment explaining it. > > This also removes the exact thing @Efrat19 flagged as feeling weird (the `@GuardedBy` split between `ensureConsumer`/`createConsumer`). > > The `KafkaPartitionSplitReaderLazyConsumerCreationTest` case for "wakeUp before creation is applied on creation" would need to flip to asserting that a pre-creation `wakeUp()` is safely dropped instead of deferred. cc @dalelane since https://github.com/apache/flink-connector-kafka/pull/293 introduces another consumer access from the main thread -- 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]
