如何用PySpark将预聚合Pandas DataFrame转为嵌套结构并写入MongoDB
问题描述
我有一个已完成预聚合的Pandas DataFrame,结构如下:
| InterfaceName | StartDate | StartHour | DocumentCount | TotalRowCount |
|---|---|---|---|---|
| Interface_A | 2023-04-01 | 0 | 5 | 4384 |
| Interface_A | 2023-04-01 | 1 | 58 | 57168 |
| Interface_B | 2023-04-01 | 1 | 1 | 136 |
| Interface_C | 2023-04-01 | 1 | 1 | 131 |
| Interface_A | 2023-04-02 | 0 | 58 | 57168 |
| Interface_B | 2023-04-02 | 0 | 1 | 131 |
| Interface_C | 2023-04-02 | 0 | 1 | 136 |
| Interface_A | 2023-04-02 | 1 | 2 | 1657 |
| Interface_B | 2023-04-02 | 1 | 2 | 1539 |
| Interface_C | 2023-04-02 | 1 | 2 | 1657 |
需要用PySpark将其转换为如下指定的嵌套Schema,然后写入MongoDB的结构化集合:
root |-- StartDate: date (nullable = true) |-- StartHour: integer (nullable = true) | |-- InterfaceSummary: struct (nullable = false) | | |-- InterfaceName: string (nullable = true) | | |-- DocumentCount: string (nullable = true) | | |-- TotalRowCount: string (nullable = true)
解决方案
1. 将Pandas DataFrame转为PySpark DataFrame
先把现有Pandas DataFrame导入PySpark环境,同步字段类型:
from pyspark.sql import SparkSession from pyspark.sql.types import DateType, IntegerType, StringType # 初始化SparkSession spark = SparkSession.builder \ .appName("ConvertToNestedSchema") \ .getOrCreate() # 假设你的Pandas DataFrame名为pd_df spark_df = spark.createDataFrame(pd_df) # 调整字段类型以匹配目标Schema spark_df = spark_df.withColumn("StartDate", spark_df["StartDate"].cast(DateType())) \ .withColumn("StartHour", spark_df["StartHour"].cast(IntegerType())) \ .withColumn("DocumentCount", spark_df["DocumentCount"].cast(StringType())) \ .withColumn("TotalRowCount", spark_df["TotalRowCount"].cast(StringType()))
2. 构造嵌套结构体并分组
按StartDate和StartHour分组,把每个接口的信息打包成InterfaceSummary结构体,再收集同组的所有结构体:
from pyspark.sql.functions import struct, collect_list # 构造InterfaceSummary结构体 struct_df = spark_df.withColumn("InterfaceSummary", struct( "InterfaceName", "DocumentCount", "TotalRowCount" )) # 分组并收集同日期小时下的所有接口摘要 nested_df = struct_df.groupBy("StartDate", "StartHour") \ .agg(collect_list("InterfaceSummary").alias("InterfaceSummary"))
注:从示例数据看,同一日期小时下存在多个接口,所以用collect_list将结构体收集为数组更贴合实际场景,若需求为单个结构体需确认数据逻辑是否允许。
3. 验证Schema结构
打印转换后的Schema确认是否符合要求:
nested_df.printSchema()
预期输出:
root |-- StartDate: date (nullable = true) |-- StartHour: integer (nullable = true) |-- InterfaceSummary: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- InterfaceName: string (nullable = true) | | |-- DocumentCount: string (nullable = true) | | |-- TotalRowCount: string (nullable = true)
4. 写入MongoDB
配置MongoDB连接参数,将转换后的DataFrame写入指定集合:
nested_df.write.format("mongo") \ .option("uri", "mongodb://<username>:<password>@<host>:<port>/<database>.<collection>") \ .mode("append") # 可选模式:append、overwrite、ignore等 .save()
替换<username>、<password>、<host>、<port>、<database>、<collection>为实际MongoDB连接信息。
若Spark环境未安装MongoDB连接器,启动Spark时需添加依赖(版本需匹配Spark和MongoDB版本):
--packages org.mongodb.spark:mongo-spark-connector_2.12:10.1.1
内容的提问来源于stack exchange,提问作者Ben Halicki
相关产品推荐
相关产品推荐

