【问题标题】:How to make sure that write csv is complete?如何确保写入 csv 完成?
【发布时间】:2023-04-03 05:35:01
【问题描述】:
我正在将数据集写入 CSV,如下所示:
df.coalesce(1)
.write()
.format("csv")
.option("header", "true")
.mode(SaveMode.Overwrite)
.save(sink);
sparkSession.streams().awaitAnyTermination();
我如何确保当流式传输作业终止时,输出正确完成?
如果我过早/过晚终止,接收器文件夹会被覆盖并且为空。
附加信息:特别是如果主题没有消息,我的 spark 作业仍在运行并用空文件覆盖结果。
【问题讨论】:
标签:
java
apache-spark
spark-structured-streaming
【解决方案1】:
我如何确保当流式传输作业终止时,输出正确完成?
Spark Structured Streaming 的工作方式是流式查询(作业)连续运行,并且“当流式作业终止时,输出正确完成”。
我要问的问题是流式查询是如何终止的。这是StreamingQuery.stop 还是Ctrl-C / kill -9?
如果流式查询以强制方式终止 (Ctrl-C / kill -9),那么,您将得到您所要求的 - 部分执行无法确保输出正确,因为该过程(流式查询)被强制关闭。
使用StreamingQuery.stop,流式查询将优雅地终止并写出当时的所有内容。
我有一个问题,如果我过早/过晚终止,接收器文件夹会被覆盖,并且该文件夹是空的。
如果您终止得太早/太晚,由于流式查询无法完成其工作,您还能期待什么。你应该优雅地stop 它,你会得到预期的输出。
附加信息:特别是如果主题没有消息,我的 spark 作业仍在运行并用空文件覆盖结果。
这是一个有趣的观察结果,需要进一步探索。
如果没有没有消息要处理,则不会触发批处理,因此没有作业,因此没有“用空文件覆盖结果。”(因为没有任务会被执行)。
【解决方案2】:
首先,我看到您没有使用过writeStream 我不太确定您的工作如何成为流媒体工作。
现在,回答您的问题 1,您可以使用 StreamingQueryListener 来监控流式查询的进度。让另一个 StreamingQuery 从输出位置读取。监视它。将文件放在输出位置后,使用StreamingQueryListener 中的查询名称和输入记录计数来优雅地stop 任何查询。 awaitAnyTermination 应该停止您的 spark 应用程序。以下代码可能会有所帮助。
spark.streams.addListener(new StreamingQueryListener() {
override def onQueryStarted(event: QueryStartedEvent) {
//logger message to show that the query has started
}
override def onQueryProgress(event: QueryProgressEvent) {
synchronized {
if(event.progress.name.equalsIgnoreCase("QueryName"))
{
recordsReadCount = recordsReadCount + event.progress.numInputRows
//Logger messages to show continuous progress
}
}
}
override def onQueryTerminated(event: QueryTerminatedEvent) {
synchronized {
//logger message to show the reason of termination.
}
}
})
回答你的第二个问题,我也不认为这是可能的,正如 Jacek 的回答中提到的那样。