将转换步骤组织为函数
将笔记本拆分为提取、转换和加载函数,让每个函数都接收并返回一个 DataFrame,以便轻松测试。
将转换步骤组织为函数 是 CoddyKit 上的免费 Pandas & NumPy Academy 课时。 这是第 1 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 Pandas & NumPy Academy 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 Pandas & NumPy Academy 课程共包含 4 节课。
超越笔记本
Jupyter 笔记本非常适合探索,但不适合生产数据流程。代码分散在 50 个单元格中,使用全局状态且没有测试,这种结构十分脆弱:修改一个单元格可能会悄悄破坏另一个单元格。专业的做法是将每个转换提取为一个命名函数,接收一个 DataFrame 并返回一个 DataFrame。这种关注点分离是构建可维护数据工程代码的基础。
# Notebook-style (fragile)
df = pd.read_csv('orders.csv')
df = df.dropna(subset=['revenue'])
df = df[df['quantity'] > 0]
df['revenue_per_unit'] = df['revenue'] / df['quantity']
# Function-style (robust)
def extract(path):
return pd.read_csv(path)
def transform(df):
return (
df.dropna(subset=['revenue'])
.query('quantity > 0')
.assign(revenue_per_unit=lambda d: d['revenue'] / d['quantity'])
)
df = transform(extract('orders.csv'))ETL 模式:提取、转换、加载
ETL 模式将数据流程分为三个阶段:提取(从源读取)、转换(清洗并丰富数据)和加载(写入目标位置)。每个阶段都是一个独立函数。这种分离使您可以轻松更换数据源(CSV 或数据库)、修改清洗逻辑,或更改输出格式,而无需改动另外两个阶段。每个生产流程都应遵循这一结构。
def extract(config):
return pd.read_csv(config['input_path'], parse_dates=['order_date'])
def transform(df, config):
return (
df
.dropna(subset=config['required_cols'])
.query('quantity > 0')
.assign(revenue=lambda d: d['quantity'] * d['unit_price'])
)
def load(df, config):
df.to_parquet(config['output_path'], index=False)
print(f'Saved {len(df)} rows.')单一职责原则
每个转换函数都应该恰好完成一件事。名为 clean_data() 的函数如果同时删除空值、限制离群值、解析日期并编码类别,就很难测试和调试。相反,请将 drop_nulls()、cap_outliers()、parse_dates() 和 encode_categories() 分别编写为独立函数。这种粒度使您可以轻松跳过、替换或重新排列任意单个步骤。
def drop_null_rows(df, required_cols):
return df.dropna(subset=required_cols)
def remove_returns(df):
return df[df['quantity'] > 0]
def compute_revenue(df):
return df.assign(revenue=lambda d: d['quantity'] * d['unit_price'])
def add_date_features(df):
return df.assign(
year=lambda d: d['order_date'].dt.year,
month=lambda d: d['order_date'].dt.month
)使用 pipe() 串联步骤
使用 pipe() 连接遵循单一职责原则的函数,将完整的转换阶段构建为易读的链。这个链读起来就像一份操作配方:每一行代表一个步骤,数据沿箭头从上到下流动。您可以注释掉或重新排序任意步骤,而无需重命名变量。最终结果是一个干净且经过丰富处理的 DataFrame,已准备好进入加载阶段。
import pandas as pd
REQUIRED = ['order_id', 'revenue', 'order_date']
def transform(raw_df):
return (
raw_df
.pipe(drop_null_rows, required_cols=REQUIRED)
.pipe(remove_returns)
.pipe(compute_revenue)
.pipe(add_date_features)
)
df_clean = transform(pd.read_csv('orders.csv', parse_dates=['order_date']))
print(df_clean.shape)返回行数以便审计
每个转换函数都应能够选择性地记录接收到的行数和返回的行数。在函数体前后分别记录行数,是一种轻量级审计追踪方式,可以让您快速了解一次流程运行中行数据在哪个位置被删除。将这些计数存入列表,并在每次运行结束时将其打印为报告。
audit_log = []
def audited(func):
def wrapper(df, *args, **kwargs):
before = len(df)
result = func(df, *args, **kwargs)
after = len(result)
audit_log.append({'step': func.__name__, 'in': before, 'out': after, 'dropped': before - after})
return result
return wrapper
@audited
def drop_null_rows(df, required_cols):
return df.dropna(subset=required_cols)使用小型 DataFrame 编写可测试函数
命名函数最大的优势是可测试性。编写一个代表真实边界情况的小型测试 DataFrame,并断言函数的输出。使用一行包含空值的数据测试 drop_null_rows 函数,验证该行会被删除;再使用一行不含空值的数据进行测试,验证该行会被保留。针对转换函数的单元测试可以在流程代码发生变化时及时发现回归问题。
import pandas as pd
def test_drop_null_rows():
test_df = pd.DataFrame({
'order_id': [1, 2, 3],
'revenue': [100.0, None, 200.0]
})
result = drop_null_rows(test_df, required_cols=['revenue'])
assert len(result) == 2, 'Should have 2 non-null rows'
assert result['revenue'].isna().sum() == 0, 'No nulls in revenue'
print('test_drop_null_rows PASSED')
test_drop_null_rows()将函数组织到模块中
随着流程不断扩大,请将函数拆分到 Python 模块文件中:extract.py、transform.py、load.py 和 validate.py。主脚本 pipeline.py 负责导入这些模块并协调它们。这种文件结构清晰易读,可以使用 pytest 进行测试,也可以作为 Python 软件包部署。它与使用 dbt 或 Airflow 等工具的数据工程团队所采用的标准布局一致。
# pipeline.py
# from extract import extract_orders
# from transform import transform
# from load import load_to_parquet
# from validate import validate_sales_df
# def run_pipeline(config):
# raw = extract_orders(config)
# clean = transform(raw, config)
# validate_sales_df(clean)
# load_to_parquet(clean, config)
print('Module-based pipeline structure shown above (imports commented for demo)')幂等性:安全地多次运行
设计良好的流程函数具有幂等性:对同一输入运行两次会产生相同的输出,并且不会造成副作用。避免原地修改(df.drop(..., inplace=True)),并且对于会修改列的函数,始终在开头使用 df.copy()。幂等函数可以在失败后重新运行,而不会破坏输出数据。
def add_revenue_flag(df, threshold=500):
# Use copy to avoid mutating the input
df = df.copy()
df['is_large_order'] = df['revenue'] >= threshold
return df
# Running twice gives the same result
df1 = add_revenue_flag(df_clean)
df2 = add_revenue_flag(df_clean)
assert df1.equals(df2), 'Function is not idempotent!'
print('Idempotency check passed.')使用类型提示实现自文档化
在转换函数签名中添加 Python 类型提示,明确函数契约:def transform(df: pd.DataFrame) -> pd.DataFrame。类型提示具有自文档作用——阅读函数的开发者无需查看函数体,就能准确了解它需要什么以及返回什么。类型提示还支持 mypy 等静态分析工具和 IDE 自动补全功能,在运行前发现类型错误。
from typing import List
def drop_null_rows(df: pd.DataFrame, required_cols: List[str]) -> pd.DataFrame:
return df.dropna(subset=required_cols)
def compute_revenue(df: pd.DataFrame) -> pd.DataFrame:
return df.assign(revenue=lambda d: d['quantity'] * d['unit_price'])
print('Type-annotated functions ready for production use.')使用文档字符串记录每个步骤
每个转换函数都应包含一行文档字符串,说明它的作用、所需列,以及添加或删除的列。良好的文档字符串使函数可以通过 help() 和 IDE 工具提示被发现。它们也可以作为规范文档:如果失败的断言与文档字符串矛盾,那么这是一个错误;如果文档字符串与代码不一致,那么这是文档错误。
def compute_revenue(df: pd.DataFrame) -> pd.DataFrame:
"""Add 'revenue' column as quantity * unit_price.
Requires: 'quantity' (numeric), 'unit_price' (numeric) columns.
Returns: df with new 'revenue' float column appended.
"""
return df.assign(revenue=lambda d: d['quantity'] * d['unit_price'])
help(compute_revenue)运行完整流程
依次调用三个阶段来协调完整的 ETL:提取、转换和加载。在每个阶段周围添加计时,以衡量时间花费在哪里。分别捕获每个阶段的异常,使错误消息能够指出失败的阶段。记录完整运行的开始和结束时间,以及运行是否成功,以便进行监控。
import time
CONFIG = {
'input_path': 'orders.csv',
'output_path': 'orders_clean.parquet',
'required_cols': ['order_id', 'revenue']
}
t0 = time.time()
raw = extract(CONFIG)
clean = transform(raw, CONFIG)
load(clean, CONFIG)
print(f'Pipeline completed in {time.time()-t0:.1f}s')
print(f'Audit log: {audit_log}')快速检查
测试您对本课数据分析概念的理解。
课程回顾
在本课中,您学习了:将流程构建为遵循单一职责的 ETL 函数、使用类型提示和文档字符串使函数可测试、幂等且能够自我说明,以及通过计时和审计日志记录协调完整流程。接下来,我们将学习如何使用配置字典为流程设置参数,以便在不同数据集之间复用。
常见问题解答
「将转换步骤组织为函数」课时是免费的吗?
是的 — 「将转换步骤组织为函数」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 Pandas & NumPy Academy 课程的其余内容,请升级到 CoddyKit PRO。 Pandas & NumPy Academy 课程共包含 4 节课。
「将转换步骤组织为函数」这节课中我会学到什么?
将笔记本拆分为提取、转换和加载函数,让每个函数都接收并返回一个 DataFrame,以便轻松测试。 你通过在浏览器中直接运行的动手代码来练习 Pandas & NumPy Academy,全天候 AI 导师会在你学习这节课的过程中回答你的问题。
学习 Pandas & NumPy Academy 需要有经验吗?
无需任何先前经验。CoddyKit 上的 Pandas & NumPy Academy 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 1 节课,共 4 节。
「将转换步骤组织为函数」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 Pandas & NumPy Academy 课中编写并运行代码吗?
能。每节 Pandas & NumPy Academy 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。