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

如何基于PyMongoLoader实现Pandas直接读取并规范化MongoDB嵌套JSON

如何让PyMongoLoader配合pandas实现嵌套JSON的直接规范化?

背景与当前实现

我正在为pymongo编写包装类PyMongoLoader,核心目标是让它兼容pandas,实现如下调用:

from loader import PyMongoLoader
import pandas as pd

_loader = PyMongoLoader(url="my_url", port="my_port")
df = pd.read_json(_loader, orient="records")

根据pandas文档,pd.read_json接受任何带有read方法的对象作为第一个参数。我在内部通过调用collection.find获取数据,并用bson.dumps将结果解析为字符串,当前的read方法实现如下:

def read(self, size: int = -1) -> str:
    assert self._is_valid, "MongoDB connection must be valid for retrieving data"

    self.logger.debug(
        f"Retrieving documents from {self._collection}. Expected size: \033[94m{self._collection.count_documents({})}\033[0m"
    )
    result = list(self._collection.find())

    self.logger.debug(
        f"Retrieved {len(result)} documents from remote. Size of json: {sum([len(document) for document in result])}"
    )
    return dumps(result)

遇到的问题

目前实现可正常运行,但数据库中存在类似“26344T Control Measure [5.3-10.8]”的复杂字段名,导致数据以嵌套JSON形式存储。我希望规范化这些嵌套字段,但pd.json_normalize仅接受字符串作为参数,不支持路径参数。

我想到两种解决方案,但都不理想:

  • 修改数据库字段名:能解决问题,但每次新增数据都要注意,非常麻烦。
  • 在Loader源码中补丁pd.read_json或pd.json_normalize:也能解决问题,但会破坏其他使用pandas的代码,不是合理方案。

我的问题是:是否有官方支持的直接从类文件对象规范化JSON的方式?如果没有,如何规范化传给pd.read_json的JSON,消除数据库中的嵌套问题?

补充说明

从数据库获取的JSON格式示例:

{
  "_id": {
    "$oid": "651ec788c110096a55c8d4de"
  },
  "DateTime": "20/02/2017 15:00:00",
  "546321B": {
    "measure": 0
  },
  "538612B": {
    "measure": 80
  },
  "517713B": {
    "measure": 70
  },
  "508021V": {
    "avg": 37
  }
}

期望得到的DataFrame格式:

DateTime546321B.measure538612B.measure517713B.measure508021V.avg
"20/02/2017 15:00:00"0807037

理想情况下,希望直接通过pd.read_json(loader, orient="records")得到上述结果。

解决方案

pandas目前没有官方支持直接从类文件对象(如你的PyMongoLoader)读取时自动规范化嵌套JSON的功能,最合理的方案是在PyMongoLoader的read方法中先扁平化嵌套数据,再转为JSON字符串输出。

具体实现步骤

  1. 编写一个辅助函数,将嵌套的文档结构扁平化为点分隔键的格式:
def flatten_document(doc):
    flattened = {}
    for key, value in doc.items():
        # 根据需求决定是否保留_id字段,这里选择忽略
        if key == "_id":
            continue
        # 处理嵌套字典
        if isinstance(value, dict):
            for sub_key, sub_value in value.items():
                flattened[f"{key}.{sub_key}"] = sub_value
        else:
            flattened[key] = value
    return flattened
  1. 修改read方法,在获取MongoDB查询结果后,先扁平化每个文档:
def read(self, size: int = -1) -> str:
    assert self._is_valid, "MongoDB connection must be valid for retrieving data"

    self.logger.debug(
        f"Retrieving documents from {self._collection}. Expected size: \033[94m{self._collection.count_documents({})}\033[0m"
    )
    result = list(self._collection.find())
    
    # 扁平化所有文档
    flattened_results = [flatten_document(doc) for doc in result]

    self.logger.debug(
        f"Retrieved {len(flattened_results)} flattened documents from remote."
    )
    return dumps(flattened_results)

这样修改后,pd.read_json(_loader, orient="records")会直接加载到扁平化后的DataFrame,完全符合你的期望。

内容的提问来源于stack exchange,提问作者Pollastre

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 02:17:04