PySpark如何explode展开user_ids数组并与user_id列合并为单一列
错误原因
你遇到的类型不匹配报错不是因为withColumn无法引用新创建的列,而是你在coalesce中用df["user_ids"]引用了原始DataFrame中未爆炸的array<bigint>类型的user_ids列,和user_id的bigint类型不匹配导致的。
正确实现代码
直接通过F.col()引用当前处理链路中已爆炸的列即可:
import pyspark.sql.functions as F result_df = df \ # 展开user_ids数组,保留空数组/Null的行 .withColumn("user_ids", F.explode_outer("user_ids")) \ # 优先取原有user_id,为空则取爆炸后的值 .withColumn("user_id", F.coalesce(F.col("user_id"), F.col("user_ids"))) \ # 仅保留最终需要的user_id列 .select("user_id")
如果怕列名混淆,也可以给爆炸后的结果单独起别名:
import pyspark.sql.functions as F result_df = df \ .select("user_id", F.explode_outer("user_ids").alias("exploded_uid")) \ .withColumn("user_id", F.coalesce(F.col("user_id"), F.col("exploded_uid"))) \ .select("user_id")
以上代码输出结果和你预期的完全一致。
内容的提问来源于stack exchange,提问作者qalis
相关产品推荐
相关产品推荐

