【问题标题】:structured streaming write to different parquet folders结构化流写入不同的镶木地板文件夹
【发布时间】:2019-01-20 01:29:45
【问题描述】:

我正在使用 spark 结构化流从 kafka 主题中读取事件并对其进行处理并写入 parquet。我必须根据我在事件中获得的密钥将输出写入不同的文件夹。我尝试使用结构化流示例总是指向一个特定的文件夹。我需要为每个文件夹启动一个流吗?

df.writeStream.format("parquet").option("path", "path/to/destination/dir").start()

【问题讨论】:

  • 您可以根据键对父目录进行分区,这将为每个值创建一个文件夹
  • 谢谢,但关键是在消息中,但要写入镶木地板,我需要在 writeStream.start 时指定它。但实际数据将在稍后发布。那么如何在分区中指定呢?
  • 我可以从 datafame 中读取密钥并使用分区中的值。 dataframe.writeStream.format("parquet") .option("path",path) 但它不工作

标签: apache-spark apache-spark-sql parquet spark-structured-streaming


【解决方案1】:

我能够通过创建多个 writeStreams 来实现这一点,每个 writeStreams 都特定于一个表

详情请参考structured streaming different schema in nested json

【讨论】:

    【解决方案2】:

    你可以使用 foreachBatch

    streamingDF.writeStream.foreachBatch { (batchDF: DataFrame, batchId: Long) =>
      batchDF.persist()
      batchDF.write.format(...).save(...)  // location 1
      batchDF.write.format(...).save(...)  // location 2
      batchDF.unpersist()
    }
    

    更多信息可以参考spark documenation foreach and foreachBatch

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2020-02-25
      • 1970-01-01
      • 2021-08-27
      • 2019-06-02
      • 2023-03-19
      • 2020-04-02
      • 2019-10-29
      相关资源
      最近更新 更多