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

使用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.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 06:39:02