【发布时间】:2023-04-11 08:57:02
【问题描述】:
我目前正在使用 Spark Structured Steaming 从 Kafka 主题中读取 json 数据。 json 以字符串形式存储在主题中。为了实现这一点,我提供了一个硬编码的 JSON 模式作为 StructType。我正在寻找一种在流式传输期间动态推断主题架构的好方法。
这是我的代码: (是 Kotlin,不是常用的 Scala)
spark
.readStream()
.format("kafka")
.option("kafka.bootstrap.servers", "kafka:9092")
.option("subscribe", "my_topic")
.option("startingOffsets", "latest")
.option("failOnDataLoss", "false")
.load()
.selectExpr("CAST(value AS STRING)")
.select(
from_json(
col("value"),
JsonSchemaRegistry.mySchemaStructType)
.`as`("data")
)
.select("data.*")
.writeStream()
.format("my_format")
.option("checkpointLocation", "/path/to/checkpoint")
.trigger(ProcessingTime("25 seconds"))
.start("/my/path")
.awaitTermination()
现在这是否可能以一种干净的方式进行,而无需为每个 DataFrame 再次推断它?我正在寻找一些惯用的方式。如果在结构化流中不建议进行模式推断,我将继续对我的模式进行硬编码,但要确定。 Spark 文档中提到了 spark.sql.streaming.schemaInference 选项,但我看不到如何使用它。
【问题讨论】: