Kedro 数据集与数据目录(kedro.io)完整指南:从 AbstractDataset 到 DataCatalog 的源码级解析
【免费下载链接】kedroKedro is a toolbox for production-ready data science. It uses software engineering best practices to help you create data engineering and data science pipelines that are reproducible, maintainable, and modular.项目地址: https://gitcode.com/GitHub_Trending/ke/kedro
Kedro 的kedro.io模块是整套数据管道的"数据接入层",它通过统一的AbstractDataset抽象与DataCatalog注册中心,把读写 CSV、Parquet、数据库乃至云存储的底层细节封装成一致的load/save接口。本文以 Kedro 官方 API 文档 kedro.io 模块总览 为主体骨架,结合本仓库源码(kedro/io/目录)与测试用例(tests/io/),系统讲解数据集基类、版本化机制、内存缓存、数据集工厂模式、凭据解析以及多进程共享内存目录的实现原理与实战用法。读完本文,你将能够独立编写自定义数据集、配置 catalog.yml、利用数据集工厂减少重复配置,并理解ParallelRunner等并行场景下数据目录的设计约束。
kedro.io 模块概览
kedro.io是 Kedro 用于读写各类数据集的核心包,模块入口文件 kedro/io/init.py 的模块注释明确指出:"kedro.ioprovides functionality to read and write to a number of datasets. At the core of the library is theAbstractDatasetclass."(kedro.io 提供读写多种数据集的功能,库的核心是AbstractDataset类)。
根据官方 API 文档 kedro.io.md 中的成员清单,该模块对外暴露的公共构件如下表所示:
| 名称 | 类型 | 说明 |
|---|---|---|
| kedro.io.AbstractDataset | Class | 所有 Kedro 数据集的基类 |
| kedro.io.AbstractVersionedDataset | Class | 版本化数据集的基类 |
| kedro.io.CachedDataset | Class | 在内存中缓存数据的数据集包装器 |
| kedro.io.DataCatalog | Class | 管理 Kedro 管道中使用的数据集 |
| kedro.io.CatalogProtocol | Class | 定义在目录中管理数据集的通用接口 |
| kedro.io.SharedMemoryDataCatalog | Class | 在共享内存上下文中管理数据集 |
| kedro.io.SharedMemoryCatalogProtocol | Class | 扩展CatalogProtocol以支持共享内存场景 |
| kedro.io.CatalogConfigResolver | Class | 基于数据集工厂模式与凭据解析数据集配置 |
| kedro.io.MemoryDataset | Class | 在内存中存储数据的数据集 |
| kedro.io.Version | Class | 表示数据集版本信息 |
| kedro.io.DatasetAlreadyExistsError | Exception | 数据集已存在时抛出 |
| kedro.io.DatasetError | Exception | 通用数据集错误 |
| kedro.io.DatasetNotFoundError | Exception | 数据集未找到时抛出 |
此外,__init__.py中还导出了SharedMemoryDataset(共享内存数据集,供SharedMemoryDataCatalog内部使用)。这些类共同构成了 Kedro 数据管理的完整体系:AbstractDataset定义协议,具体数据集(如kedro_datasets中的CSVDataset、ParquetDataset)实现协议,DataCatalog负责按名称注册与调度,CatalogConfigResolver负责从配置解析出数据集实例。
AbstractDataset:所有数据集的基类
AbstractDataset定义在 kedro/io/core.py 中,是所有 Kedro 数据集的抽象基类,同时是泛型类AbstractDataset[_DI, _DO]:_DI是save时输入数据的类型,_DO是load时返回数据的类型。
必须实现的方法
任何自定义数据集都必须实现以下三个抽象方法:
load() -> _DO:从数据源读取数据并返回;save(data: _DI) -> None:将数据写入数据源;_describe() -> dict:返回数据集的关键属性字典,用于日志与repr展示(不含None值)。
同时有两个可选方法:
_exists() -> bool:判断目标数据是否已存在。基类默认实现只打印警告并返回False,exists()方法会包装它并捕获异常转换为DatasetError;_release():释放缓存数据,基类默认为空操作(pass),release()方法负责包装。
自动装饰机制:init_subclass
AbstractDataset通过__init_subclass__(core.py 中第 422 行起)在子类定义时自动完成三件事:
- 捕获初始化参数:包装子类的
__init__,通过getcallargs记录传入的实参到self._init_args,供后续_init_config()与to_config()使用; - 兼容
_load/_save命名:如果子类实现的是_load/_save而非load/save,会自动别名到load/save; - 自动包裹错误处理:用
_load_wrapper/_save_wrapper装饰load/save,统一捕获底层异常并包装成带上下文的DatasetError。
其中_save_wrapper还会拦截data is None的情况,直接抛出DatasetError("Saving 'None' to a 'Dataset' is not allowed"),防止把None写入存储。这意味着子类只需实现_load/_save即可获得统一的日志与异常包装,测试用例tests/io/test_core.py中的MyDataset正是用_load/_save实现的典型写法。
持久性与 _EPHEMERAL 标志
AbstractDataset默认设置_EPHEMERAL = False,表示数据集是持久化的;而MemoryDataset、SharedMemoryDataset、CachedDataset等纯内存实现会把_EPHEMERAL设为True。DataCatalog.to_config()在序列化目录时会跳过内存型数据集,正是依据这一标志(见 data_catalog.py 中_is_memory_dataset的调用处)。
from_config 工厂方法与自定义数据集示例
from_config(core.py 第 260 行起)是解析 catalog.yml 配置并实例化数据集的统一入口。它先调用parse_dataset_definition解析出数据集类与配置字典,再用class_obj(**config)实例化;若实例化失败,会给出包含数据集名与类型全限定名的DatasetError。
官方文档给出的自定义数据集示例(同时出现在 core.py 的 docstring 中):
from pathlib import Path, PurePosixPath import pandas as pd from kedro.io import AbstractDataset class MyOwnDataset(AbstractDataset[pd.DataFrame, pd.DataFrame]): def __init__(self, filepath, param1, param2=True): self._filepath = PurePosixPath(filepath) self._param1 = param1 self._param2 = param2 def load(self) -> pd.DataFrame: return pd.read_csv(self._filepath) def save(self, df: pd.DataFrame) -> None: df.to_csv(str(self._filepath)) def _exists(self) -> bool: return Path(self._filepath.as_posix()).exists() def _describe(self): return dict(param1=self._param1, param2=self._param2)对应的 catalog.yml 配置:
my_dataset: type: <path-to-my-own-dataset>.MyOwnDataset filepath: data/01_raw/my_data.csv param1: <param1-value> # param1 is a required argument # param2 will be True by defaultparse_dataset_definition:类型解析的底层逻辑
parse_dataset_definition(core.py 第 621 行起)承担配置解析的核心职责,其关键行为包括:
- 配置必须包含
type键,值为类的全限定名(如kedro_datasets.pandas.CSVDataset)或类对象; - 解析时按
_DEFAULT_PACKAGES = ["kedro.io.", "kedro_datasets.", ""]三个前缀依次尝试导入,因此 catalog.yml 中既可以写全限定名,也可以只写pandas.CSVDataset这样的短名; - 若最终找不到类,会抛出
DatasetError,并附带提示"请确认是否安装了 kedro-datasets(pip install kedro-datasets)"; - 解析出的类必须是
AbstractDataset的子类,否则报DatasetError; - 配置中的
versioned: true标志会触发版本注入(见下文版本化章节); validator键属于 catalog 级配置,会被移除并警告,不会传给数据集构造函数;- 出于向后兼容,遇到以
Set结尾的类型名(旧拼写)会发出警告:"Since kedro-datasets 2.0, 'Dataset' is spelled with a lowercase 's'"。
版本化数据集:AbstractVersionedDataset 与 Version
Version:加载版本与保存版本的载体
Version定义在 core.py 第 601 行,是一个namedtuple("Version", ["load", "save"]):
- 若
Version.load为None,则加载最新可用版本; - 若
Version.save为None,保存版本将自动生成,格式为YYYY-MM-DDThh.mm.ss.sssZ(UTC 时间戳)。
时间戳由generate_timestamp()(core.py 第 590 行)生成,格式常量VERSION_FORMAT = "%Y-%m-%dT%H.%M.%S.%fZ",微秒部分会被截去。
AbstractVersionedDataset:版本化数据集基类
AbstractVersionedDataset继承自AbstractDataset,是所有支持版本化数据集的基类。其构造函数签名(core.py 第 833 行):
def __init__( self, filepath: PurePosixPath, version: Version | None, exists_function: Callable[[str], bool] | None = None, glob_function: Callable[[str], list[str]] | None = None, ):filepath:POSIX 格式的文件路径;version:Version实例,None表示关闭版本化;exists_function/glob_function:可注入的文件系统探测函数,默认分别为Path.exists与glob.iglob。注入机制让同一套版本化逻辑可以复用到本地文件系统之外的存储(如通过fsspec访问 S3/GCS),测试 tests/io/test_core.py 中的MyVersionedDataset正是传入self._fs.exists与self._fs.glob来支持任意协议。
版本化目录结构:_get_versioned_path(version)(第 950 行)把路径构造成filepath / version / filepath.name,即:
data/model.pkl/2024-01-15T10.30.00.000Z/model.pkl版本解析流程:
resolve_load_version():若显式指定了load版本则直接返回;否则通过_fetch_latest_load_version()用glob列出所有版本并取字典序最大(最新)者,结果会缓存以避免重复文件系统操作;resolve_save_version():若显式指定了save版本则直接返回;否则用generate_timestamp()生成并缓存;_get_save_path()会检查目标版本路径是否已存在,若存在则抛出DatasetError(版本化数据集不允许覆盖保存);_is_unsafe_version会拒绝包含路径分隔符、.、..的版本字符串,防止路径穿越。
一致性警告:AbstractVersionedDataset._save_wrapper(第 1013 行)在保存后会比对save_version与load_version,若不一致(例如显式指定了某个中间数据集的 load 版本),会发出_CONSISTENCY_WARNING警告,提示这种"保存版本与加载版本不一致"的做法会引发数据不一致风险,应尽量避免为中间数据集固定 load 版本。
list_versions 版本审计:list_versions(full_path=True)(第 970 行)返回按时间倒序排列的所有版本。full_path=False时仅返回版本时间戳字符串,例如:
['2024-01-15T10.30.00.000Z', '2024-01-14T09.15.00.000Z']该功能可用于版本历史审计、追踪数据变更或实现自定义版本选择逻辑。
版本化数据集的 catalog 配置
my_dataset: type: <path-to-my-own-dataset>.MyOwnDataset filepath: data/01_raw/my_data.csv versioned: true param1: <param1-value> # param1 is a required argument # param2 will be True by default注意:versioned: true会触发parse_dataset_definition把Version(load_version, save_version)注入配置(core.py 第 719-725 行);而配置中直接出现的version键属于保留字,会被移除并告警。此外,HTTP(S) 协议不支持版本化——get_protocol_and_path会在协议为 http(s) 且传入了版本时抛出DatasetError(见 core.py 第 1081 行)。
MemoryDataset:内存数据集
MemoryDataset定义在 kedro/io/memory_dataset.py,用于在 Python 进程内存中读写数据,是管道节点间传递中间结果(特别是ParallelRunner未显式配置的输出)的默认实现。其_EPHEMERAL属性为True,表示数据不持久化。
构造函数:
MemoryDataset(data=_EMPTY, copy_mode=None, metadata=None)copy_mode 拷贝模式是理解MemoryDataset的关键。可取值为"deepcopy"、"copy"、"assign"三种,若不指定则由_infer_copy_mode根据数据类型自动推断:
pandas.DataFrame、numpy.ndarray→"copy"(浅拷贝);- 其他"DataFrame"类型或
ibis.Table→"assign"(直接引用,零拷贝); - 其余类型 →
"deepcopy"(深拷贝)。
这种设计在"避免意外共享可变对象"与"大数据不重复拷贝的开销"之间取得平衡。若传入非法copy_mode,_copy_with_mode会抛出DatasetError,提示合法值为deepcopy, copy, assign。注意_EPHEMERAL是实例属性在__init__中被设为True。
CachedDataset:内存缓存包装器
CachedDataset定义在 kedro/io/cached_dataset.py,是一个包装器数据集:它把被包装的数据集与一个MemoryDataset缓存组合,load时优先命中缓存,从而避免反复访问慢速存储介质。
其 catalog.yml 配置方式(必须把versioned标志声明在包装器上,而不是被包装数据集上):
test_ds: type: CachedDataset versioned: true dataset: type: pandas.CSVDataset filepath: example.csv实现要点:
load()逻辑:缓存命中则直接读缓存,否则从底层数据集加载并写入缓存;save()逻辑:同时写底层数据集与缓存;_exists():缓存或底层任一存在即返回True;_SINGLE_PROCESS = True:该类无法与ParallelRunner一起使用(因为进程间 pickle 序列化会清空缓存,源码在__getstate__中会打印 "clearing cache to pickle" 警告);如需并行,官方注释建议改用ThreadRunner或SequentialRunner;versioned标志若写在了被包装数据集内,_from_config会抛出ValueError,引导用户把版本化声明放到CachedDataset层。
DataCatalog:数据集管理的核心
DataCatalog定义在 kedro/io/data_catalog.py,是 Kedro 数据管理的"注册中心",向程序任意位置提供统一的load/save能力。官方文档描述它为:"A centralized registry for managing datasets in a Kedro project"(管理 Kedro 项目数据集的集中式注册表)。
两种构造方式
方式一:直接传入数据集实例
from kedro.io import DataCatalog, MemoryDataset datasets = { "cars": MemoryDataset(data={"type": "car", "capacity": 5}), "planes": MemoryDataset(data={"type": "jet", "capacity": 200}), } catalog = DataCatalog(datasets=datasets) cars_data = catalog.load("cars") catalog.save("planes", {"type": "propeller", "capacity": 100})方式二:从配置工厂创建
DataCatalog.from_config(catalog, credentials, load_versions, save_version)接受配置字典与凭据字典。其中type指定数据集类,credentials键引用凭据字典中的条目:
config = { "cars": { "type": "pandas.CSVDataset", "filepath": "cars.csv", "save_args": {"index": False} }, "boats": { "type": "pandas.CSVDataset", "filepath": "s3://aws-bucket-name/boats.csv", "credentials": "boats_credentials", "save_args": {"index": False} } } credentials = { "boats_credentials": { "client_kwargs": { "aws_access_key_id": "<your key id>", "aws_secret_access_key": "<your secret>" } } } catalog = DataCatalog.from_config(config, credentials) df = catalog.load("cars") catalog.save("boats", df)from_config会校验load_versions中引用的数据集名是否存在于配置或模式中,否则抛出DatasetNotFoundError。若同一数据集名同时出现在datasets参数与配置中,直接传入的实例优先,配置项会被跳过并记录警告(见 data_catalog.py 第 283-291 行)。
懒加载机制:_LazyDataset
DataCatalog采用懒加载策略提升性能:从配置注册数据集时只创建_LazyDataset占位对象(持有 name、config、load_version、save_version),真正实例化推迟到首次访问时。get()/__getitem__访问时调用materialize()完成实例化。这也解释了为什么DataCatalog.__init__可以快速完成——它不必为 catalog.yml 中的每个数据集立即构造对象。
__setitem__支持三类值:AbstractDataset实例(直接注册)、_LazyDataset(懒注册)、其他原始数据(自动包装为MemoryDataset):
catalog = DataCatalog() catalog["data_df"] = df # 原始数据自动包装为 MemoryDataset catalog["data_csv_dataset"] = csv_dataset # 数据集实例直接注册常用 API 速览
load(name, version=None)/save(name, data):加载/保存数据,未找到数据集时抛DatasetNotFoundError;load可传version参数指定具体版本(仅对版本化数据集生效);exists(name):检查数据集输出是否存在;release(name):释放数据集缓存数据;confirm(name):确认数据集(如事务性写入提交),数据集无confirm方法时抛DatasetError;keys()/values()/items()/__len__/__contains__:目录的字典式遍历接口(懒数据集与已实例化数据集都会列出);filter(name_regex, type_regex, by_type):按名称正则、类型正则或具体类型过滤数据集名,支持预编译re.Pattern;get_type(name):获取数据集的全限定类型名(如kedro.io.memory_dataset.MemoryDataset),且不会把按模式解析的数据集加入目录;to_config():把目录序列化为(catalog, credentials, load_versions, save_version)四元组,可配合from_config完成"目录保存—重建"的往返(round-trip)。序列化时会跳过内存型数据集,并把validator声明重新注入。
版本管理:load_versions 与 save_version
DataCatalog支持目录级版本管理:
load_versions:数据集名到具体加载版本的映射,对未启用版本化的数据集无影响;save_version:所有启用版本化数据集共用的保存版本。要求:a) 大小写不敏感且符合操作系统文件名的限制;b) 按字典序排序时总是最新版本。
_validate_versions(data_catalog.py 第 1322 行)会同步目录与数据集的版本:若某版本化数据集显式指定了 save 版本且与目录的 save_version 冲突,抛出VersionAlreadyExistsError。测试 tests/io/test_data_catalog.py 中的test_redefine_save_version_via_catalog、test_set_load_and_save_versions等用例验证了这些版本同步行为。
数据集校验:validator 与 validation_enabled
DataCatalog把validator键从数据集配置中剥离(通过ValidatorSpec.from_dataset_config解析),在load/save时统一调用验证器。相关行为:
- 构造参数
validation_enabled(默认True)控制是否启用验证; - 环境变量
KEDRO_DATASET_VALIDATION可覆盖该标志(0/false/off关闭,1/true/on开启),且每次操作时实时读取,方便在不停机的情况下"熔断"验证; _save_validated集合记录已验证过的保存,配合skip_load_after_save可在同一次 save 后跳过冗余的 load 校验(详见 data_catalog.py 中_validate与_is_validation_enabled的实现)。
CatalogConfigResolver:数据集工厂模式与凭据解析
CatalogConfigResolver定义在 kedro/io/catalog_config_resolver.py,负责基于**数据集工厂模式(dataset factory patterns)**与凭据字典动态生成数据集配置。DataCatalog内部持有一个 resolver,通过config_resolver属性对外暴露。
模式匹配与优先级
catalog 配置中带{}占位符的键被视为模式,例如{namespace}.int_{name}。resolver 会:
- 提取所有模式并按特异性排序:优先匹配花括号外字符更多的模式,其次占位符更多的模式,再按字母序(
_sort_patterns,见 catalog_config_resolver.py 第 225 行); resolve_pattern(ds_name)依次尝试:已解析配置 → 数据集模式 → 用户自定义 catch-all 模式 → 运行时模式;- 模式配置中的占位符(如
filepath: "{name}.csv")会被format_map替换为实际值。若配置中使用了模式名中不存在的占位符键,_validate_pattern_config会抛出DatasetError; - catch-all 模式限制:特异性为 0 的 catch-all 模式(如
{name})整个 catalog 只允许一个,多个会触发DatasetError。
官方示例:
config = { "{namespace}.int_{name}": { "type": "pandas.CSVDataset", "filepath": "{name}.csv", "credentials": "db_credentials", } } credentials = {"db_credentials": {"user": "username", "pass": "pass"}} resolver = CatalogConfigResolver(config=config, credentials=credentials) resolved_config = resolver.resolve_pattern("data.int_customers") # {'type': 'pandas.CSVDataset', 'filepath': 'customers.csv', # 'credentials': {'user': 'username', 'pass': 'pass'}}凭据解析
_resolve_credentials会递归遍历配置,把credentials: "boats_credentials"这样的字符串引用替换为凭据字典中的实际值;若引用的凭据不存在,抛出带指引的KeyError。_unresolve_credentials则执行反向操作,把内联凭据抽离为<数据集名>_credentials引用键——这正是DataCatalog.to_config()实现"凭据与配置分离"的机制。
运行时模式 {default}
DataCatalog.default_runtime_patterns = {"{default}": {"type": "kedro.io.MemoryDataset"}}(见 data_catalog.py 第 211 行)。该模式兜底匹配所有未显式配置、也未命中任何用户模式的数据集名,将其实例化为MemoryDataset。这解释了 Kedro 管道中节点间隐式传递数据的机制——未在 catalog 中声明的中间数据集默认在内存中流转。
SharedMemoryDataCatalog 与并行执行
协议层:CatalogProtocol 与 SharedMemoryCatalogProtocol
CatalogProtocol(core.py 第 1148 行)是@runtime_checkable的Protocol,定义了目录应有的通用接口:__contains__、keys、values、items、get、save、load、release、confirm、exists、from_config等。SharedMemoryCatalogProtocol在它之上增加set_manager_datasets(manager)与validate_catalog()两个方法,用于多进程共享内存场景。
SharedMemoryDataCatalog
SharedMemoryDataCatalog继承DataCatalog,专为ParallelRunner等多进程场景设计:
- 默认运行时模式改为
{"{default}": {"type": "kedro.io.SharedMemoryDataset"}},即未显式配置的输出默认使用共享内存数据集; set_manager_datasets(manager):把multiprocessing.managers.SyncManager注入所有SharedMemoryDataset;validate_catalog():逐一检查数据集是否可 pickle 序列化。含_SINGLE_PROCESS = True标志或无法序列化的数据集会被收集,最终抛出AttributeError列出全部不兼容数据集;若数据集挂在云存储协议(_protocol非file)上且无法 pickle,会额外发出UserWarning,建议改用ThreadRunner或SequentialRunner。
SharedMemoryDataset
SharedMemoryDataset(kedro/io/shared_memory_dataset.py)通过SyncManager代理一个共享的MemoryDataset,使数据能够在多个进程间同步访问。其save方法在底层序列化失败时会抛出DatasetError,提示"ParallelRunner 隐式内存数据集只能用于可序列化的数据"。
异常体系
kedro.io的异常都定义在 core.py 第 161-202 行,形成清晰的继承层次:
| 异常 | 父类 | 抛出场景 |
|---|---|---|
DatasetError | Exception | 数据集读写失败时的通用错误(AbstractDataset实现应提供有指导性的错误信息) |
DatasetNotFoundError | DatasetError | 尝试使用目录中不存在的数据集 |
DatasetAlreadyExistsError | DatasetError | 向目录添加已存在的数据集 |
VersionNotFoundError | DatasetError | 版本化数据集没有可用加载版本(或权限不足无法访问版本目录) |
VersionAlreadyExistsError | DatasetError | 向目录添加数据集时其保存版本与目录已设置的保存版本冲突 |
另外,parse_dataset_definition还会对配置中的非法字符做校验:validate_on_forbidden_chars禁止字符串值包含空格与分号。_redact_url_credentials等工具函数则会在错误信息与repr输出中脱敏 URL 中的凭据信息(如预签名 URL 的签名参数),防止密钥泄露到日志中。
源码与测试验证路径
想要深入理解kedro.io的行为,推荐从以下仓库路径入手:
- 核心实现:kedro/io/core.py(
AbstractDataset、AbstractVersionedDataset、Version、parse_dataset_definition、CatalogProtocol、异常类) - 目录实现:kedro/io/data_catalog.py(
DataCatalog、SharedMemoryDataCatalog、_LazyDataset) - 配置解析:kedro/io/catalog_config_resolver.py(
CatalogConfigResolver、数据集工厂模式、凭据解析) - 内存与缓存:kedro/io/memory_dataset.py、kedro/io/cached_dataset.py、kedro/io/shared_memory_dataset.py
- 模块导出:kedro/io/init.py
- 测试用例:tests/io/test_core.py(自定义数据集与版本化数据集的标准写法)、tests/io/test_data_catalog.py(目录生命周期、版本同步、模式匹配、validator 行为)、tests/io/test_cached_dataset.py、tests/io/test_shared_memory_dataset.py
在实战中,配置层面的更多用法可继续参阅仓库文档 数据目录说明、数据集工厂配置 以及 分区与增量数据集。掌握kedro.io的类体系与解析流程,是编写自定义数据集、排查 catalog 配置问题以及理解 Kedro 数据流的关键一步。
【免费下载链接】kedroKedro is a toolbox for production-ready data science. It uses software engineering best practices to help you create data engineering and data science pipelines that are reproducible, maintainable, and modular.项目地址: https://gitcode.com/GitHub_Trending/ke/kedro
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考