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

Databricks中Spark Streaming读Kafka遇权限集群不支持数据源V2错误该如何解决?

Databricks读取Kafka流报错问题排查

问题场景

在Databricks环境中运行以下PySpark代码读取Kafka流:

kafka = spark.readStream\
    .format("kafka")\
    .option("kafka.sasl.mechanism", "SCRAM-SHA-512")\
    .option("kafka.security.protocol", "SASL_SSL")\
    .option("kafka.sasl.jaas.config", f'org.apache.kafka.common.security.scram.ScramLoginModule required  username="{user_stg}" password="{pass_stg}"')\
    .option("kafka.bootstrap.servers", "b-1.dataservices-msk-st.****.amazonaws.com:9096")\
    .option("subscribe", "app-***-events")\
    .option("startingOffsets", "earliest").load()

执行后返回错误:

Java.lang.SecurityException: Data source V2 streaming is not supported on table acl or credential passthrough clusters. 
StreamingRelationV2 org.apache.spark.sql.kafka010.KafkaSourceProvider@11002bae, kafka, org.apache.spark.sql.kafka010.KafkaSourceProvider$KafkaTable@35ae434, 
[kafka.sasl.mechanism=SCRAM-SHA-512, subscribe=app--events, kafka.sasl.jaas.config=*********(redacted), kafka.bootstrap.servers=b-1.dataservices-msk-st.****.amazonaws.com:9096, startingOffs

错误原因

该错误的核心原因是当前Databricks集群启用了表ACL或**凭证传递(Credential Passthrough)**安全配置,而Spark Kafka数据源的V2流式读取机制不兼容这两种集群模式。Databricks在这类安全增强型集群上限制了V2流式数据源的使用,以避免权限管理冲突和潜在的安全风险。

解决办法

  • 禁用集群的表ACL或凭证传递功能:若业务场景允许,修改集群配置,关闭表ACL权限控制(在集群权限设置中取消启用)或凭证传递功能(在集群编辑页面的「高级选项」->「Spark」中关闭相关开关),重启集群后重新执行代码。
  • 强制使用Kafka数据源V1:在读取流的配置中添加.option("useDeprecatedOffsetFetching", "true")参数,强制Spark采用旧版V1数据源,绕过V2的兼容性限制。修改后的代码示例:
kafka = spark.readStream\
    .format("kafka")\
    .option("kafka.sasl.mechanism", "SCRAM-SHA-512")\
    .option("kafka.security.protocol", "SASL_SSL")\
    .option("kafka.sasl.jaas.config", f'org.apache.kafka.common.security.scram.ScramLoginModule required  username="{user_stg}" password="{pass_stg}"')\
    .option("kafka.bootstrap.servers", "b-1.dataservices-msk-st.****.amazonaws.com:9096")\
    .option("subscribe", "app-***-events")\
    .option("startingOffsets", "earliest")\
    .option("useDeprecatedOffsetFetching", "true")\  # 新增参数
    .load()
  • 使用Databricks官方Kafka连接器:改用Databricks适配的Kafka连接器,该连接器针对表ACL和凭证传递集群做了专门优化,支持流式读取操作。可通过Delta Live Table或结构化流结合连接器配置实现。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 02:15:42