【问题标题】:How to save Spark Stream data to file如何将 Spark Stream 数据保存到文件
【发布时间】:2020-06-04 05:29:18
【问题描述】:

我是 Spark 的新手,目前正在解决与在 Context 时间之后将 Spark Stream 的结果保存到文件相关的问题。所以问题是:我希望一个查询运行 60 秒,并将它在这段时间内读取的所有输入保存到一个文件中,并且还能够定义文件名以供将来处理。

最初我认为下面的代码是可行的方法:

sc.socketTextStream("localhost", 12345)
                .foreachRDD(rdd -> {
                    rdd.saveAsTextFile("./test");
                });

但是,在运行之后,我意识到它不仅为每个输入读取保存了一个不同的文件 - (想象我在该端口上以随机速度生成随机数),所以如果在一秒钟内它读取 1 文件会包含 1 个数字,但如果它读取更多文件将包含它们,而不是仅写入一个包含 60 年代时间范围内所有值的文件 - 而且我无法命名文件,因为 saveAsTextFile 中的参数 是所需的目录。

所以想问一下有没有spark原生的解决方案这样我就不用“java技巧”来解决了,像这样:

sc.socketTextStream("localhost", 12345)
                .foreachRDD(rdd -> {
                    PrintWriter out = new PrintWriter("./logs/votes["+dtf.format(LocalDateTime.now().minusMinutes(2))+","+dtf.format(LocalDateTime.now())+"].txt");
                    List<String> l = rdd.collect();
                    for(String voto: l)
                        out.println(voto + "    "+dtf.format(LocalDateTime.now()));
                    out.close();
                });

我搜索了类似问题的 spark 文档,但找不到解决方案:/ 谢谢你的时间:)

【问题讨论】:

  • 收集绝不是把戏
  • 我的意思是使用 java 默认 PrintWriter 将字符串保存到文件,而不是使用(我认为必须存在的)火花解决方案。 TBH 我在理解 foreachRDD 的工作原理时遇到了一些麻烦,因为在上面使用 saveAsTextFile 显示的案例中,它只保存一个值,但在其他情况下它适用于所有数据
  • Spark 就是这样

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


【解决方案1】:

以下是使用新 Spark API 使用套接字流数据的模板。

import org.apache.spark.sql.streaming.{OutputMode, Trigger}

object ReadSocket {

  def main(args: Array[String]): Unit = {
    val spark = Constant.getSparkSess

    //Start reading from socket
    val dfStream = spark.readStream
      .format("socket")
      .option("host","127.0.0.1") // Replace it your socket host
      .option("port","9090")
      .load()

    dfStream.writeStream
      .trigger(Trigger.ProcessingTime("1 minute")) // Will trigger 1 minute
      .outputMode(OutputMode.Append) // Batch will processed for the data arrived in last 1 minute
      .foreachBatch((ds,id) => { //
        ds.foreach(row => { // Iterdate your data set
          //Put around your File generation logic
          println(row.getString(0)) // Thats your record
        })
      }).start().awaitTermination()
  }

}

代码解释请阅读内联cmets代码

Java 版本

import org.apache.spark.api.java.function.VoidFunction2;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Encoders;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.streaming.StreamingQueryException;
import org.apache.spark.sql.streaming.Trigger;

public class ReadSocketJ {

    public static void main(String[] args) throws StreamingQueryException {
        SparkSession spark = Constant.getSparkSess();


        Dataset<Row> lines = spark
                .readStream()
                .format("socket")
                .option("host", "127.0.0.1") // Replace it your socket host
                .option("port", "9090")
                .load();

        lines.writeStream()
                .trigger(Trigger.ProcessingTime("5 seconds"))
                .foreachBatch((VoidFunction2<Dataset<Row>, Long>) (v1, v2) -> {
                    v1.as(Encoders.STRING())
                            .collectAsList().forEach(System.out::println);
                }).start().awaitTermination();


    }
}

【讨论】:

  • 我正在使用 Java 开发:/ 但是你能解释一下这个概念,以便我尝试适应它吗?
  • 非常感谢您的更新!但问题是如何保存到文件(并且能够定义文件名,因为 saveAsText 只允许我选择目录)。 Spark没有为此提供本机解决方案吗?必须使用 java 解决方案来存储字符串似乎很奇怪
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-11-26
  • 2011-07-20
  • 2016-09-01
  • 1970-01-01
  • 1970-01-01
  • 2011-08-17
相关资源
最近更新 更多