PySpark基于其他DataFrame列条件拆分数据子集的问题
解决Spark DataFrame按其他表列拆分的问题
问题根源说明
filter+isin报错原因:Spark的isin()方法无法直接接收另一个DataFrame的列作为参数,它需要的是本地可迭代集合(如列表、元组),或者使用子查询语法,直接传入DataFrame列会导致解析失败。merge报错原因:merge是Pandas DataFrame的方法,Spark DataFrame没有这个属性,Spark中实现数据关联要使用join方法。
解决方案一:提取ID本地集合后筛选
先将Movies和Series中的ID提取为本地列表,再用isin()过滤Ratings:
# 提取电影ID列表 movie_ids = [row.ID for row in Movies.collect()] # 筛选电影评分数据 Ratings_movies = Ratings.filter(Ratings.ID.isin(movie_ids)) # 提取剧集ID列表 series_ids = [row.ID for row in Series.collect()] # 筛选剧集评分数据 Ratings_series = Ratings.filter(Ratings.ID.isin(series_ids)) # 查看结果 Ratings_movies.show() Ratings_series.show()
注意:如果Movies/Series数据量极大,
collect()会将分布式数据拉取到Driver节点,可能引发内存溢出,这种场景更适合用方案二。
解决方案二:用Spark Join实现分布式筛选
通过内关联匹配ID,再保留Ratings的目标列,全程在分布式环境执行:
# 关联Movies得到电影评分数据 Ratings_movies = Ratings.join(Movies, on="ID", how="inner").select(Ratings["ID"], Ratings["Rating"]) # 关联Series得到剧集评分数据 Ratings_series = Ratings.join(Series, on="ID", how="inner").select(Ratings["ID"], Ratings["Rating"]) # 查看结果 Ratings_movies.show() Ratings_series.show()
最终结果验证
执行后会得到你期望的两个子集:
- Ratings_movies:
| ID | Rating | | -- | ------ | | 2 | 9 | | 4 | 10 | | 5 | 2 |
- Ratings_series:
| ID | Rating | | -- | ------ | | 1 | 7 | | 3 | 5 | | 6 | 9 |
内容的提问来源于stack exchange,提问作者AWDn0n
相关产品推荐
相关产品推荐

