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

Databricks写入Kafka遭遇集群授权失败问题求助

问题描述

我正在使用Databricks将数据发送至Kafka,代码如下:

(df.write
  .format('kafka')
  .option('kafka.bootstrap.servers', 'glider.srvs.cloudkafka.com:9094')
  .option('kafka.security.protocol', 'SASL_SSL')
  .option('kafka.sasl.mechanism', 'SCRAM-SHA-512')
  .option("kafka.sasl.jaas.config", "kafkashaded.org.apache.kafka.common.security.scram.ScramLoginModule required username='my_username' password='my_password';")
  .option('topic', '<my_topic>')
  .save())

但执行失败,报错信息如下:

An error occurred while calling o2296.save.
: org.apache.spark.SparkException: Job aborted due to stage failure: Task 3 in stage 31.0 failed 1 times, most recent failure: Lost task 3.0 in stage 31.0 (TID 251) (ip-10-172-216-86.us-west-2.compute.internal executor driver): kafkashaded.org.apache.kafka.common.KafkaException: Cannot execute transactional method because we are in an error state
...
Caused by: [CIRCULAR REFERENCE: kafkashaded.org.apache.kafka.common.errors.ClusterAuthorizationException: Cluster authorization failed.]

已确认凭证与安全协议完全正确,因为使用Python的confluent_kafka库可正常发送数据。

我尝试更换登录模块:

.option("kafka.sasl.jaas.config", "kafkashaded.org.apache.kafka.common.security.plain.PlainLoginModule required username='my_username' password='my_password';")

但结果依旧报错。请问问题出在哪里?是否需要调整某些配置项?


问题分析与解决方案

核心问题是Spark Kafka Connector默认启用事务性写入,而你的CloudKafka集群可能未开启事务支持,或当前用户无事务相关权限,导致授权失败。以下是具体调整方案:

  • 禁用事务性写入:
    在配置中添加事务ID为空值,或直接关闭幂等性来禁用事务逻辑:

    (df.write
      .format('kafka')
      .option('kafka.bootstrap.servers', 'glider.srvs.cloudkafka.com:9094')
      .option('kafka.security.protocol', 'SASL_SSL')
      .option('kafka.sasl.mechanism', 'SCRAM-SHA-512')
      .option("kafka.sasl.jaas.config", "kafkashaded.org.apache.kafka.common.security.scram.ScramLoginModule required username='my_username' password='my_password';")
      .option('topic', '<my_topic>')
      .option('kafka.transactional.id', '')  # 禁用事务
      .save())
    

    或直接设置:

    .option('kafka.enable.idempotence', 'false')
    
  • 简化JAAS配置类路径:
    部分Spark Kafka Connector版本无需kafkashaded.前缀,尝试简化登录模块配置:

    .option("kafka.sasl.jaas.config", "org.apache.kafka.common.security.scram.ScramLoginModule required username='my_username' password='my_password';")
    
  • 确认主题与集群权限:
    确保当前用户对目标主题拥有WRITE权限,同时检查是否具备集群级DESCRIBE权限(Spark Connector初始化时需获取集群元数据)。

  • 指定写入模式:
    若无需Exactly-Once语义,显式设置追加模式,避免触发事务逻辑:

    (df.write
      .format('kafka')
      .mode("append")
      # 其他配置项...
      .save())
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 08:12:53