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]

Reply via email to