PySpark连接Elasticsearch报错:无法检测ES版本求助
连接PySpark到Elasticsearch并插入DataFrame时触发以下错误:
Py4JJavaError: An error occurred while calling o130.save.
: org.elasticsearch.hadoop.EsHadoopIllegalArgumentException: Cannot detect ES version - typically this happens if the
network/Elasticsearch cluster is not accessible or when targeting a
WAN/Cloud instance without the proper setting 'es.nodes.wan.only'
相关代码片段:
from pyspark.sql import SparkSession from pyspark.sql.functions import * appName = "mysql example" master = "local" spark = SparkSession.builder.master(master).appName(appName)\ .config("spark.jars", "postgresql-42.2.6.jar,elasticsearch-spark-30_2.12-8.8.1.jar").getOrCreate() category_product_df = dataframe_from_table("category") category_product_df.write.format("org.elasticsearch.spark.sql") \ .option("es.resource", "wael/test") \ .option("es.port", "9200") \ .option("es.nodes", "elastic:changeme@localhost") \ .option("es.nodes.wan.only", "true") \ .save()
环境详情:
- PySpark版本:3.4.0
- Scala版本:2.12.17
- Elasticsearch版本:7.15.0
- Elasticsearch-Spark连接器版本:elasticsearch-spark-30_2.12-8.8.1.jar
已确认网络连通正常,主机、端口和认证凭据无误,且已设置es.nodes.wan.only=true,但错误仍存在。
修复版本兼容性问题
Elasticsearch和其Spark连接器跨大版本会存在兼容性冲突,你当前用的ES是7.15.0,但连接器是8.8.1,大版本不匹配直接导致版本检测失败。建议更换为与ES版本一致的连接器,比如elasticsearch-spark-30_2.12-7.15.0.jar。修正认证参数格式
不要在es.nodes中直接拼接用户名和密码,这种格式无法被连接器正确解析,分开配置认证信息:.option("es.nodes", "localhost") \ .option("es.net.http.auth.user", "elastic") \ .option("es.net.http.auth.pass", "changeme")确认协议与端口配置
如果ES启用了HTTPS,需添加es.net.ssl=true,同时确认端口是否为9243(HTTPS默认);如果是HTTP,可尝试在es.nodes中带上完整协议前缀:.option("es.nodes", "http://localhost:9200")调整索引资源格式
ES 7.x及以上版本已废弃文档类型,es.resource无需指定类型,直接写索引名即可,比如将wael/test改为wael。强制指定ES版本
若上述操作无效,可跳过自动版本检测,强制指定ES版本:.option("es.version", "7.15.0")
内容的提问来源于stack exchange,提问作者Flàyn

