PySpark DataFrame:如何执行双查询及选取指定列值?
解决PySpark选取指定列+多条件过滤,以及多查询执行的问题
一、正确实现“选取指定两列且满足双条件”的方法
你之前写的df.where((df['E'] ==0 ).where(df['C']=='non'))其实能得到符合条件的行,但写法不够直观,而且没有指定要选取的列。要同时满足条件并提取指定列,有几种更清晰的实现方式:
方式1:用逻辑与(&)连接条件,配合select选列
注意PySpark中逻辑与要用&,而且每个条件必须加括号(因为运算符优先级问题),然后用select()指定要保留的列:
# 先过滤符合条件的行,再选取E和C列 result_df = df.where((df['E'] == 0) & (df['C'] == 'non')).select('E', 'C') # 或者用filter,和where功能完全等价 result_df = df.filter((df['E'] == 0) & (df['C'] == 'non')).select('E', 'C')
方式2:链式调用where(更易读)
连续调用where()就相当于逻辑与,这种写法不用处理运算符优先级,可读性更强:
result_df = df.where(df['E'] == 0).where(df['C'] == 'non').select('E', 'C')
方式3:SQL风格查询(如果习惯SQL语法)
先把DataFrame注册成临时视图,再用SQL语句查询:
# 注册临时视图(视图名可以自定义) df.createOrReplaceTempView("my_data") # 执行SQL查询,直接指定列和条件 result_df = spark.sql("SELECT E, C FROM my_data WHERE E = 0 AND C = 'non'")
二、在DataFrame中执行两个查询的方法
这里分两种常见场景来解释:
场景1:执行两个独立的查询(分别获取不同结果)
PySpark的DataFrame是惰性求值的,你可以直接定义多个基于源DataFrame的新DataFrame,直到调用show()、count()这类action操作时才会真正计算。比如:
# 查询1:获取E=0且C='non'的E、C列 query1 = df.where((df['E'] == 0) & (df['C'] == 'non')).select('E', 'C') # 查询2:获取E>0且C='yes'的E、C列 query2 = df.where((df['E'] > 0) & (df['C'] == 'yes')).select('E', 'C') # 触发计算并查看结果 query1.show() query2.show()
场景2:执行两个查询并合并结果
如果需要把两个查询的结果合并成一个DataFrame,可以用union()(要求两个查询的列结构完全一致):
# 两个查询的列要保持一致 query1 = df.where((df['E'] == 0) & (df['C'] == 'non')).select('E', 'C') query2 = df.where((df['E'] == 1) & (df['C'] == 'yes')).select('E', 'C') # 合并两个查询结果 combined_result = query1.union(query2) combined_result.show()
如果两个查询的列结构不同,或者需要关联结果,可以用join()来实现,比如根据某列关联两个查询的结果。
内容的提问来源于stack exchange,提问作者sara jones
相关产品推荐
相关产品推荐

