【发布时间】: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