PySpark统计两个不同表DataFrame的CountDistinct值报错如何解决
问题原因
pyspark.sql.functions.countDistinct() 的入参要求为列名字符串、列对象,或是多个列的组合,不支持直接传入DataFrame类型。你代码中ppp.select(["T1_"+c for c in impacted_columns.key1.split("-")])的返回结果是DataFrame,直接作为参数传入countDistinct()就触发了类型不匹配错误。
解决方法
你需要统计多列联合的去重行数,有两种常用写法,功能完全等价:
写法1:使用distinct() + count()(更直观)
直接对选中的key列去重后统计行数:
# 统计ppp的key去重数量 pppc = ppp.select(["T1_"+c for c in impacted_columns.key1.split("-")]).distinct().count() # 统计plu的key去重数量 pluc = plu.select(["T2_" + c for c in impacted_columns.key2.split("-")]).distinct().count()
写法2:使用countDistinct聚合
将多个key列解包后传入countDistinct(),再通过DataFrame的agg方法执行聚合操作,最后提取计数值:
# 统计ppp的key去重数量 key_cols1 = [F.col("T1_"+c) for c in impacted_columns.key1.split("-")] pppc = ppp.agg(F.countDistinct(*key_cols1).alias("distinct_cnt")).first()[0] # 统计plu的key去重数量 key_cols2 = [F.col("T2_"+c) for c in impacted_columns.key2.split("-")] pluc = plu.agg(F.countDistinct(*key_cols2).alias("distinct_cnt")).first()[0]
内容的提问来源于stack exchange,提问作者wmbff
相关产品推荐
相关产品推荐

