【问题标题】:Has anyone been able to use elasticsearch xpack sql with Spark?有没有人能够将 elasticsearch xpack sql 与 Spark 一起使用?
【发布时间】:2019-01-31 00:04:05
【问题描述】:

我正在尝试使用 PySpark 从 elasticsearch 读取数据。通常我会将查询设置为沿线的内容(请参阅下面的查询)并将 es.resource 设置为索引,例如“my_index/doc”,我可以将数据读入 spark:

q ="""{
          "query": {
              "match_all": {}
          }  
      }"""

但是最近我尝试了 _xpack/sql 与 kibana 和 JDBC 与其他 SQL 客户端,它们在获取数据方面工作得很好。但是,当我尝试在我的 pyspark 代码中引用 _xpack 时,出现以下错误:

Py4JJavaError: An error occurred while calling 
z:org.apache.spark.api.python.PythonRDD.newAPIHadoopRDD.
: org.elasticsearch.hadoop.rest.EsHadoopInvalidRequest: 
org.elasticsearch.hadoop.rest.EsHadoopRemoteException: 
invalid_index_name_exception: Invalid index name [_xpack], must not start with '_'.
null

有没有人尝试过使用 _xpack 或者知道如何从 Elasticsearch hadoop 插件执行 Elasticsearch SQL 查询?

您将在下面找到我试图用来在 pyspark 上执行的代码摘录,提前致谢!

q = """{"query": "select * from eg_flight limit 1"}"""

es_read_conf = {
    "es.nodes" : "192.168.1.71,192.168.1.72,192.168.1.73",
    "es.port" : "9200",
    "es.resource" :  "_xpack/sql",
    "es.query" : q
}

es_rdd = sc.newAPIHadoopRDD(
    inputFormatClass="org.elasticsearch.hadoop.mr.EsInputFormat",
    keyClass="org.apache.hadoop.io.NullWritable", 
    valueClass="org.elasticsearch.hadoop.mr.LinkedMapWritable", 
    conf=es_read_conf)

【问题讨论】:

    标签: apache-spark elasticsearch pyspark apache-spark-sql


    【解决方案1】:

    我认为不支持此功能。 PySpark 中的另一种解决方案是使用 JDBC 驱动程序,我确实尝试过。我尝试了以下方法:

    es_df = spark.read.jdbc(url="jdbc:es://http://192.168.1.71:9200", table = "(select * from eg_flight) mytable")
    

    我收到以下错误:

    Py4JJavaError: An error occurred while calling o2488.jdbc.
    : java.sql.SQLFeatureNotSupportedException: Found 1 problem(s)
    line 1:8: Unexecutable item
    
    ...
    

    另一种方法是使用核心 Python 和请求来执行此操作,但我不建议将它用于大型数据集。

    import requests as r
    import json
    
    
    es_template = {
        "query": "select * from eg_flight"
    }
    
    es_link = "http://192.168.1.71:9200/_xpack/sql"
    headers = {'Content-type': 'application/json'}
    
    
    if __name__ == "__main__":
    
        load = r.post(es_link, data=json.dumps(es_template), headers=headers)
        if load.status_code == 200:
            load = load.json()
            #do something with it
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2012-01-02
      • 1970-01-01
      • 2018-08-14
      • 2010-10-29
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多