【问题标题】:Dynamic query preparation and execution in sparkspark中的动态查询准备和执行
【发布时间】:2019-04-28 19:42:49
【问题描述】:

在 Spark 中,这个 json 在 dataframe(DF) 中,现在我们必须导航到表(在基于 cust 的 json 中),我们必须读取表的第一个块并且必须准备 sql 查询。 例如:SELECT CUST_NAME FROM CUST WHERE CUST_ID =112

我们必须在数据库中执行此查询并将结果存储在 json 文件中。

{
     "cust": "Retails",
     "tables": [
        {
             "Name":"customer",
             "table_NAME":"cust",
             "param1":"cust_id",  
             "val":"112",
             "op":"cust_name"
        },
        {
             "Name":"sales",
             "table_NAME":"sale",
             "param1":"country",  
             "val":"ind",
             "op":"monthly_sale"
         }]
}

 root |-- cust: string (nullable = true) 
      |-- tables: array (nullable = true) 
      | |-- element: struct (containsNull = true) 
      | | |-- Name: string (nullable = true) 
      | | |-- op: string (nullable = true) 
      | | |-- param1: string (nullable = true) 
      | | |-- table_NAME: string (nullable = true) 
      | | |-- val: string (nullable = true) 

第二个表块相同。 例如:SELECT MONTHLY_SALE FROM SALE WHERE COUNTRY = 'IND'

必须在数据库中执行此查询,并且必须将此结果存储在上述 json 文件中。

执行此操作的最佳方法是什么?有什么想法吗?

【问题讨论】:

  • 你能发布 df.printSchema 的输出吗?
  • root |-- cust: 字符串 (nullable = true) |-- 表格: 数组 (nullable = true) | |-- 元素:结构 (containsNull = true) | | |-- 名称:字符串(可为空=真)| | |-- 操作:字符串(可为空=真)| | |-- 参数1:字符串(可为空=真)| | |-- 表名:字符串(可为空=真)| | |-- val: string (nullable = true)

标签: apache-spark apache-spark-sql apache-spark-mllib


【解决方案1】:

这是我实现这一目标的方式。对于整个解决方案,我使用了 spark-shell。以下是一些先决条件:

  1. json-serde下载这个jar

  2. 将 zip 文件解压到任意位置

  3. 现在使用这个命令运行 spark-shell

    spark-shell --jars path/to/jars/json-serde-cdh5-shim-1.3.7.3.jar,path/to/jars/json-serde-1.3.7.3.jar,path/to/jars/json-1.3.7.3.jar
    

您的 Json 文档:

{
 "cust": "Retails",
 "tables": [
    {
         "Name":"customer",
         "table_NAME":"cust",
         "param1":"cust_id",  
         "val":"112",
         "op":"cust_name"
    },
    {
         "Name":"sales",
         "table_NAME":"sale",
         "param1":"country",  
         "val":"ind",
         "op":"monthly_sale"
     }]
}

折叠版:

{"cust": "Retails","tables":[{"Name":"customer","table_NAME":"cust","param1":"cust_id","val":"112","op":"cust_name"},{"Name":"sales","table_NAME":"sale","param1":"country","val":"ind","op":"monthly_sale"}]}

我已将此 json 放入此 /tmp/sample.json

现在进入 spark-sql 部分:

  1. 基于 json 架构创建表

    sql("CREATE TABLE json_table(cust string,tables array<struct<Name: string,table_NAME:string,param1:string,val:string,op:string>>) ROW FORMAT SERDE 'org.openx.data.jsonserde.JsonSerDe'")
    
  2. 现在将json数据加载到表中

    sql("LOAD DATA LOCAL INPATH  '/tmp/sample.json' OVERWRITE INTO TABLE json_table")
    
  3. 现在我将使用 hive 横向视图概念Lateral view

    val ans=sql("SELECT myCol FROM json_table LATERAL VIEW explode(tables) myTable as myCol").collect
    
  4. 返回结果的架构:

        ans.printSchema
        root
         |-- table: struct (nullable = true)
         |    |-- Name: string (nullable = true)
         |    |-- table_NAME: string (nullable = true)
         |    |-- param1: string (nullable = true)
         |    |-- val: string (nullable = true)
         |    |-- op: string (nullable = true)
    
  5. ans.show 的结果

         ans.show
         +--------------------+
         |               table|
         +--------------------+
         |[customer,cust,cu...|
         |[sales,sale,count...|
         +--------------------+
    
  6. 现在我假设可以有两种类型的数据,例如cust_idNumber 类型,countryString强>类型。我正在添加一种方法来根据数据的值来识别数据的类型。例如

    def isAllDigits(x: String) = x forall Character.isDigit
    

    注意:您可以使用自己的方式来识别它

7.现在基于json数据创建查询

    ans.foreach(f=>{
val splitted_string=f.toString.split(",")
val op=splitted_string(4).substring(0,splitted_string(4).size-2)
val table_NAME=splitted_string(1)
val param1 = splitted_string(2)
val value = splitted_string(3)
if(isAllDigits(value)){
println("SELECT " +op+" FROM "+ table_NAME+" WHERE "+param1+"="+value)
}else{
println("SELECT " +op+" FROM "+ table_NAME+" WHERE "+param1+"='"+value+"'")
}
})

这是我得到的结果:

SELECT cust_name FROM cust WHERE cust_id=112
SELECT monthly_sale FROM sale WHERE country='ind'

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2019-01-02
    • 2016-08-29
    • 1970-01-01
    • 1970-01-01
    • 2014-07-11
    • 2018-11-20
    • 2016-07-29
    相关资源
    最近更新 更多