Kedro 数据集与数据目录(kedro.io)完整指南:从 AbstractDataset 到 DataCatalog 的源码级解析 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/kedroKedro 的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.AbstractDatasetClass所有 Kedro 数据集的基类kedro.io.AbstractVersionedDatasetClass版本化数据集的基类kedro.io.CachedDatasetClass在内存中缓存数据的数据集包装器kedro.io.DataCatalogClass管理 Kedro 管道中使用的数据集kedro.io.CatalogProtocolClass定义在目录中管理数据集的通用接口kedro.io.SharedMemoryDataCatalogClass在共享内存上下文中管理数据集kedro.io.SharedMemoryCatalogProtocolClass扩展CatalogProtocol以支持共享内存场景kedro.io.CatalogConfigResolverClass基于数据集工厂模式与凭据解析数据集配置kedro.io.MemoryDatasetClass在内存中存储数据的数据集kedro.io.VersionClass表示数据集版本信息kedro.io.DatasetAlreadyExistsErrorException数据集已存在时抛出kedro.io.DatasetErrorException通用数据集错误kedro.io.DatasetNotFoundErrorException数据集未找到时抛出此外__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判断目标数据是否已存在。基类默认实现只打印警告并返回Falseexists()方法会包装它并捕获异常转换为DatasetError_release()释放缓存数据基类默认为空操作passrelease()方法负责包装。自动装饰机制init_subclassAbstractDataset通过__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_configcore.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, param2True): 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(param1self._param1, param2self._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_definitioncore.py 第 621 行起承担配置解析的核心职责其关键行为包括配置必须包含type键值为类的全限定名如kedro_datasets.pandas.CSVDataset或类对象解析时按_DEFAULT_PACKAGES [kedro.io., kedro_datasets., ]三个前缀依次尝试导入因此 catalog.yml 中既可以写全限定名也可以只写pandas.CSVDataset这样的短名若最终找不到类会抛出DatasetError并附带提示请确认是否安装了 kedro-datasetspip install kedro-datasets解析出的类必须是AbstractDataset的子类否则报DatasetError配置中的versioned: true标志会触发版本注入见下文版本化章节validator键属于 catalog 级配置会被移除并警告不会传给数据集构造函数出于向后兼容遇到以Set结尾的类型名旧拼写会发出警告Since kedro-datasets 2.0, Dataset is spelled with a lowercase s。版本化数据集AbstractVersionedDataset 与 VersionVersion加载版本与保存版本的载体Version定义在 core.py 第 601 行是一个namedtuple(Version, [load, save])若Version.load为None则加载最新可用版本若Version.save为None保存版本将自动生成格式为YYYY-MM-DDThh.mm.ss.sssZUTC 时间戳。时间戳由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, ):filepathPOSIX 格式的文件路径versionVersion实例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_pathTrue)第 970 行返回按时间倒序排列的所有版本。full_pathFalse时仅返回版本时间戳字符串例如[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_modeNone, metadataNone)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或SequentialRunnerversioned标志若写在了被包装数据集内_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(datasetsdatasets) 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 行。懒加载机制_LazyDatasetDataCatalog采用懒加载策略提升性能从配置注册数据集时只创建_LazyDataset占位对象持有 name、config、load_version、save_version真正实例化推迟到首次访问时。get()/__getitem__访问时调用materialize()完成实例化。这也解释了为什么DataCatalog.__init__可以快速完成——它不必为 catalog.yml 中的每个数据集立即构造对象。__setitem__支持三类值AbstractDataset实例直接注册、_LazyDataset懒注册、其他原始数据自动包装为MemoryDatasetcatalog DataCatalog() catalog[data_df] df # 原始数据自动包装为 MemoryDataset catalog[data_csv_dataset] csv_dataset # 数据集实例直接注册常用 API 速览load(name, versionNone)/save(name, data)加载/保存数据未找到数据集时抛DatasetNotFoundErrorload可传version参数指定具体版本仅对版本化数据集生效exists(name)检查数据集输出是否存在release(name)释放数据集缓存数据confirm(name)确认数据集如事务性写入提交数据集无confirm方法时抛DatasetErrorkeys()/values()/items()/__len__/__contains__目录的字典式遍历接口懒数据集与已实例化数据集都会列出filter(name_regex, type_regex, by_type)按名称正则、类型正则或具体类型过滤数据集名支持预编译re.Patternget_type(name)获取数据集的全限定类型名如kedro.io.memory_dataset.MemoryDataset且不会把按模式解析的数据集加入目录to_config()把目录序列化为(catalog, credentials, load_versions, save_version)四元组可配合from_config完成目录保存—重建的往返round-trip。序列化时会跳过内存型数据集并把validator声明重新注入。版本管理load_versions 与 save_versionDataCatalog支持目录级版本管理load_versions数据集名到具体加载版本的映射对未启用版本化的数据集无影响save_version所有启用版本化数据集共用的保存版本。要求a) 大小写不敏感且符合操作系统文件名的限制b) 按字典序排序时总是最新版本。_validate_versionsdata_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_enabledDataCatalog把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会抛出DatasetErrorcatch-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(configconfig, credentialscredentials) 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 与 SharedMemoryCatalogProtocolCatalogProtocolcore.py 第 1148 行是runtime_checkable的Protocol定义了目录应有的通用接口__contains__、keys、values、items、get、save、load、release、confirm、exists、from_config等。SharedMemoryCatalogProtocol在它之上增加set_manager_datasets(manager)与validate_catalog()两个方法用于多进程共享内存场景。SharedMemoryDataCatalogSharedMemoryDataCatalog继承DataCatalog专为ParallelRunner等多进程场景设计默认运行时模式改为{{default}: {type: kedro.io.SharedMemoryDataset}}即未显式配置的输出默认使用共享内存数据集set_manager_datasets(manager)把multiprocessing.managers.SyncManager注入所有SharedMemoryDatasetvalidate_catalog()逐一检查数据集是否可 pickle 序列化。含_SINGLE_PROCESS True标志或无法序列化的数据集会被收集最终抛出AttributeError列出全部不兼容数据集若数据集挂在云存储协议_protocol非file上且无法 pickle会额外发出UserWarning建议改用ThreadRunner或SequentialRunner。SharedMemoryDatasetSharedMemoryDatasetkedro/io/shared_memory_dataset.py通过SyncManager代理一个共享的MemoryDataset使数据能够在多个进程间同步访问。其save方法在底层序列化失败时会抛出DatasetError提示ParallelRunner 隐式内存数据集只能用于可序列化的数据。异常体系kedro.io的异常都定义在 core.py 第 161-202 行形成清晰的继承层次异常父类抛出场景DatasetErrorException数据集读写失败时的通用错误AbstractDataset实现应提供有指导性的错误信息DatasetNotFoundErrorDatasetError尝试使用目录中不存在的数据集DatasetAlreadyExistsErrorDatasetError向目录添加已存在的数据集VersionNotFoundErrorDatasetError版本化数据集没有可用加载版本或权限不足无法访问版本目录VersionAlreadyExistsErrorDatasetError向目录添加数据集时其保存版本与目录已设置的保存版本冲突另外parse_dataset_definition还会对配置中的非法字符做校验validate_on_forbidden_chars禁止字符串值包含空格与分号。_redact_url_credentials等工具函数则会在错误信息与repr输出中脱敏 URL 中的凭据信息如预签名 URL 的签名参数防止密钥泄露到日志中。源码与测试验证路径想要深入理解kedro.io的行为推荐从以下仓库路径入手核心实现kedro/io/core.pyAbstractDataset、AbstractVersionedDataset、Version、parse_dataset_definition、CatalogProtocol、异常类目录实现kedro/io/data_catalog.pyDataCatalog、SharedMemoryDataCatalog、_LazyDataset配置解析kedro/io/catalog_config_resolver.pyCatalogConfigResolver、数据集工厂模式、凭据解析内存与缓存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),仅供参考