如何使用PySpark读取MongoDB中的多个Collections?
读取MongoDB多集合并执行Join与聚合操作
读取多个Collections
复用已创建的my_spark实例即可读取其他集合,无需重复初始化SparkSession,以下是两种常用方式:
方法1:通过uri指定目标集合
# 读取collA df_a = my_spark.read.format("mongo") \ .option("uri", "mongodb://127.0.0.1/mydb.collA") \ .load() # 读取collB df_b = my_spark.read.format("mongo") \ .option("uri", "mongodb://127.0.0.1/mydb.collB") \ .load() # 读取collC df_c = my_spark.read.format("mongo") \ .option("uri", "mongodb://127.0.0.1/mydb.collC") \ .load()
方法2:分开指定数据库与集合参数
# 读取collA df_a = my_spark.read.format("mongo") \ .option("spark.mongodb.input.database", "mydb") \ .option("spark.mongodb.input.collection", "collA") \ .load() # 读取collB df_b = my_spark.read.format("mongo") \ .option("spark.mongodb.input.database", "mydb") \ .option("spark.mongodb.input.collection", "collB") \ .load() # 读取collC df_c = my_spark.read.format("mongo") \ .option("spark.mongodb.input.database", "mydb") \ .option("spark.mongodb.input.collection", "collC") \ .load()
执行Join操作与聚合计算
假设三个集合通过user_id字段关联,先完成多表关联,再基于关联结果执行聚合:
# 关联df_a与df_b joined_ab = df_a.join(df_b, on="user_id", how="inner") # 继续关联df_c得到全量关联数据 joined_all = joined_ab.join(df_c, on="user_id", how="inner") # 示例聚合:统计每个用户的累计订单数(假设df_b含order_count字段) agg_result = joined_all.groupBy("user_id") \ .sum("order_count") \ .withColumnRenamed("sum(order_count)", "total_orders") # 查看聚合结果 agg_result.show()
关键注意点
- 确认关联字段在各集合中存在且数据类型一致,避免关联失败
- 根据业务场景选择合适的Join类型(inner/left/right/full)
- 聚合前可先过滤无效数据,提升计算效率
内容的提问来源于stack exchange,提问作者yuser099881232
相关产品推荐
相关产品推荐

