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

如何用PySpark将预聚合Pandas DataFrame转为嵌套结构并写入MongoDB

问题描述

我有一个已完成预聚合的Pandas DataFrame,结构如下:

InterfaceNameStartDateStartHourDocumentCountTotalRowCount
Interface_A2023-04-01054384
Interface_A2023-04-0115857168
Interface_B2023-04-0111136
Interface_C2023-04-0111131
Interface_A2023-04-0205857168
Interface_B2023-04-0201131
Interface_C2023-04-0201136
Interface_A2023-04-02121657
Interface_B2023-04-02121539
Interface_C2023-04-02121657

需要用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 00:42:09