【发布时间】:2018-10-08 18:34:21
【问题描述】:
我有一张蜂巢桌。我想创建动态火花 SQL 查询。在火花提交时,我指定规则名称。基于规则名称的查询应该生成。在提交火花时,我必须指定规则名称。例如:
sparks-submit <RuleName> IncorrectAge
它应该触发我的 scala 目标代码:
select tablename, filter, condition from all_rules where rulename="IncorrectAge"
我的表:规则(输入表)
|---------------------------------------------------------------------------|
| rowkey| rule_name|rule_run_status| tablename |condition|filter |level|
|--------------------------------------------------------------------------|
| 1 |IncorrectAge| In_Progress | VDP_Vendor_List| age>18 gender=Male|NA|
|---------------------------------------------------------------------------
|2 | Customer_age| In_Progress | Customer_List | age<25 gender=Female|NA|
|----------------------------------------------------------------------------
我获取规则名称:
select tablename, filter, condition from all_rules where rulename="IncorrectAge";
执行此查询后,我得到如下结果:
|----------------------------------------------|
|tablename | filter | condition |
|----------------------------------------------|
|VDP_Vendor_List | gender=Male | age>18 |
|----------------------------------------------|
现在我想动态进行 spark sql 查询
select count(*) from VDP_Vendor_List // first column --tablename
select count(*) from VDP_Vendor_List where gender=Male --tablename and filter
select * from EMP where gender=Male AND age >18 --tablename, filter, condition
我的代码-Spark 2.2 版本代码:
import org.apache.spark.sql.{ Row, SparkSession }
import org.apache.log4j._
object allrules {
def main(args: Array[String]) {
val spark = SparkSession.builder().master("local[*]")
.appName("Spark Hive")
.enableHiveSupport().getOrCreate();
import spark.implicits._
val sampleDF = spark.read.json("C:/software/sampletableCopy.json") // for testing purpose i converted hive table to json data
sampleDF.registerTempTable("sampletable")
val allrulesDF = spark.sql("SELECT * FROM sampletable")
allrulesDF.show()
val TotalCount: Long = allrulesDF.count()
println("==============> Total count ======>" + allrulesDF.count())
val df1 = allrulesDF.select(allrulesDF.col("tablename"),allrulesDF.col("condition"),allrulesDF.col("filter"),allrulesDF.col("rule_name"))
df1.show()
val df2= df1.where(df1.col("rule_name").equalTo("IncorrectAge")).show()
println(df2)
// var table_name = ""
// var condition =""
// var filter = "";
// df1.foreach(row=>{
// table_name = row.get(1).toString();
// condition = row.get(2).toString();
// filter = row.get(3).toString();
// })
}
}
【问题讨论】:
-
我认为您的问题并不完全清楚。您能否在末尾添加几行来准确说明您要查找的内容?
-
从 all_rules 中选择表名、过滤器、条件 where rulename="IncorrectAge";这里 IncorrectAge 是我的规则名称。我正在使用 3 个属性(表名、过滤器、条件)。对于第一个查询,我只使用一个属性。例如,从 VDP_Vendor_List // 第一列 --tablename 中选择 count()。在第二个查询中,我使用了 2 个属性。例如- 从 VDP_Vendor_List 中选择 count(),其中 gender=Male --tablename 并过滤。
-
欢迎来到 Stack Overflow!其他用户将您的问题标记为低质量和需要改进。我重新措辞/格式化您的输入,使其更易于阅读/理解。请查看我的更改以确保它们反映您的意图。但我认为你的问题仍然无法回答。 你现在应该edit你的问题,添加缺失的细节(见minimal reproducible example)。如果您对我有其他问题或反馈,请随时给我留言。
-
请:A)永远不要在评论中添加更多信息,他们应该进入问题。 B) 尝试提出一个仅包含与您的问题相关的信息的最小问题。
标签: scala apache-spark apache-spark-sql