如何在PySpark中结合withColumn与for循环处理带编号列的DataFrame
问题解决:PySpark批量生成新列(withColumn循环报错修复)
错误原因
你使用withColumn(*[...])的方式有误,因为withColumn一次只能添加单个列,它仅接受两个参数:新列名、列计算表达式。而你通过*解包列表中的元组,会一次性传入多个参数(比如2组列的话会传4个参数),超出方法的参数预期,因此抛出TypeError: col should be Column。
正确实现方式
方式1:循环迭代调用withColumn
通过循环依次为每一组days/code列生成对应的new列,每次调用withColumn添加一列,逐步更新DataFrame:
from pyspark.sql import functions as F # 示例DataFrame df = spark.createDataFrame( [(100, 100, 'A', 'A'), (1000, 200, 'A', 'A'), (1000, 300, 'B', 'A'), (1000, 1000, 'B', 'B')], "days1 int, days2 int, code1 string, code2 string") # 初始化临时DataFrame,用于逐步添加新列 temp_df = df # 循环处理1-12列,这里示例用range(2)对应2列,实际改为range(12)即可 for i in range(2): n = i + 1 temp_df = temp_df.withColumn( f'new{n}', F.when(F.col(f'days{n}') > 100, F.col(f'code{n}')).otherwise('X') ) # 查看结果 temp_df.select('days1', 'code1', 'new1', 'days2', 'code2', 'new2').show()
运行结果:
+-----+-----+----+-----+-----+----+ |days1|code1|new1|days2|code2|new2| +-----+-----+----+-----+-----+----+ | 100| A| X| 100| A| X| | 1000| A| A| 200| A| A| | 1000| B| B| 300| A| A| | 1000| B| B| 1000| B| B| +-----+-----+----+-----+-----+----+
方式2:用select批量生成新列
如果更偏好一次性生成所有列,可以结合原列和新列表达式,通过select完成:
from pyspark.sql import functions as F df = spark.createDataFrame( [(100, 100, 'A', 'A'), (1000, 200, 'A', 'A'), (1000, 300, 'B', 'A'), (1000, 1000, 'B', 'B')], "days1 int, days2 int, code1 string, code2 string") # 获取所有原列 original_cols = df.columns # 批量生成新列的表达式 new_cols = [ F.when(F.col(f'days{i+1}') > 100, F.col(f'code{i+1}')).otherwise('X').alias(f'new{i+1}') for i in range(2) ] # 合并原列与新列,执行select result_df = df.select(original_cols + new_cols) result_df.show()
内容的提问来源于stack exchange,提问作者Chuck
相关产品推荐
相关产品推荐

