【问题标题】:Reading schema of streaming Dataframe in Spark Structured Streaming [duplicate]在 Spark Structured Streaming 中读取流数据帧的模式 [重复]
【发布时间】:2021-04-26 04:33:16
【问题描述】:

我是 Apache Spark 结构化流的新手。我正在尝试从事件中心读取一些事件(以 XML 格式)并尝试从嵌套的 XML 创建新的 Spark DF。

我正在使用https://github.com/databricks/spark-xml 中描述的代码示例,并且在批处理模式下运行良好,但在结构化 Spark 流中却没有。

spark-xml Github 库的代码块

import com.databricks.spark.xml.functions.from_xml
import com.databricks.spark.xml.schema_of_xml
import spark.implicits._
val df = ... /// DataFrame with XML in column 'payload' 
val payloadSchema = schema_of_xml(df.select("payload").as[String])
val parsed = df.withColumn("parsed", from_xml($"payload", payloadSchema))

我的批处理代码

val df = Seq(
  (8, "<AccountSetup xmlns:xsi=\"test\"><Customers test=\"a\">d</Customers><tag1>7</tag1> <tag2>4</tag2> <mode>0</mode> <Quantity>1</Quantity></AccountSetup>"),
  (64, "<AccountSetup xmlns:xsi=\"test\"><Customers test=\"a\">d</Customers><tag1>6</tag1> <tag2>4</tag2>  <mode>0</mode> <Quantity>1</Quantity></AccountSetup>"),
  (27, "<AccountSetup xmlns:xsi=\"test\"><Customers test=\"a\">d</Customers><tag1>4</tag1> <tag2>4</tag2> <mode>3</mode> <Quantity>1</Quantity></AccountSetup>")
).toDF("number", "body")
)


val payloadSchema = schema_of_xml(df.select("body").as[String])
val parsed = df.withColumn("parsed", from_xml($"body", payloadSchema))

val final_df = parsed.select(parsed.col("parsed"))
display(final_df.select("parsed.*"))

我试图对 Spark Structured Streaming 执行相同的逻辑,如下代码:

结构化流代码

import com.databricks.spark.xml.functions.from_xml
import com.databricks.spark.xml.schema_of_xml
import org.apache.spark.eventhubs.{ ConnectionStringBuilder, EventHubsConf, EventPosition }
import spark.implicits._
import org.apache.spark.sql.functions._
import org.apache.spark.sql.types._


val streamingInputDF = 
  spark.readStream
    .format("eventhubs")
    .options(eventHubsConf.toMap)
    .load()

val payloadSchema = schema_of_xml(streamingInputDF.select("body").as[String])
val parsed = streamingSelectDF.withColumn("parsed", from_xml($"body", payloadSchema))
val final_df = parsed.select(parsed.col("parsed"))

display(final_df.select("parsed.*"))

val payloadSchema = schema_of_xml(streamingInputDF.select("body").as[String]) 指令的代码部分抛出错误 Queries with streaming sources must be executed with writeStream.start();;

更新

尝试过


val streamingInputDF = 
  spark.readStream
    .format("eventhubs")
    .options(eventHubsConf.toMap)
    .load()
    .select(($"body").cast("string"))

val body_value = streamingInputDF.select("body").as[String]
body_value.writeStream
    .format("console")
    .start()

spark.streams.awaitAnyTermination()


val payloadSchema = schema_of_xml(body_value)
val parsed = body_value.withColumn("parsed", from_xml($"body", payloadSchema))
val final_df = parsed.select(parsed.col("parsed"))

现在没有遇到错误,但 Databricks 保持“等待状态”

谢谢!!

【问题讨论】:

    标签: xml scala apache-spark spark-structured-streaming azure-eventhub


    【解决方案1】:

    如果您的代码在批处理模式下工作,则它没有任何问题。

    不仅将源转换为流(通过使用readStreamload)很重要,而且还需要将接收器部分转换为流。

    您收到的错误消息只是提醒您还要查看水槽部分。您的 Dataframe final_df 实际上是一个 streaming Dataframe,必须通过 start 启动。

    结构化流式传输指南为您提供了所有可用Output Sinks 的概览,最简单的方法是将结果打印到控制台。

    总而言之,您需要在程序中添加以下内容:

    final_df.writeStream
        .format("console")
        .start()
    
    spark.streams.awaitAnyTermination()
    

    【讨论】:

    • 非常感谢迈克的回复。我真的需要更多关于 Spark Structured Streaming 的详细知识。无论如何,有些东西对我不起作用,为什么final_df.writeStream?如果我运行schema_of_xml(streamingInputDF.select("body").as[String]),在我到达 final_df 之前,它已经因该错误而失败了
    • 因为您的 streamingInputDF 是一个流数据帧,并且您将它用于您的 payloadSchema,而 payloadSchema 没有“writeStream”和“start”。
    • 如果其他 Stackoverflow 帖子中的答案没有帮助,我建议创建一个最小的可重复示例并打开一个新问题,专门指出您遇到的错误/问题。
    • 我不知道我可以在多大程度上使它更可重现,它只是有一个 XML 文档(如批处理代码示例中)并将其作为 Stream 读取。也许我误解了你。至于另一篇文章,我看到它非常针对缓存,我想我仍然不知道很多关于结构化流的基本概念。
    • 非常感谢 Mike,我会这样做的,我会尝试用更简单的例子 :)
    猜你喜欢
    • 2019-10-26
    • 2021-04-24
    • 1970-01-01
    • 2023-03-08
    • 2017-04-25
    • 1970-01-01
    • 2020-09-12
    • 2021-12-05
    • 2019-01-07
    相关资源
    最近更新 更多