如何基于PyMongoLoader实现Pandas直接读取并规范化MongoDB嵌套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格式:
| DateTime | 546321B.measure | 538612B.measure | 517713B.measure | 508021V.avg |
|---|---|---|---|---|
| "20/02/2017 15:00:00" | 0 | 80 | 70 | 37 |
理想情况下,希望直接通过pd.read_json(loader, orient="records")得到上述结果。
解决方案
pandas目前没有官方支持直接从类文件对象(如你的PyMongoLoader)读取时自动规范化嵌套JSON的功能,最合理的方案是在PyMongoLoader的read方法中先扁平化嵌套数据,再转为JSON字符串输出。
具体实现步骤
- 编写一个辅助函数,将嵌套的文档结构扁平化为点分隔键的格式:
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
- 修改
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

