PySpark SQL Filter与Join异常:执行后返回空DataFrame问题咨询
解决Spark DataFrame Join后为空的问题
我来帮你分析下为什么第一次运行代码会得到空的DataFrame,以及怎么修复这个问题。
问题根源
你第一次的代码里,HiveContext的初始化方式不对!在Spark中,HiveContext必须依赖一个SparkContext实例才能正确初始化,直接hc = HiveContext()是不完整的,这会导致后续创建的DataFrame在执行join操作时无法正确解析数据关联条件,最终返回空结果。
另外还有个潜在小问题:两个DataFrame有完全同名的列(比如id1、id2、id3),虽然这不是导致空结果的直接原因,但后续处理容易引发歧义,最好在join时给DataFrame指定别名来区分。
修正后的代码
下面是能得到预期结果的修正版本:
from pyspark.sql import functions as F from pyspark.sql import Row, HiveContext from pyspark import SparkContext # 先初始化或获取已有的SparkContext实例 sc = SparkContext.getOrCreate() # 用SparkContext实例正确初始化HiveContext hc = HiveContext(sc) rows1 = [Row(id1 = '2', id2 = '1', id3 = 'a'), Row(id1 = '3', id2 = '2', id3 = 'a'), Row(id1 = '4', id2 = '3', id3 = 'b')] df1 = hc.createDataFrame(rows1) df2 = df1.filter(F.col("id3")=="a") # 给DataFrame指定别名,明确关联条件的列归属 df3 = df1.alias("df1").join(df2.alias("df2"), F.col("df1.id2") == F.col("df2.id1"), "inner") # 查看结果 df3.show()
预期输出
执行修正后的代码后,你会得到以下符合预期的结果:
+---+---+---+---+---+---+ |id1|id2|id3|id1|id2|id3| +---+---+---+---+---+---+ | 3| 2| a| 2| 1| a| | 4| 3| b| 3| 2| a| +---+---+---+---+---+---+
额外优化建议
如果你使用的是Spark 2.0及以上版本,更推荐用SparkSession替代HiveContext,它是Spark统一的API入口,用法更简洁且功能更全面:
from pyspark.sql import SparkSession, functions as F, Row # 初始化支持Hive的SparkSession spark = SparkSession.builder.appName("JoinExample").enableHiveSupport().getOrCreate() rows1 = [Row(id1 = '2', id2 = '1', id3 = 'a'), Row(id1 = '3', id2 = '2', id3 = 'a'), Row(id1 = '4', id2 = '3', id3 = 'b')] df1 = spark.createDataFrame(rows1) df2 = df1.filter(F.col("id3")=="a") df3 = df1.alias("df1").join(df2.alias("df2"), F.col("df1.id2") == F.col("df2.id1"), "inner") df3.show()
内容的提问来源于stack exchange,提问作者WEIHANG LIU
相关产品推荐
相关产品推荐

