PySpark代码报错'Column is not iterable':DataFrame对比问题排查
问题:PySpark对比DataFrame时出现"Column is not iterable"错误
我尝试对两个PySpark DataFrame进行带计算的对比,输出包含"Type of Change"和"Number of Occurences"两列的表格。为此合并了两个DataFrame,给列名添加"_New"和"_Old"后缀以区分,还编写了一个函数来对比列、统计次数并输出DataFrame,但运行代码时返回错误"Column is not iterable"。
可复现代码如下:
import pyspark.sql.functions as f from pyspark.sql import SparkSession from pyspark.sql.functions import col, when from pyspark.sql.functions import sum as fsum from pyspark.sql.types import StructType, StructField data1 = [ ("Postcode1", "100", "150"), ("Postcode2", "200", "250"), ("Postcode3", "300", "350"), ("Postcode4", "400", "450"), ] data2 = [ ("Postcode1", "150", "150"), ("Postcode2", "200", "200"), ("Postcode3", "350", "350"), ("Postcode4", "400", "450"), ] Columns = ["Postcode", "Count1", "Count2"] rdd1 = spark.sparkContext.parallelize(data1) rdd2 = spark.sparkContext.parallelize(data2) df1 = spark.createDataFrame(rdd1, schema=Columns) df2 = spark.createDataFrame(rdd2, schema=Columns) ColumnNames_New = [ f.col("Postcode"), f.col("Count1").alias("Count1_New"), f.col("Count2").alias("Count2_New"), ] df1_NewData = df1.select(ColumnNames_New) ColumnNames_Old = [ f.col("Postcode"), f.col("Count1").alias("Count1_Old"), f.col("Count2").alias("Count2_Old"), ] df2_OldData = df2.select(ColumnNames_Old) # Joining two dataframes on postcode Comparison_DF = df1_NewData.join(df2_OldData, "Postcode") Columns_to_Compare = [f.col("Count1"), f.col("Count2")] # Forming blank dataframe for change report ChangeReport_RDD = spark.sparkContext.emptyRDD() ChangeReport_Columns = [ StructField("Type of Change", StringType(), True), StructField("Number of Occurences", IntegerType(), True), ] ChangeReport_DF = spark.createDataFrame([], schema=StructType(ChangeReport_Columns)) def Form_Change_Report( ColumnList, Comparison_Dataframe: df1_NewData, ChangeReport_Dataframe: ChangeReport_DF, ): for x in ColumnList: Comparison_Dataframe = Comparison_Dataframe.withColumn( x + "_Change", when(col(x + "_Old") == col(x + "_New"), 0).otherwise(1) ) Change_DF = Comparison_Dataframe.select( lit("Change to " + x).alias("Type of Change"), fsum(x + "_Change").alias("Number of Occurences"), ) ChangeReport_Dataframe = ChangeReport_Dataframe.unionByName(Change_DF) return ChangeReport_Dataframe ChangeReport_DF = Form_Change_Report(Columns_to_Compare, Comparison_DF, ChangeReport_DF)
错误原因
- 列名传入错误:
Columns_to_Compare传入的是Column对象列表,而非字符串列名。循环中x + "_Old"会尝试拼接Column对象和字符串,触发"Column is not iterable"错误。 - 类型标注错误:函数参数的类型标注用了DataFrame实例(
df1_NewData、ChangeReport_DF),这不符合语法规范,应使用pyspark.sql.DataFrame或直接移除错误标注。 - 函数未导入:
lit函数未导入,会引发NameError。 - 统计方式错误:使用
select搭配sum函数无法正确聚合数据,需改用agg方法。
修正后的代码
import pyspark.sql.functions as f from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, lit from pyspark.sql.functions import sum as fsum from pyspark.sql.types import StructType, StructField, StringType, IntegerType data1 = [ ("Postcode1", "100", "150"), ("Postcode2", "200", "250"), ("Postcode3", "300", "350"), ("Postcode4", "400", "450"), ] data2 = [ ("Postcode1", "150", "150"), ("Postcode2", "200", "200"), ("Postcode3", "350", "350"), ("Postcode4", "400", "450"), ] Columns = ["Postcode", "Count1", "Count2"] # 初始化SparkSession(Databricks中可省略,本地运行需添加) spark = SparkSession.builder.appName("ChangeReport").getOrCreate() rdd1 = spark.sparkContext.parallelize(data1) rdd2 = spark.sparkContext.parallelize(data2) df1 = spark.createDataFrame(rdd1, schema=Columns) df2 = spark.createDataFrame(rdd2, schema=Columns) ColumnNames_New = [ f.col("Postcode"), f.col("Count1").alias("Count1_New"), f.col("Count2").alias("Count2_New"), ] df1_NewData = df1.select(ColumnNames_New) ColumnNames_Old = [ f.col("Postcode"), f.col("Count1").alias("Count1_Old"), f.col("Count2").alias("Count2_Old"), ] df2_OldData = df2.select(ColumnNames_Old) # 关联两个DataFrame Comparison_DF = df1_NewData.join(df2_OldData, "Postcode") # 改为传入字符串列名列表 Columns_to_Compare = ["Count1", "Count2"] # 创建空的变更报告DataFrame ChangeReport_Columns = [ StructField("Type of Change", StringType(), True), StructField("Number of Occurences", IntegerType(), True), ] ChangeReport_DF = spark.createDataFrame([], schema=StructType(ChangeReport_Columns)) def Form_Change_Report(ColumnList, Comparison_Dataframe, ChangeReport_Dataframe): for x in ColumnList: # 计算变更标记列 Comparison_Dataframe = Comparison_Dataframe.withColumn( f"{x}_Change", when(col(f"{x}_Old") == col(f"{x}_New"), 0).otherwise(1) ) # 统计变更次数(改用agg方法聚合) Change_DF = Comparison_Dataframe.agg( lit(f"Change to {x}").alias("Type of Change"), fsum(f"{x}_Change").alias("Number of Occurences"), ) # 合并到报告DataFrame ChangeReport_Dataframe = ChangeReport_Dataframe.unionByName(Change_DF) return ChangeReport_Dataframe # 生成变更报告 ChangeReport_DF = Form_Change_Report(Columns_to_Compare, Comparison_DF, ChangeReport_DF) # 查看结果 ChangeReport_DF.show()
输出结果
+-----------------+-------------------+ | Type of Change|Number of Occurences| +-----------------+-------------------+ |Change to Count1| 2| |Change to Count2| 1| +-----------------+-------------------+
内容的提问来源于stack exchange,提问作者Danylo Kuznetsov
相关产品推荐
相关产品推荐

