【问题标题】:How to write outputs of spark streaming application to a single file如何将火花流应用程序的输出写入单个文件
【发布时间】:2019-12-24 09:03:05
【问题描述】:

我正在使用火花流从 Kafka 读取数据并传递到 py 文件进行预测。它返回预测以及原始数据。它将原始数据及其预测保存到文件中,但是它为每个 RDD 创建了一个文件。 我需要一个包含所有收集到的数据的文件,直到我停止程序才能保存到一个文件中。

我试过 writeStream 它甚至不会创建一个文件。 我尝试使用 append 将其保存到镶木地板,但它会为每个 RDD 创建多个文件,即 1 个文件。 我尝试使用附加模式写入多个文件作为输出。 下面的代码创建一个文件夹 output.csv 并将所有文件输入其中。

 def main(args: Array[String]): Unit = {
    val ss = SparkSession.builder()
      .appName("consumer")
      .master("local[*]")
      .getOrCreate()

    val scc = new StreamingContext(ss.sparkContext, Seconds(2))


    val kafkaParams = Map[String, Object](
        "bootstrap.servers" -> "localhost:9092",
        "key.deserializer"-> 
"org.apache.kafka.common.serialization.StringDeserializer",
        "value.deserializer"> 
"org.apache.kafka.common.serialization.StringDeserializer",
        "group.id"-> "group5" // clients can take
      )
mappedData.foreachRDD(
      x =>
    x.map(y =>       
ss.sparkContext.makeRDD(List(y)).pipe(pyPath).toDF().repartition(1)
.write.format("csv").mode("append").option("truncate","false")
.save("output.csv")
          )
    )
scc.start()
scc.awaitTermination()

我只需要获取 1 个文件,其中包含在流式传输时一一收集的所有语句。

任何帮助将不胜感激,谢谢您的期待。

【问题讨论】:

    标签: apache-spark apache-spark-sql streaming spark-streaming csv-write-stream


    【解决方案1】:

    您无法修改 hdfs 中的任何文件,一旦它被写入。如果您希望实时写入文件(每 2 秒将来自流式作业的数据块附加到同一文件中),则根本不允许这样做,因为 hdfs 文件是不可变的。如果可能的话,我建议你尝试编写一个读取多个文件的读取逻辑。

    但是,如果您必须从单个文件中读取,我建议您在将输出写入单个 csv/parquet 文件夹后使用两种方法中的一种,使用“Append”SaveMode(这将为每个块创建部分文件你每 2 秒写一次)。

    1. 您可以在此文件夹的顶部创建一个配置单元表,从该表中读取数据。
    2. 您可以在 spark 中编写一个简单的逻辑来读取包含多个文件的文件夹,然后使用 reparation(1) 或 coalesce(1) 将其作为单个文件写入另一个 hdfs 位置,然后从该位置读取数据。见下文:

      spark.read.csv("oldLocation").coalesce(1).write.csv("newLocation")
      

    【讨论】:

    • #Ana 好吧,你的意思是我不能将东西附加到 spark 流应用程序的单个文件中,我必须编写另一个代码来定期从 spark 流写入的文件中读取,并且追加到另一个文件是吗?
    • 是和不是。是的,您需要编写另一个代码。不,此代码也不会在任何文件中追加数据,因为在 HDFS 中无法追加,相反,它会通过组合先前位置的较小文件来创建一个大的新文件。
    • 按照 Raja 的建议在示例中更新 repartition 以合并,这在减少分区数量时会表现得更好。
    【解决方案2】:

    重新分区 - 建议在增加分区数量的同时使用重新分区,因为它涉及所有数据的洗牌。

    coalesce - 建议在减少分区数量的同时使用合并。例如,如果您有 3 个分区并且您想将其减少到 2 个分区,Coalesce 会将第 3 个分区的数据移动到分区 1 和 2。分区 1 和 2 将保留在同一个 Container 中。但是重新分区将在所有分区中打乱数据,以便网络使用执行者之间的距离会很高,并且会影响性​​能。

    在减少分区数量的同时,性能方面的合并性能优于重新分区。

    所以在编写使用选项作为合并时。 例如:df.write.coalesce

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2017-04-28
      • 1970-01-01
      • 2015-10-13
      • 2020-06-21
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多