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

调用t.addCustomDisplayData触发Py4JJavaError,Spark读取Kafka求助

Spark读取Kafka触发Py4JJavaError错误的排查

执行Spark读取Kafka的代码时出现Py4JJavaError错误,错误提示为An error occurred while calling t.addCustomDisplayData,涉事代码如下:

df = ( spark.read .format('kafka') 
.option('kafka.bootstrap.servers', BOOTSTRAP_SERVER) 
.option('kafka.security.protocol','SASL_SSL') 
.option('kafka.sasl.mechanism','PLAIN') 
.option('kafka.sasl.jaas.config', f"{JAAS_MODULE} required username= '{CLUSTER_API_KEY}' password='{CLUSTER_API_SECRET}';") 
.option('subscribe','invoices') .load() ) 

问题根源:JAAS配置格式错误

  • 代码末尾的JAAS配置字符串多了一个转义双引号",导致配置解析失败,进而触发后续的addCustomDisplayData错误
  • username和password值两侧的单引号属于冗余内容,JAAS配置规范里不需要给账号密码额外加单引号

修正后的代码:

df = (spark.read
      .format('kafka')
      .option('kafka.bootstrap.servers', BOOTSTRAP_SERVER)
      .option('kafka.security.protocol', 'SASL_SSL')
      .option('kafka.sasl.mechanism', 'PLAIN')
      .option('kafka.sasl.jaas.config', f"{JAAS_MODULE} required username={CLUSTER_API_KEY} password={CLUSTER_API_SECRET};")
      .option('subscribe', 'invoices')
      .load())

额外注意事项

  • 确认JAAS_MODULE变量值正确,比如Confluent Kafka场景下通常为org.apache.kafka.common.security.plain.PlainLoginModule
  • 若CLUSTER_API_KEY或CLUSTER_API_SECRET包含特殊字符,需做转义处理

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 08:49:54