You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

配置启用X-Pack的Elasticsearch Structured Streaming及认证设置

问题描述

我正尝试使用已安装X-Pack的Elasticsearch (ES) 6.1.1,通过Spark Structured Streaming 2.2.1将数据写入ES(ES中已存在对应索引)。我的代码如下:

val exceptions = spark 
 .readStream 
 .text(path) 
val advancedQuery = exceptions 
 .writeStream 
 .format("org.elasticsearch.spark.sql") 
 .trigger(Trigger.ProcessingTime(10.seconds)) 
 .option("checkpointLocation", "/checkpoint") 
val runningQuery = advancedQuery.start("spark/exc") 
runningQuery.awaitTermination

但运行时抛出异常,核心错误信息如下:

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'
...
Caused by: org.elasticsearch.hadoop.rest.EsHadoopInvalidRequest: missing authentication token for REST request [/] null

请问该如何设置所需的认证信息?

解决方案

从错误栈能看出来,问题出在两个核心点:一是访问带X-Pack的ES集群时缺少认证令牌,二是因为认证失败导致Spark无法连接集群、检测ES版本。你需要在Spark Streaming的写入配置中补充以下必要配置:

  • 添加X-Pack认证信息
    在writeStream的配置链中加入用户名和密码的配置项,替换为你的ES集群实际认证信息(默认X-Pack的管理员用户名是elastic):

    .option("es.net.http.auth.user", "your-es-username")
    .option("es.net.http.auth.pass", "your-es-password")
    
  • 指定ES集群节点地址
    错误提示无法检测ES版本,本质是因为Spark没能成功连接到ES集群,所以需要显式指定ES节点的IP和端口:

    .option("es.nodes", "your-es-node-ip:9200")
    

    如果是多节点集群,可以用逗号分隔多个节点,比如"node1.example.com:9200,node2.example.com:9200"。

  • (可选)WAN/云环境ES实例配置
    如果你的ES集群部署在云环境或者公网WAN中,需要额外添加es.nodes.wan.only配置,确保Spark以WAN模式连接:

    .option("es.nodes.wan.only", "true")
    

修改后的完整代码示例

val exceptions = spark 
 .readStream 
 .text(path) 

val advancedQuery = exceptions 
 .writeStream 
 .format("org.elasticsearch.spark.sql") 
 .trigger(Trigger.ProcessingTime(10.seconds)) 
 .option("checkpointLocation", "/checkpoint")
 // 配置X-Pack认证
 .option("es.net.http.auth.user", "elastic")
 .option("es.net.http.auth.pass", "your-secure-password")
 // 指定ES集群节点
 .option("es.nodes", "127.0.0.1:9200")
 // 若为云/WAN环境则启用下面的配置
 // .option("es.nodes.wan.only", "true")

val runningQuery = advancedQuery.start("spark/exc") 
runningQuery.awaitTermination()
  • 版本兼容性提醒
    一定要确保你的elasticsearch-spark依赖版本和ES集群版本严格一致(都是6.1.1),Spark 2.2.1和ES 6.1.1的组合是兼容的。如果用Maven管理依赖,参考以下配置:
    <dependency>
        <groupId>org.elasticsearch</groupId>
        <artifactId>elasticsearch-spark-20_2.11</artifactId>
        <version>6.1.1</version>
    </dependency>
    

内容的提问来源于stack exchange,提问作者Matthias Mueller

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 08:30:45