Spark-es连接器向Azure Elasticsearch写数据时遇连接异常求助
解决Spark-ES写入Azure Elasticsearch时的EsHadoopNoNodesLeftException问题
针对你写入2000万条数据到Azure Elasticsearch时,在完成1300万后出现的连接失败问题,提供以下解决思路:
1. 排查Elasticsearch集群状态与限流
- 登录Azure门户查看ES集群监控指标:重点关注CPU使用率、内存占用、磁盘IO、请求吞吐量和错误率,确认是否因写入压力过大触发集群限流或过载保护。
- 调整ES批量写入参数,降低单次批量的负载:通过
es.batch.size.bytes(如设为10mb)和es.batch.size.entries(如设为1000)控制每次发送到ES的数据量,避免压垮集群。
2. 增强写入的容错与重试机制
- 添加ES HTTP请求重试配置:设置
es.http.retries为5(重试次数)、es.http.timeout为30s(延长超时时间)、es.http.retries.wait为5s(重试间隔),应对临时的集群繁忙或网络波动。 - 配置Spark任务重试:设置
spark.task.maxFailures为3或更高,让单个写入失败的分区任务自动重试,而非直接终止整个作业。
3. 优化网络与SSL连接配置
- 确认网络连通性:检查Spark集群与Azure ES之间的防火墙/NSG规则是否临时阻断连接,公网访问场景下排查是否存在网络带宽瓶颈或运营商临时故障。
- 验证SSL证书信任:如果Spark节点无法信任Azure ES的SSL证书,可尝试设置
es.net.ssl.cert.allow.self.signed为true(仅测试用,生产环境建议配置信任store),或指定es.net.ssl.truststore.location指向合法的信任证书路径。
4. 调整索引分片与Spark并行度
- 优化ES索引分片数:2000万数据建议设置10-20个主分片,避免热点分片导致集群负载不均,可通过修改索引模板或重建索引调整。
- 匹配Spark并行度:设置
spark.sql.shuffle.partitions的值与ES分片数相近,让写入请求均匀分布到各个分片,避免并行度过高压垮集群或过低导致单任务负载过大。
5. 实现断点续写避免重复写入
- 由于已成功写入1300万数据,可通过
id字段(你已配置es.mapping.id)过滤已写入的数据:先从ES中导出已存在的id集合,在Spark数据集中过滤掉这些id后再执行写入,避免重复数据和重复计算。
优化后的写入代码示例
data .write .format("org.elasticsearch.spark.sql") .option("es.nodes", node) .option("es.port", port) .option("es.net.http.auth.user", username) .option("es.net.http.auth.pass", password) .option("es.net.ssl", "true") .option("es.nodes.wan.only", "true") .option("es.mapping.id", "id") // 容错重试配置 .option("es.http.retries", "5") .option("es.http.timeout", "30s") .option("es.http.retries.wait", "5s") // 批量写入控制 .option("es.batch.size.bytes", "10mb") .option("es.batch.size.entries", "1000") .mode(writingMode) .save(index)
内容的提问来源于stack exchange,提问作者mham
相关产品推荐
相关产品推荐

