PySpark代码与Hive SQL同数据查询行数差异排查求助
问题排查与修复方案
行数差异核心原因
1. isin方法参数误用
PySpark中isin("new")会将字符串拆分为单个字符的迭代序列(即['n','e','w']),实际判断逻辑变为**uid_type是否等于'n'/'e'/'w'**,而非Hive SQL中要求的等于"new"。同理isin("old")会匹配uid_type为'o'/'l'/'d'的行,导致tbl2、tbl3的匹配行数远超预期,直接引发第一次JOIN后数据大规模膨胀。
2. 链式左连接的基数倍增
即使修复isin参数,如果tbl2/tbl3中存在多条匹配同一条tbl1行的记录,PySpark的连续左连接会基于第一次JOIN后的膨胀数据集执行第二次JOIN,最终行数变为N个tbl2匹配 × M个tbl3匹配的乘积。而Hive结果行数与原始tbl1一致,说明业务数据中每个(direct_uid, item_id)对应最多一个uid_type='new'和uid_type='old'的记录,但Spark未显式处理多匹配场景时会产生笛卡尔积。
修复方案
方案一:提前去重(通用安全版)
先对tbl2、tbl3按(direct_uid, item_id)分组去重,确保每个键只保留一个related_uid,再与原始tbl1左连接,彻底避免基数膨胀:
from pyspark.sql import functions as F # 预处理tbl2:取每个(direct_uid, item_id)下第一个uid_type='new'的related_uid tbl2_df = ( input_table_df .filter(F.col("uid_type") == "new") .groupBy("direct_uid", "item_id") .agg(F.first("related_uid").alias("tbl2_related_uid")) ) # 预处理tbl3:取每个(direct_uid, item_id)下第一个uid_type='old'的related_uid tbl3_df = ( input_table_df .filter(F.col("uid_type") == "old") .groupBy("direct_uid", "item_id") .agg(F.first("related_uid").alias("tbl3_related_uid")) ) # 与原始tbl1左连接 self_joined_tbl_df = ( input_table_df.alias("tbl1") .join(tbl2_df.alias("tbl2"), on=["direct_uid", "item_id"], how="left") .join(tbl3_df.alias("tbl3"), on=["direct_uid", "item_id"], how="left") .select( F.col("tbl1.item_id"), F.col("tbl1.item_type"), F.col("tbl1.item_name"), F.col("tbl1.p_date"), F.col("tbl1.uid"), F.col("tbl1.direct_uid"), F.col("tbl1.related_uid"), F.col("tbl1.uid_type"), F.coalesce(F.col("tbl2.tbl2_related_uid"), F.col("tbl3.tbl3_related_uid")).alias("new_related_uid") ) )
方案二:仅修复isin参数(适用于数据唯一场景)
如果业务数据保证每个(direct_uid, item_id, uid_type)是唯一的,仅需将isin的参数改为列表即可:
from pyspark.sql import functions as F self_joined_tbl_df = ( input_table_df.alias("tbl1") .join( input_table_df.alias("tbl2"), on=( (F.col("tbl1.direct_uid") == F.col("tbl2.direct_uid")) & (F.col("tbl1.item_id") == F.col("tbl2.item_id")) & (F.col("tbl2.uid_type").isin(["new"])) # 改为列表参数 ), how="left" ) .join( input_table_df.alias("tbl3"), on=( (F.col("tbl1.direct_uid") == F.col("tbl3.direct_uid")) & (F.col("tbl1.item_id") == F.col("tbl3.item_id")) & (F.col("tbl3.uid_type").isin(["old"])) # 改为列表参数 ), how="left" ) .select( F.col("tbl1.item_id"), F.col("tbl1.item_type"), F.col("tbl1.item_name"), F.col("tbl1.p_date"), F.col("tbl1.uid"), F.col("tbl1.direct_uid"), F.col("tbl1.related_uid"), F.col("tbl1.uid_type"), F.coalesce(F.col("tbl2.related_uid"), F.col("tbl3.related_uid")).alias("new_related_uid") ) )
内容的提问来源于stack exchange,提问作者WarBoy
相关产品推荐
相关产品推荐

