PySpark动态生成列的聚合问题:分组合并行与动态列聚合
解决PySpark动态列聚合的UNRESOLVED_COLUMN错误
问题原因
你的错误在于尝试将字符串拼接的SQL表达式传给col(),但PySpark的agg()方法需要接收PySpark Column对象,而非原生字符串表达式。这种写法会导致Spark无法解析你构造的字符串,从而抛出UNRESOLVED_COLUMN_WITH_SUGGESTION错误。
正确实现方案
针对你的需求:
bookID列聚合为列表:使用collect_list即可,这部分你的写法是正确的。- 动态列聚合规则(组内任意行值为1则设为1):由于动态列的值只有0或1,直接用
max()聚合即可——组内只要有一行值为1,最大值就是1;全为0则最大值为0,完美匹配需求。
完整代码示例
from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, lit, collect_list, max # 初始化SparkSession spark = SparkSession.builder.appName("BookAggregation").getOrCreate() # 构造原始数据 data = [ ("firstname1", "lastname1", "address1", -97), ("firstname1", "lastname1", "address1", -23), ("firstname2", "lastname2", "address2", -23), ("firstname2", "lastname2", "address2", -76) ] columns = ["name", "last", "address", "bookID"] df = spark.createDataFrame(data, schema=columns) # 动态添加列 books = [-23, -44, -97, -32, -57, -76] for book in books: df = df.withColumn(str(book), when(df['bookID'] == book, lit(1)).otherwise(0)) # 构造聚合表达式 agg_exps = [collect_list("bookID").alias("bookID_values")] for book in books: col_name = str(book) # 用max聚合实现"任意行有1则为1"的逻辑 agg_exps.append(max(col(col_name)).alias(col_name)) # 执行分组聚合 agg_df = df.groupBy("name", "last", "address").agg(*agg_exps) # 查看结果 agg_df.show()
替代方案(用sum+when实现)
如果你更倾向于用sum判断的方式,也可以这样写:
agg_exps = [collect_list("bookID").alias("bookID_values")] for book in books: col_name = str(book) agg_exps.append( when(sum(col(col_name)) > 0, lit(1)).otherwise(lit(0)).alias(col_name) )
输出结果
执行后会得到你期望的输出:
+----------+---------+--------+-------------+---+---+---+---+---+---+ | name| last| address|bookID_values|-23|-44|-97|-32|-57|-76| +----------+---------+--------+-------------+---+---+---+---+---+---+ |firstname1|lastname1|address1| [-97, -23]| 1| 0| 1| 0| 0| 0| |firstname2|lastname2|address2| [-23, -76]| 1| 0| 0| 0| 0| 1| +----------+---------+--------+-------------+---+---+---+---+---+---+
内容的提问来源于stack exchange,提问作者Mojdeh Ebrahimi
相关产品推荐
相关产品推荐

