Azure Databricks连接Confluent Cloud超时异常排查与解决
Azure Databricks写入Confluent Cloud Kafka超时问题排查与解决
问题现象
在Azure Databricks Notebook中向Confluent Cloud的Kafka集群写入数据时,出现超时异常,目标Topic确认已存在且网络无问题。异常信息如下:Job aborted due to stage failure: Topic spark_poc_topic not present in metadata after 60000 ms. Caused by: TimeoutException: Topic spark_poc_topic not present in metadata after 60000 ms
问题代码
相关Scala写入代码:
df2.selectExpr("key", "value") .write .format("kafka") .option("kafka.bootstrap.servers", "xxxx.azure.confluent.cloud:9092") .option("kafka.security.protocol", "SASL_PLAINTEXT") .option("kafka.sasl.jaas.config", "kafkashaded.org.apache.kafka.common.security.plain.PlainLoginModule required username='{}' password='{}';".format("INxxxxxxxI", "n/LeO+aEJbxxxxx")) .option("kafka.ssl.endpoint.identification.algorithm", "https") .option("kafka.sasl.mechanism", "PLAIN") .option("topic", "topic123") .save()
问题根因
问题出在**kafka.sasl.jaas.config的配置写法**:使用单引号包裹用户名和密码,结合format方法拼接的方式会导致JAAS配置解析异常,无法完成身份验证,最终因无法获取Topic元数据引发超时。
解决方案
修改kafka.sasl.jaas.config的配置方式,改用双引号包裹用户名和密码,通过字符串拼接传入Confluent Cloud的API密钥与密码:
.option("kafka.sasl.jaas.config", "kafkashaded.org.apache.kafka.common.security.plain.PlainLoginModule required username=\"" + confluentApiKey + "\" password=\"" + confluentSecret + "\";")
(注:confluentApiKey和confluentSecret需提前定义为对应Confluent Cloud的API密钥、密码变量)
内容的提问来源于stack exchange,提问作者191180rk
相关产品推荐
相关产品推荐

