【发布时间】:2019-04-07 22:47:03
【问题描述】:
我必须生成 3000000 个文件作为 spark 作业的输出。
我有两个输入文件:
File 1 -> Size=3.3 Compressed, No.Of Records=13979835
File 2 -> Size=1.g Compressed, No.Of Records=6170229
Spark Job 正在执行以下操作:
- 读取此文件并根据公共列 1 加入它们。 -> 数据帧-A
- DataFrame-A 基于一列的分组结果2 -> DataFrame-B
- 从 DataFrame-B 使用 array_join 作为聚合列,并用 '\n' 字符分隔该列。 -> DataFrame-C
-
DataFrame-C 按 column2 分区写入结果。
val DF1 = sparkSession.read.json("FILE1") // |ID |isHighway|isRamp|pvId |linkIdx|ffs |length | val DF12 = sparkSession.read.json("FILE2") // |lId |pid | val joinExpression = DF1.col("pvId") === DF2.col("lId") val DFA = DF.join(tpLinkDF, joinExpression, "inner").select(col("ID").as("SCAR"), col("lId"), col("length"), col("ffs"), col("ar"), col("pid")).orderBy("linkIdx") val DFB = DFA.select(col("SCAR"),concat_ws(",", col("lId"), col("length"),col("ffs"), col("ar"), col("pid")).as("links")).groupBy("SCAR").agg(collect_list("links").as("links")) val DFC = DFB.select(col("SCAR"), array_join(col("links"), "\n").as("links")) DFC.write.format("com.databricks.spark.csv").option("quote", "\u0000").partitionBy("SCAR").mode(SaveMode.Append).format("csv").save("/tmp")
我必须生成 3000000 个文件作为 spark 作业的输出。
【问题讨论】:
-
为什么需要这个?小文件问题。
-
这是一种要求,其他系统需要读取这些小文件(不是全部,但随着对该文件的请求到达,每个文件都在文件名中包含一些 id,因此使用该 ID 请求该文件接收系统必须读取该文件)并在 45 秒内实时给出结果。
-
老实说听起来像是一场灾难。
标签: scala apache-spark