You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Azure Databricks连接Confluent Kafka超时:节点分配失败求排查

排查Azure Databricks与Confluent Kafka集成时的Stream读取超时及认证问题

问题现象

  • 使用SASL_PLAINTEXT创建readStream时触发超时错误:
    java.util.concurrent.ExecutionException: kafkashaded.org.apache.kafka.common.errors.TimeoutException: Timed out waiting for a node assignment. Call: describeTopics
    
  • 尝试切换SSL连接并添加.option("kafka.ssl.endpoint.identification.algorithm", "https")后,出现认证错误

排查步骤

1. 验证网络连通性

  • 即使Azure DBX与Confluent Kafka同属一个资源组,仍需确认Confluent的bootstrap server端口(SASL_PLAINTEXT默认9092、SASL_SSL默认9093)是否在Databricks集群的安全组/网络规则中开放,允许集群IP范围访问。
  • 在Databricks集群中执行命令测试连通性:
    nc -zv <bootstrap-server-host> <port>
    
  • 若使用Confluent Cloud,需检查IP白名单配置,确认Databricks集群的公网IP已加入白名单(公网连接场景)。

2. 核对SASL认证配置细节

  • 确认kafka.sasl.jaas.config中的用户名是Confluent的API Key,密码是API Secret,不要混淆两者。
  • 检查kafka.security.protocol与Confluent Kafka的监听配置匹配:若Confluent端仅启用SASL_SSL,使用SASL_PLAINTEXT会导致无法建立连接,触发超时。
  • 调整kafka.request.timeout.ms至更大值(例如30000),避免因网络延迟导致超时。

3. 修正SSL连接配置

  • 使用SSL时,kafka.security.protocol需设置为SASL_SSL,而非单纯的SSL。
  • kafka.ssl.endpoint.identification.algorithm值需与Confluent证书的CN域名匹配,通常设置为HTTPS(大小写不敏感)。
  • 若Confluent使用自签名证书,需将证书上传至Databricks集群,并添加以下配置:
    .option('kafka.ssl.truststore.location', '/path/to/truststore.jks')
    .option('kafka.ssl.truststore.password', 'truststore-password')
    
  • 标准SASL_SSL配置示例:
    df = spark  \ 
      .readStream \ 
      .format('kafka') \ 
      .option('kafka.bootstrap.servers', '...:...') \ 
      .option('subscribe', 'transaction_stream') \ 
      .option('group.id', 'dbx-streaming') \ 
      .option('startingOffsets', 'earliest') \ 
      .option('kafka.security.protocol', 'SASL_SSL') \ 
      .option('kafka.sasl.mechanism', 'PLAIN') \ 
      .option('kafka.sasl.jaas.config', \ 
      """kafkashaded.org.apache.kafka.common.security.plain.PlainLoginModule required username="<API-KEY>" password="<API-SECRET>";""") \ 
      .option('kafka.request.timeout.ms', '30000') \ 
      .option("kafka.enable.idempotence", "false") \ 
      .load()
    

4. 检查版本兼容性

  • 确认Databricks运行时自带的Kafka客户端版本与Confluent Kafka版本兼容(例如Confluent 7.4对应Kafka 3.4),版本差异过大可能导致协议不兼容,引发超时或认证失败。
  • 在Databricks Notebook中执行spark.version查看运行时版本,再核对对应Kafka依赖版本。

5. 查看Confluent日志

  • 登录Confluent控制平台,查看Broker日志,确认是否有连接请求记录,以及认证失败的具体原因(如无效API Key、IP被拒绝等),直接定位问题根源。

内容的提问来源于stack exchange,提问作者SimonNagy

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.14 15:05:06