Pyspark写入Elasticsearch遇EsHadoopIllegalArgumentException问题求助
在使用PySpark将DataFrame写入Elasticsearch集群时遇到如下错误:
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'
我的代码实现如下:
df1.write.format("org.elasticsearch.spark.sql")\ .option("es.nodes", host)\ .option("es.port", port)\ .option("es.net.http.auth.user", username)\ .option("es.net.http.auth.pass", password)\ .option("es.resource", indexName)\ .option("es.net.ssl.keystore.location", pathToCAFile))\ .mode('overwrite')\ .save()
已尝试添加es.nodes.wan.only配置但无效;通过curl和Python均可正常连接Elasticsearch集群,且使用的EsHadoop-jar版本与集群版本一致,仍无法解决问题。
补全SSL配置参数
仅配置es.net.ssl.keystore.location不足以建立完整的SSL连接,需补充以下参数:.option("es.net.ssl", "true")\ .option("es.net.ssl.truststore.location", pathToCAFile)\ .option("es.net.ssl.truststore.type", "PEM") # 若CA证书为PEM格式,必须指定自签名证书场景下,必须明确开启SSL并指定信任存储,否则连接器会拒绝握手。
修正
es.resource格式es.resource需符合索引名/文档类型格式(ES 7.x+可省略文档类型,直接写索引名),确认你的indexName格式正确,示例:# ES 6.x .option("es.resource", "my_target_index/_doc") # ES 7+ .option("es.resource", "my_target_index")强制指定ES版本跳过检测
若已知集群版本,可直接配置es.version跳过自动版本检测,示例:.option("es.version", "7.17.0") # 替换为你的ES集群版本检查节点地址格式
es.nodes参数不要带http://或https://前缀,仅保留域名或IP,连接器会根据SSL配置自动适配协议。验证Spark集群节点的网络连通性
确保Spark集群所有节点都能访问ES集群的目标端口(如9200/443),可在Spark节点上执行curl命令测试,排除节点间网络隔离问题。确认权限配置
检查ES账号是否拥有目标索引的写入权限,同时确认ES的ACL规则允许Spark集群IP段访问。
内容的提问来源于stack exchange,提问作者Piyush Jain

