使用Python Great Expectations移除Pandas DataFrame无效数据并存入PostgreSQL
使用Great Expectations过滤Pandas无效数据并写入PostgreSQL实现方案
可行性说明
Great Expectations本身的核心能力是数据校验,没有直接提供删除原DataFrame行的内置方法,但可以通过获取校验结果中的无效行索引,自行实现数据过滤和分流写入,完全可以满足你的需求。
具体实现步骤
步骤1:依赖安装
先确保你已经安装了需要的库:
pip install great-expectations pandas sqlalchemy psycopg2-binary
步骤2:单校验规则实现逻辑
- 构造Great Expectations的Pandas数据源,执行校验规则
- 从校验结果中提取无效行的索引
- 拆分出有效数据集和无效数据集
- 分别写入目标位置:有效数据保留做后续处理,无效数据写入PostgreSQL错误表
示例代码如下:
import pandas as pd import great_expectations as ge from sqlalchemy import create_engine # 1. 初始化PostgreSQL连接,替换为你自己的数据库配置 pg_engine = create_engine('postgresql://用户名:密码@数据库地址:端口/数据库名') # 2. 加载你的原始DataFrame,这里用示例数据演示 raw_df = pd.DataFrame({ 'id': [1,2,3,4,5,6,7,8,9,10], 'age': [25, None, 30, None, 40, 18, None, 22, 35, None] }) # 3. 转换为Great Expectations的Pandas数据集 ge_df = ge.dataset.PandasDataset(raw_df) # 4. 执行校验规则,这里以age非空为例 validation_result = ge_df.expect_column_values_to_not_be_null('age') # 5. 提取无效行的索引 invalid_indexes = validation_result['result']['unexpected_index_list'] # 6. 拆分有效数据和无效数据 valid_df = raw_df.drop(invalid_indexes).reset_index(drop=True) invalid_df = raw_df.loc[invalid_indexes].reset_index(drop=True) # 7. 将无效数据写入PostgreSQL的错误表,if_exists参数根据你的需求调整为append/replace/fail invalid_df.to_sql( name='你的错误表名', con=pg_engine, if_exists='append', index=False ) # 后续可直接使用valid_df进行业务处理
步骤3:多校验规则扩展方案
如果你需要同时执行多个校验规则,只需要在每次执行校验后收集所有无效行索引,最后统一去重即可,示例:
# 存储所有无效行索引的集合 all_invalid_indexes = set() # 规则1:age非空 res1 = ge_df.expect_column_values_to_not_be_null('age') all_invalid_indexes.update(res1['result']['unexpected_index_list']) # 规则2:age在18到60之间 res2 = ge_df.expect_column_values_to_be_between('age', min_value=18, max_value=60) all_invalid_indexes.update(res2['result']['unexpected_index_list']) # 规则3:id唯一 res3 = ge_df.expect_column_values_to_be_unique('id') all_invalid_indexes.update(res3['result']['unexpected_index_list']) # 后续拆分逻辑和上述单规则逻辑一致 valid_df = raw_df.drop(list(all_invalid_indexes)).reset_index(drop=True) invalid_df = raw_df.loc[list(all_invalid_indexes)].reset_index(drop=True)
注意事项
- 如果你需要记录每行具体违反了什么校验规则,可以在收集无效索引的时候同步关联规则信息,写入错误表的时候额外增加违规规则字段即可
- 写入PostgreSQL的时候建议提前建好错误表的结构,避免
to_sql自动建表的字段类型不符合预期 - 对于大体积DataFrame,建议分批次执行校验和写入,避免内存溢出
内容的提问来源于stack exchange,提问作者Florin P.
相关产品推荐
相关产品推荐

