【问题标题】:Why does Spark write always the same amount of files to HDFS?为什么 Spark 总是向 HDFS 写入相同数量的文件?
【发布时间】:2019-02-22 03:26:28
【问题描述】:

我有一个用 Scala 编写的 Spark 流应用程序,在 CDH 中运行。应用程序从 Kafka 读取数据并将数据写入 HDFS。在向HDFS写入数据之前,我执行了partitionBy,所以数据是分区写入的。每个分区在写入时获得 3 个文件。我还使用coalesce 来控制我的数据的分区数。我的期望是coalesce 命令设置的分区数将设置HDFS 中输出目录中的文件数,但是尽管coalesce 命令设置了分区数,但文件数始终为3。我尝试使用 3 个执行器和 6 个执行器运行,但每个分区中的文件数仍然是 3。

这就是我将数据写入 HDFS 的方式:

//Some code
val ssc = new StreamingContext(sc, Seconds(1))
val stream = KafkaUtils.createDirectStream[String, String](
             ssc,
             PreferConsistent,
             Subscribe[String,String](topics, kafkaParams))
val sparkExecutorsCount = sc.getConf.getInt("spark.executor.instances", 1)
//Some code
stream.foreachRDD { rdd =>
    if(!rdd.isEmpty()) {
        val offsetRanges = rdd.asInstanceOf[HasOffsetRanges].offsetRanges
        val data = rdd.map(kafkaData => (getKey(kafkaData.value()), kafkaData.value()))
        val columns = Array("key", "value")
        data.toDF(columns: _*).coalesce(sparkExecutorsCount)
            .write.mode(SaveMode.Append)
            .partitionBy("key").text(MY_PATH)

       stream.asInstanceOf[CanCommitOffsets].commitAsync(offsetRanges)
    } else {
        //handle empty RDD
    }
}

请告知如何让我的 Spark 应用程序将其他数量的文件写入输出目录。谢谢

【问题讨论】:

  • 为3,MY_PATH下的总文件数或文件夹数。你能发布 hdfs dfs -ls -R MY_PATH 吗?
  • @alexeipab 3 是MY_PATH下每个子目录(分区)的文件数。它看起来像这样:MY_PATH/key=1 MY_PATH/key=1/file1.txt MY_PATH/key=1/file2.txt MY_PATH/key=1/file3.txt MY_PATH/key=2 MY_PATH/key=2/file1.txt MY_PATH/key=2/file2.txt MY_PATH/key=3/file3.txt
  • 几个附加问题:1) 当您生成输出时 sparkExecutorsCount == 3 还是 1? 2)每个文件的时间戳是多少,所有文件都一样吗?
  • 1) 我尝试使用 sparkExecutorsCount == 3 和 sparkExecutorsCount == 6。输出是相同的。我从未尝试过 sparkExecutorsCount == 1。2)是的,所有文件都有相同的时间戳

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


【解决方案1】:

coalesce 不会重新洗牌键上的数据,它连接分区而不跨分区重新分配记录。在您的示例中, partitionBy 不是在 Dataframe 上调用,而是在由 .write 函数返回的 DataFrameWriter 上调用。在这种情况下,key 列看起来有 3 个值,因此 3 个文件夹(key=1,key=2,key=3)和每个文件夹中具有相同时间戳的 3 个文件可以解释为Dataframe 至少有 3 个分区,因为每个分区都会有一个写入器运行,它必须输出到 3 个文件夹(key=1,key=2,key=3)。我怀疑“sparkExecutorsCount == 6”没有影响可能是因为 Kafka 只为您提供了 3 个分区,在这种情况下合并没有影响。

要在每个密钥文件夹中只保存 1 个文件,您可以尝试 coalesce(1) 或使用 repartition($"key") 代替它并保留现有的partitionBy

data.toDF(columns: _*).repartition($"key")
        .write.mode(SaveMode.Append)
        .partitionBy("key").text(MY_PATH)

或

data.toDF(columns: _*).repartition(sparkExecutorsCount, $"key")
        .write.mode(SaveMode.Append)
        .partitionBy("key").text(MY_PATH)

【讨论】:

  • 我提供的文件结构只是部分的:因为“key”有 1000 个值,我没有全部指定。抱歉,如果我的文件结构打印输出误导了您。我不知道 DataFrame 是否有 3 个分区。我该如何检查和控制这个? Kafka 为我提供了与 sparkExecutorsCount 相同数量的分区。我在创建 Kafka 主题时手动控制它。
猜你喜欢
  • 1970-01-01
  • 2020-12-03
  • 1970-01-01
  • 2016-01-30
  • 1970-01-01
  • 2020-07-10
  • 2016-01-11
  • 1970-01-01
  • 2020-09-15
相关资源
最近更新 更多