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

