【问题标题】:Running custom Apache Phoenix SQL query in PySpark在 PySpark 中运行自定义 Apache Phoenix SQL 查询
【发布时间】:2019-07-24 06:04:24
【问题描述】:

有人可以提供一个使用 pyspark 的示例,说明如何运行自定义 Apache Phoenix SQL 查询并将该查询的结果存储在 RDD 或 DF 中。注意:我正在寻找一个自定义查询,而不是要读入 RDD 的整个表。

从 Phoenix 文档中,我可以使用这个来加载整个表格:

table = sqlContext.read \
        .format("org.apache.phoenix.spark") \
        .option("table", "<TABLENAME>") \
        .option("zkUrl", "<hostname>:<port>") \
        .load() 

我想知道使用自定义 SQL 的对应等价物是什么

sqlResult =  sqlContext.read \
             .format("org.apache.phoenix.spark") \
             .option("sql", "select * from <TABLENAME> where <CONDITION>") \
             .option("zkUrl", "<HOSTNAME>:<PORT>") \
             .load()

谢谢。

【问题讨论】:

    标签: apache-spark pyspark apache-spark-sql spark-dataframe phoenix


    【解决方案1】:

    这可以使用 Phoenix 作为 JDBC 数据源来完成,如下所示:

    sql = '(select COL1, COL2 from TABLE where COL3 = 5) as TEMP_TABLE'
    
    df = sqlContext.read.format('jdbc')\
           .options(driver="org.apache.phoenix.jdbc.PhoenixDriver", url='jdbc:phoenix:<HOSTNAME>:<PORT>', dbtable=sql).load()
    
    df.show() 
    

    但需要注意的是,如果 SQL 语句中有列别名,那么 .show() 语句会抛出异常(如果您使用 .select() 选择没有别名的列,它将起作用) ,这可能是 Phoenix 中的一个错误。

    【讨论】:

    • 这是问题的答案还是部分问题?
    • 两者。它使用 JDBC 来实现我想要做的事情,但使用 Phoenix Spark 选项会更好,因此我尝试了它以及相应的错误消息。
    • 问题应该在第一篇文章中编辑,因为这是答案部分。 stackoverflow 不像普通的论坛。
    【解决方案2】:

    在这里您需要使用 .sql 来处理自定义查询。这是语法

    dataframe = sqlContext.sql("select * from <table> where <condition>")
    dataframe.show()
    

    【讨论】:

    【解决方案3】:

    对于 Spark2,我对 .show() 函数没有问题,并且我没有使用 .select() 函数来打印来自 Phoenix 的 DataFrame 的所有值。 因此,请确保您的 sql 查询已在括号内,请看我的示例:

     val sql = " (SELECT  P.PERSON_ID as PERSON_ID, P.LAST_NAME as LAST_NAME, C.STATUS as STATUS FROM PERSON P INNER JOIN CLIENT C ON C.CLIENT_ID = P.PERSON_ID) "
              val dft = dfPerson.sparkSession.read.format("jdbc")
                .option("driver", "org.apache.phoenix.jdbc.PhoenixDriver")
                .option("url", "jdbc:phoenix:<HOSTNAME>:<PORT>")
                .option("useUnicode", "true")
                .option("continueBatchOnError", "true")
                .option("dbtable", sql)
                .load()
    dft.show();
    

    它告诉我:

    +---------+--------------------+------+
    |PERSON_ID|           LAST_NAME|STATUS|
    +---------+--------------------+------+
    |     1005|             PerDiem|Active|
    |     1008|NAMEEEEEEEEEEEEEE...|Active|
    |     1009|           Admission|Active|
    |     1010|            Facility|Active|
    |     1011|                MeUP|Active|
    +---------+--------------------+------+
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2019-10-29
      • 2020-08-07
      • 2016-02-06
      • 1970-01-01
      • 2019-07-28
      • 1970-01-01
      • 2022-01-18
      相关资源
      最近更新 更多