【发布时间】: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