代码之家  ›  专栏  ›  技术社区  ›  TMichel

pyspark-读取拼花后优化分区数

  •  1
  • TMichel  · 技术社区  · 8 年前

    在由 year month spark.default.parallelism 设置为 4 ,假设我要创建一个数据框架,由2017年的11~12个月和2018年的1~3个月两个来源组成 A B .

    df = spark.read.parquet(
        "A.parquet/_YEAR={2017}/_MONTH={11,12}",
        "A.parquet/_YEAR={2018}/_MONTH={1,2,3}",
        "B.parquet/_YEAR={2017}/_MONTH={11,12}",
        "B.parquet/_YEAR={2018}/_MONTH={1,2,3}",
    )
    

    如果我得到了分区的数目,spark使用了 spark.default.parallelism并行性 默认情况下:

    df.rdd.getNumPartitions()
    Out[4]: 4
    

    考虑到 df 我需要表演 join groupBy 每个周期的操作,并且这些数据或多或少均匀地分布在每个周期上(每个周期大约1000万行):

    问题

    • 重新分区是否会提高后续操作的性能?
    • 如果是,如果我有10个不同的时段(a和b都是每年5个),我是否应该按时段数重新分区并显式地引用要重新分区的列( df.repartition(10,'_MONTH','_YEAR') )?
    1 回复  |  直到 8 年前
        1
  •  2
  •   Alper t. Turker    8 年前

    重新分区是否会提高后续操作的性能?

    通常不会。抢先重新分区数据的唯一原因是避免在相同的情况下进一步洗牌。 Dataset 用于基于相同条件的多个联接

    如果是这样,如果我有10个不同的时段(a和b都是每年5个),我是否应该按时段数重新分区,并显式地引用要重新分区的列(df.repartition(10,''u month',''u year')?

    我们一步一步来:

    • 我应该按句号重新划分吗

      实践者不能保证级别和分区之间的1:1关系,所以唯一要记住的是,不能有比唯一键更多的非空分区,所以使用大得多的值是没有意义的。

    • 并显式引用要重新分区的列

      如果你 repartition 后来 join groupBy 对这两部分使用相同的列集是唯一合理的解决方案。

    总结

    repartitoning 在两种情况下,连接前是有意义的:

    • 如果有多个后续 joins

      df_ = df.repartition(10, "foo", "bar")
      df_.join(df1, ["foo", "bar"])
      ...
      df_.join(df2, ["foo", "bar"])
      
    • 当所需的 输出 分区不同于 spark.sql.shuffle.partitions (而且没有广播连接)

      spark.conf.get("spark.sql.shuffle.partitions")
      # 200
      spark.conf.set("spark.sql.autoBroadcastJoinThreshold", -1)
      
      df1_ = df1.repartition(11, "foo", "bar")
      df2_ = df2.repartition(11, "foo", "bar")
      
      df1_.join(df2_, ["foo", "bar"]).rdd.getNumPartitions()
      # 11
      
      df1.join(df2, ["foo", "bar"]).rdd.getNumPartitions()
      # 200
      

      这可能比:

      spark.conf.set("spark.sql.shuffle.partitions", 11)
      df1.join(df2, ["foo", "bar"]).rdd.getNumPartitions()
      spark.conf.set("spark.sql.shuffle.partitions", 200)