【问题标题】:how to change json schema in spark spark streamning events without interrupting the streamning Job?如何在不中断流式传输作业的情况下更改 spark spark 流式事件中的 json 架构?
【发布时间】:2021-10-29 01:27:08
【问题描述】:

我有一个用例,我需要在不中断流式传输作业的情况下更改 JSON 的架构。我正在使用一个 conf 文件,其中提到了所有必需的架构。我已经尝试过使用单独的流管道持久化和取消持久化缓存和广播变量,但仍然没有运气。提前感谢您的帮助!

【问题讨论】:

    标签: scala apache-spark pyspark apache-spark-sql spark-streaming


    【解决方案1】:

    除了将数据集读取为 json 之外,您还可以尝试将其读取为文本,然后根据从 HDFS 或 DB 中的配置文件外部来的架构映射它。

    所以不要做类似的事情,

    val df = spark.readStream.format("json").load(.. path ..) 
    

    做,

    import sparkSession.implicits._
    
    val df = spark.readStream
    .format("text").load( .. path .. )
    .select("value")
    .as[String]
    .mapPartitions(partStrings => {
        val currentSchema = readSchemaFromFile(???)
        partStrings.map(str => parseJSON(currentSchema, str))
    })
    

    mapPartitions 阻止对每条记录进行架构查找。

    【讨论】:

      猜你喜欢
      • 2018-05-27
      • 1970-01-01
      • 2020-07-17
      • 2018-06-29
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-10-17
      相关资源
      最近更新 更多