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
相关产品推荐
相关产品推荐

