PySpark中如何用Join替代Isin实现高效条件值填充?
问题:PySpark中通过Join实现多列匹配的条件填充
我需要在PySpark DataFrame中,对满足「某几列的值存在于另一个DataFrame对应列中」条件的行填充新列值,但由于使用*.collect().distinct()和.isin()*的耗时远高于join,因此希望改用join或broadcast实现条件填充。
Pandas中的实现方式
在pandas中我会这样实现:
df.loc[(df.A.isin(df2.A)) | (df.B.isin(df2.B)), 'new_column'] = 'new_value'
我的尝试与问题
我尝试了以下PySpark方法,但通过.count()发现行数异常减少,未得到正确结果:
count_first = df.count() dfA_1 = df.join(df2, 'A', 'leftanti') \ .withColumn('new_column', F.lit(None).cast(StringType())) dfA_2= df.join(df2, 'A', 'inner') \ .withColumn('new_column', F.lit('new_value')) df = dfA_1.unionByName(dfA_2) count_second = df.count() count_first - count_second
请问如何在PySpark中用join正确实现该需求?
内容的提问来源于stack exchange,提问作者euh
相关产品推荐
相关产品推荐

