【发布时间】:2019-08-06 15:52:15
【问题描述】:
我有一个 kafka 流,我将在 spark 中处理。我想将此流的输出写入文件。但是,我想按天对这些文件进行分区,所以每天它都会开始写入一个新文件。可以做这样的事情吗?我希望它继续运行,当新的一天发生时,它将切换到写入新文件。
val streamInputDf = spark.readStream.format("kafka")
.option("kafka.bootstrapservers", "XXXX")
.option("subscribe", "XXXX")
.load()
val streamSelectDf = streamInputDf.select(...)
streamSelectDf.writeStream.format("parquet)
.option("path", "xxx")
???
【问题讨论】:
-
为什么不将Kafka中的数据直接消费到Spark中呢?
-
因为出于审计目的,我们必须像我们拥有的所有其他数据流一样运行此数据流(每天在设定的时间)。所以我所要做的就是获取稍后将处理的数据。 @罗宾莫法特。我想过用
.trigger(ProcessingTime("24 hours"))做writeStream,但我不知道如何让文件被写入,真正改变 -
你提到了 Kafka,但实际上 Kafka 本身应该只是充当消息总线。如果您使用的是 Cloudera/Hortonworks Data Flow 平台,您可以使用 NiFi 将数据移入/移出 Kafka,否则您可以使用 Spark 或 Kafka Connect 等工具来填补这个角色。
-
@DennisJaheruddin 正确。对不起,我忘了提。我将使用 spark 将这些数据处理到文件中。但是我不确定如何按日期对消息进行分区,以便在新的一天发生时将文件放入一个新文件中(对原始帖子的小编辑)
标签: apache-spark apache-kafka spark-streaming