c-broggy commented on issue #39789:
URL: https://github.com/apache/beam/issues/39789#issuecomment-5867728308
@stankiewicz this is what I have found :
Enabled --
sdk_harness_log_level_overrides='{"org.apache.kafka.common.network":"DEBUG","org.apache.kafka.common.security":"DEBUG"}'
on the default (unpinned) beam_java25_sdk harness, and the debug logging
surfaces the actual cause behind the TimeoutException: Timeout expired while
fetching topic metadata:
java.lang.UnsupportedOperationException: getSubject is not supported
at java.base/javax.security.auth.Subject.getSubject(Subject.java:277)
at
org.apache.kafka.common.security.oauthbearer.internals.OAuthBearerSaslClientCallbackHandler.handleCallback(OAuthBearerSaslClientCallbackHandler.java:100)
at
org.apache.kafka.common.security.oauthbearer.internals.OAuthBearerSaslClient.evaluateChallenge(OAuthBearerSaslClient.java:93)
at
org.apache.kafka.common.security.authenticator.SaslClientAuthenticator.createSaslToken(SaslClientAuthenticator.java:536)
at
org.apache.kafka.common.security.authenticator.SaslClientAuthenticator.sendInitialToken(SaslClientAuthenticator.java:335)
at
org.apache.kafka.common.security.authenticator.SaslClientAuthenticator.authenticate(SaslClientAuthenticator.java:276)
at
org.apache.kafka.common.network.KafkaChannel.prepare(KafkaChannel.java:181)
at
org.apache.kafka.common.network.Selector.pollSelectionKeys(Selector.java:563)
...
at
org.apache.kafka.clients.consumer.internals.TopicMetadataFetcher.getTopicMetadata(TopicMetadataFetcher.java:102)
at
org.apache.kafka.clients.consumer.KafkaConsumer.partitionsFor(KafkaConsumer.java:1434)
at
org.apache.beam.sdk.io.kafka.KafkaIO$Read$GenerateKafkaSourceDescriptor.processElement(KafkaIO.java:2169)
The OAUTHBEARER token fetch itself succeeds ("Successfully logged in."), but
the SASL handshake then fails immediately with the exception above, the
connection is torn down, and it retries in a loop — which eventually surfaces
as the metadata-fetch timeout everyone's seeing.
This looks like a JDK-version compatibility issue, not anything cluster- or
config-specific. JDK 24+ finalized JEP 486 (permanently disabling the Security
Manager), which turns Subject.getSubject(AccessControlContext) into a hard
UnsupportedOperationException. Kafka's OAuthBearerSaslClientCallbackHandler
still calls that legacy API to retrieve the current Subject during the SASL
handshake — so any environment where the SDK harness runs on JDK 24+ (e.g.
beam_java25_sdk) will hit this on every OAUTHBEARER auth attempt against Kafka.
Consistent with that: pinning sdk_harness_container_image_overrides to a
harness on JDK 21 or earlier (predates JEP 486) avoids the issue entirely —
clean connect and auth, no disconnects, no timeouts. So the fix/workaround is
specifically about the JDK version the harness runs on, not the Beam version or
SDK harness image name per se — any harness bundling JDK 24+ should reproduce
this with SASL/OAUTHBEARER Kafka auth, and any harness on JDK ≤21 should not.
--
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]