【问题标题】:Spark Difference in repartition and spark.sql.shuffle.partitionSpark 重新分区和 spark.sql.shuffle.partition 的区别
【发布时间】:2019-08-26 20:48:21
【问题描述】:

我正在使用 --conf spark.sql.shuffle.partitions=100 运行一个 spark 程序

在应用程序中我有以下内容

Dataset<Row> df_partitioned = df.repartition(df.col("enriched_usr_id"));
df_partitioned = df_partitioned.sortWithinPartitions(df_partitioned.col("transaction_ts"));
df_partitioned.mapPartitions(
    SparkFunctionImpl.mapExecuteUserLogic(), Encoders.bean(Transformed.class));

我有大约 500 万用户,我想为每个用户排序数据并为每个用户执行一些逻辑。

我的问题是,这会将数据划分为 500 万个分区还是 100 个分区,以及每个用户如何执行。

【问题讨论】:

  • spark.sql.shuffle.partitions 用于决定涉及洗牌时的分区数量,即在连接期间等。

标签: java apache-spark dataframe apache-spark-sql


【解决方案1】:

df.repartition(df.col("enriched_usr_id")) 将使用enriched_usr_id 将数据划分为100 个分区(spark.sql.shuffle.partitions),这意味着多个用户将在同一个分区中。

【讨论】:

    猜你喜欢
    • 2019-03-02
    • 1970-01-01
    • 1970-01-01
    • 2020-10-11
    • 2019-11-13
    • 2022-08-03
    • 2017-04-05
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多