You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark基于其他DataFrame列条件拆分数据子集的问题

解决Spark DataFrame按其他表列拆分的问题

问题根源说明

  1. filter+isin报错原因:Spark的isin()方法无法直接接收另一个DataFrame的列作为参数,它需要的是本地可迭代集合(如列表、元组),或者使用子查询语法,直接传入DataFrame列会导致解析失败。
  2. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.08 22:45:38