【问题标题】:Spark Structured Streaming writestream doesn't write file until I stop the job在我停止工作之前,Spark Structured Streaming writestream 不会写入文件
【发布时间】:2019-07-22 02:56:41
【问题描述】:

我在一个经典用例上使用 Spark 结构化流:我想从 kafka 主题中读取数据并将流以 parquet 格式写入 HDFS。

这是我的代码:

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
import org.apache.spark.sql.streaming.Trigger
import org.apache.spark.sql.types.{ArrayType, DataTypes, StructType}

object TestKafkaReader extends  App{
  val spark = SparkSession
    .builder
    .appName("Spark-Kafka-Integration")
    .master("local")
    .getOrCreate()
  spark.sparkContext.setLogLevel("ERROR")
  import spark.implicits._

  val kafkaDf = spark
    .readStream
    .format("kafka")
    .option("kafka.bootstrap.servers","KAFKA_BROKER_IP:PORT")
    //.option("subscribe", "test")
    .option("subscribe", "test")
    .option("startingOffsets", "earliest")
    .load()

  val moviesJsonDf = kafkaDf.selectExpr("CAST(value AS STRING)")

  // movie struct
  val struct = new StructType()
    .add("title", DataTypes.StringType)
    .add("year", DataTypes.IntegerType)
    .add("cast", ArrayType(DataTypes.StringType))
    .add("genres", ArrayType(DataTypes.StringType))

  val moviesNestedDf = moviesJsonDf.select(from_json($"value", struct).as("movie"))
  // json flatten
  val movieFlattenedDf = moviesNestedDf.selectExpr("movie.title", "movie.year", "movie.cast","movie.genres")


  // convert to parquet and save to hdfs
  val query = movieFlattenedDf
    .writeStream
    .outputMode("append")
    .format("parquet")
    .queryName("movies")
    .option("checkpointLocation", "src/main/resources/chkpoint_dir")
    .start("src/main/resources/output")
    .awaitTermination()
  }

上下文:

  • 我直接从 intellij 运行它(使用本地 spark 已安装)
  • 我设法毫无问题地从 kafka 读取并写入 控制台(使用控制台模式)
  • 暂时我想写文件 在本地机器上(但我确实尝试过 HDFS 集群,问题是 一样)

我的问题:

在工作期间,它没有在文件夹中写入任何内容,我必须手动停止工作才能最终看到文件。

我想这可能与.awaitTermination() 有关 有关信息,我尝试删除此选项,但如果没有删除,我会收到错误消息,并且作业根本无法运行。

也许我没有设置正确的选项,但在阅读了很多次文档并在 Google 上搜索后,我没有找到任何东西。

你能帮帮我吗?

谢谢

编辑:

  • 我使用的是 spark 2.4.0
  • 我尝试了 64/128mb 格式 => 在我停止工作之前没有任何变化没有文件

【问题讨论】:

  • 你的 spark 版本是什么?
  • 可以尝试写入 64MB/128MB 的数据吗?例如,您可能面临与此 one 类似的问题,或者将 hdfs.rollSize 设置为较小的值。我相信 Spark 中可能也有类似的配置
  • 我用 spark 版本更新帖子
  • Hello @Yrah 目前有多少数据写入输出流?将 rollSize 规范设置为 64MB,您必须写出至少 64MB 才能使写出过程生效。如果您不想等待应用程序写入 64MB,可以将其设置为一个较小的数字
  • 您可以尝试设置 parquet.block.size 选项,如下所示:block_sz = 1024 //1KB val query = movieFlattenedDf .writeStream .outputMode("append") .format("parquet") .queryName("movies") .option("parquet.block.size", block_sz) .option("checkpointLocation", "src/main/resources/chkpoint_dir") .start("src/main/resources/output") .awaitTermination()

标签: scala apache-spark apache-kafka parquet spark-structured-streaming


【解决方案1】:

是的问题解决

我的问题是,我的数据太少,spark 正在等待更多数据来写入 parquet 文件。

为了完成这项工作,我使用了来自@AlexandrosBiratsis 的评论 (改变块大小)

再次感谢@AlexandrosBiratsis 非常感谢

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-07-10
    • 2020-09-08
    • 2018-07-02
    • 2022-01-14
    • 2012-04-07
    相关资源
    最近更新 更多