PySpark中如何指定带列参数的REBALANCE分区提示
如何使用PySpark API指定携带列名的REBALANCE分区提示
问题复现
首先创建测试用DataFrame:
df = spark.range(10)
以下几种传入列名的写法都会执行失败:
- 直接传入字符串形式的列名
>>> df.hint("rebalance", "id").explain() ... pyspark.sql.utils.AnalysisException: REBALANCE Hint parameter should include columns, but id found
- 传入带表别名的限定列名字符串
>>> df.alias("df").hint("rebalance", "df.id").explain() ... pyspark.sql.utils.AnalysisException: REBALANCE Hint parameter should include columns, but df.id found
- 传入
F.col()生成的列引用对象,会直接触发Python侧的参数类型校验错误
>>> import pyspark.sql.functions as F >>> df.hint("rebalance", F.col("id")).explain() TypeError: all parameters should be in (<class 'str'>, <class 'list'>, <class 'float'>, <class 'int'>), got Column<'id'> of type <class 'pyspark.sql.column.Column'>
注意:不带任何列参数的rebalance提示可以正常运行,但这种是无分区键的轮询重平衡,会使用REBALANCE_PARTITIONS_BY_NONE模式,不符合按指定列重平衡的需求,执行计划如下:
>>> df.hint("rebalance").explain() == Physical Plan == AdaptiveSparkPlan isFinalPlan=false +- Exchange RoundRobinPartitioning(4), REBALANCE_PARTITIONS_BY_NONE, [id=#551] +- Range (0, 10, step=1, splits=4)
解决方案
该问题和使用的Spark版本强相关,不同版本对应不同的正确写法:
Spark 3.3及以上版本
Spark 3.3修复了DataFrame.hint对REBALANCE hint的参数解析bug,直接传入列名字符串即可生效:
df.hint("rebalance", "id").explain()
执行后物理计划会显示REBALANCE_PARTITIONS_BY_COL,证明已按指定列做哈希重平衡:
== Physical Plan == AdaptiveSparkPlan isFinalPlan=false +- Exchange hashpartitioning(id#0L, 4), REBALANCE_PARTITIONS_BY_COL, [id=#89] +- Range (0, 10, step=1, splits=4)
多列重平衡直接依次传入多个列名即可:
df.hint("rebalance", "col1", "col2").explain()
如果需要指定初始分区数,在列名前传入整数类型的分区参数即可:
# 按id列重平衡,初始分区数设置为8 df.hint("rebalance", 8, "id").explain()
Spark 3.2.x版本(首个支持REBALANCE hint的版本)
该版本存在两个问题:一是JVM侧hint解析器会把直接传入的字符串列名识别为字面量而非列引用,二是Python侧的hint方法做了严格的参数类型校验,禁止传入Column对象,可通过两种方式实现需求:
- 用
col()表达式包裹列名传入
传入字符串形式的列表达式,强制解析器将其识别为列引用:df.hint("rebalance", "col(id)").explain() - 直接使用Spark SQL语法
注册临时视图后通过SQL语句编写hint,不存在API层的解析问题:df.createOrReplaceTempView("test_tmp") spark.sql("SELECT /*+ REBALANCE(id) */ * FROM test_tmp").explain()
内容的提问来源于stack exchange,提问作者Martin Studer
相关产品推荐
相关产品推荐

