【问题标题】:JSON schema inference in Structured Streaming with Kafka as source以 Kafka 为源的结构化流中的 JSON 模式推断
【发布时间】: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 选项,但我看不到如何使用它。

【问题讨论】:

    标签: apache-spark apache-kafka


    【解决方案1】:

    对于 KAFKA 是不可能的。花费太多时间。对于文件源,您可以。

    来自手册:

    流数据帧/数据集的架构推断和分区

    默认情况下,来自基于文件的源的结构化流式传输需要您 指定模式,而不是依赖 Spark 来推断它 自动地。此限制可确保一致的架构 用于流式查询,即使在失败的情况下。对于临时 用例,您可以通过设置重新启用模式推断 spark.sql.streaming.schemaInference 为真。

    但对于不是 KAFKA 的文件源。

    【讨论】:

      猜你喜欢
      • 2018-06-29
      • 1970-01-01
      • 2017-09-25
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-06-04
      • 2022-10-16
      • 1970-01-01
      相关资源
      最近更新 更多