【问题标题】:Configuration for spark job to write 3000000 file as output配置 Spark 作业以将 3000000 文件作为输出写入
【发布时间】: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. 读取此文件并根据公共列 1 加入它们。 -> 数据帧-A
  2. DataFrame-A 基于一列的分组结果2 -> DataFrame-B
  3. 从 DataFrame-B 使用 array_join 作为聚合列,并用 '\n' 字符分隔该列。 -> DataFrame-C
  4. 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


【解决方案1】:

在运行了一些测试后,我想到了批量运行这项工作,例如:

  • 查询 startIdx: 0, endIndex:100000
  • 查询 startIdx: 100000, endIndex:200000
  • 查询 startIdx: 200000, endIndex:300000

等等....直到

  • 查询 startIdx: 2900000, endIndex:3000000

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-07-22
    • 1970-01-01
    • 2023-03-14
    相关资源
    最近更新 更多