【问题标题】:Can we convert data frame in Databricks to string and why do we get error Queries with streaming sources must be executed with writeStream.start()我们可以将 Databricks 中的数据帧转换为字符串吗?为什么会出现错误 必须使用 writeStream.start() 执行带有流源的查询
【发布时间】:2020-09-01 20:15:13
【问题描述】:

我正在选择作为数据框的列。我想将它转换为字符串,以便它可以用于构建 cosmos DB 动态查询。数据帧上的collect()函数抱怨流源查询必须用writeStream.start();;

val DF = AppointmentDF
            .select("*")
            .filter($"xyz" === "abc")

DF.createOrReplaceTempView("MyTable")
val column1DF = spark.sql("SELECT column1 FROM MyTable")


// This is not getting resolved
val sql="select c.abc from c where c.column = \"" + String.valueOf(column1DF) + "\""
println(sql)

Error:
org.apache.spark.sql.AnalysisException: cannot resolve '`column1DF`' given input columns: []; line 1 pos 12;


DF.collect().foreach { row =>
  println(row.mkString(","))
 } 

 Error: 
 org.apache.spark.sql.AnalysisException: Queries with streaming sources must 
 be executed with writeStream.start();;

【问题讨论】:

    标签: scala apache-spark-sql databricks


    【解决方案1】:

    数据帧是一种分布式数据结构,而不是位于您的机器中可以打印的结构。 DFcolumn1DF 的值将完全是数据帧。要将查询的所有数据带到驱动程序节点,您可以使用数据框方法collect,并从返回的行数组中提取您的值。 如果您将千兆字节的数据带入驱动程序节点的内存,则收集可能是有害的。

    【讨论】:

      【解决方案2】:

      您可以使用collect 并使用head 获取DataFrame 的第一行:

      val column1DF = spark.sql("SELECT column1 FROM MyTable").collect().head.getAs[String](0)
      val sql="select c.abc from c where c.column = \"" + column1DF + "\""
      

      【讨论】:

      • 我在 scala notebook 命令中并得到了这个。 '带有流源的查询必须使用 writeStream.start();;;'
      猜你喜欢
      • 2021-10-13
      • 2017-03-29
      • 1970-01-01
      • 2021-01-31
      • 2018-03-14
      • 2017-06-23
      • 2021-08-16
      • 1970-01-01
      • 2019-05-25
      相关资源
      最近更新 更多