【问题标题】:create new SparkSession for different query?为不同的查询创建新的 SparkSession?
【发布时间】:2019-12-26 06:29:07
【问题描述】:

我想从 elasticsearch 获取两个数据

一个用查询过滤,另一个没有过滤器。

 // with query
 session = get_spark_session(query=query)

 df = session.read.option(
     "es.resource", "analytics-prod-2019.08.02"
 ).format("org.elasticsearch.spark.sql").load()


 df.show() // empty result

 // without query
 session = get_spark_session()

 df = session.read.option(
     "es.resource", "analytics-prod-2019.08.02"
 ).format("org.elasticsearch.spark.sql").load()

 df.show() // empty result


 def get_spark_session(query=None, excludes=[]):

     conf = pyspark.SparkConf()
     conf.set("spark.driver.allowMultipleContexts", "true")
     conf.set("es.index.auto.create", "true")
     conf.set("es.nodes.discovery", "true")
     conf.set("es.scroll.size", 10000)
     conf.set("es.read.field.exclude", excludes)
     conf.set("spark.driver.extraClassPath", "/usr/local/elasticsearch-hadoop/dist/elasticsearch-spark-20_2.11-6.6.2.jar")
     if query:
         conf.set("es.query", query)


     sc = SparkSession.builder.config(conf=conf).getOrCreate()

     return sc 

问题是会话是否被重用..

当我首先运行 filtered 查询,然后运行 ​​non-filtered 查询时, 两者都给出空结果

但是当我首先运行non-filtered查询时,它会显示一些结果,而随后的filtered查询显示空结果。

 // below, I reverse the order
 // without query
 session = get_spark_session()

 df = session.read.option(
     "es.resource", "analytics-prod-2019.08.02"
 ).format("org.elasticsearch.spark.sql").load()

 df.show() // some result

 // with query
 session = get_spark_session(query=query)

 df = session.read.option(
     "es.resource", "analytics-prod-2019.08.02"
 ).format("org.elasticsearch.spark.sql").load()


 df.show() // empty result

** 编辑

所以我可以通过以下方式获得所需的结果:

def get_spark_session(query=None, excludes=[]):

    conf = pyspark.SparkConf()
    conf.set("spark.driver.allowMultipleContexts", "true")
    conf.set("es.index.auto.create", "true")
    conf.set("es.nodes.discovery", "true")
    conf.set("es.scroll.size", 10000)
    conf.set("es.read.field.exclude", excludes)
    conf.set("spark.driver.extraClassPath", "/usr/local/elasticsearch-hadoop/dist/elasticsearch-spark-20_2.11-6.6.2.jar")
    if query:
        conf.set("es.query", query)
    else:
        conf.set("es.query", "") # unset the query 

【问题讨论】:

    标签: python apache-spark elasticsearch pyspark elasticsearch-hadoop


    【解决方案1】:

    SparkSession.builder 获取一个现有的 SparkSession,或者,如果没有现有的,则根据此构建器中设置的选项创建一个新的。在您的情况下,火花配置正在被重用。从配置中删除“es.query”应该可以解决这个问题:

    def get_spark_session(query=None, excludes=[]):
         conf = pyspark.SparkConf()
         conf.unset("es.query")
         conf.set("spark.driver.allowMultipleContexts", "true")
         conf.set("es.index.auto.create", "true")
         conf.set("es.nodes.discovery", "true")
         conf.set("es.scroll.size", 10000)
         conf.set("es.read.field.exclude", excludes)
         conf.set("spark.driver.extraClassPath", "/usr/local/elasticsearch-hadoop/dist/elasticsearch-spark-20_2.11-6.6.2.jar")
         if query:
             conf.set("es.query", query)    
    
         sc = SparkSession.builder.config(conf=conf).getOrCreate()    
         return sc 
    

    【讨论】:

    • 那么我如何获得一个数据集with查询和without查询?
    • 当你首先运行过滤后的查询时,查询是在 conf.xml 中设置的。由于查询是 conf 的一部分,当您查询未过滤的数据集时,它会被重用。
    猜你喜欢
    • 2022-01-23
    • 1970-01-01
    • 1970-01-01
    • 2020-10-13
    • 2021-09-21
    • 1970-01-01
    • 2017-09-26
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多