Kedro Data Catalog 完全指南:从 catalog.yml 配置到源码级运行原理
【免费下载链接】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 中的Data Catalog是项目所有数据源的注册中心:它以catalog.yml文件为入口,将节点的输入输出名称映射为DataCatalog类中的具体数据集对象。本文以 Data Catalog 介绍 为主线,完整梳理 catalog 的核心概念、catalog.yml的全部配置维度(类型、路径、参数、凭据、版本化、多环境覆盖)、数据集工厂、懒加载机制,并结合 DataCatalog 源码 解释这些配置在底层是如何被解析和执行的。读完本文,你将能够独立完成 Kedro 项目数据层的注册、定制与调试。
Data Catalog 是什么
在 Kedro 项目中,Data Catalog 是所有可供项目使用的数据源的注册表。它由一个 YAML 文件(catalog.yml)承载,把节点输入和输出的名称作为DataCatalog类的键(key),与具体的数据集实现一一对应。节点通过名称引用数据,Kedro 再通过 Catalog 把名称解析为实际的文件、表或对象存储资源。
两个与版本相关的关键事实需要先了解:
- 从 Kedro
0.19.0起,具体的数据集实现(如CSVDataset)不再包含在 Kedro 核心包中,需要从kedro-datasets包导入; - 从
kedro-datasets2.0.0起,所有数据集名称中的大写 "S"(DataSet)改为小写 "s"(Dataset),例如CSVDataSet现在写作CSVDataset。
kedro-datasets为常见文件类型与文件系统提供了开箱即用的数据集实现,是 catalog 配置中type字段的主要来源。
catalog.yml基础:注册第一批数据集
Data Catalog 的入口文件是catalog.yml,通常位于conf/base/目录下。以下示例来自 set_up_data 教程,注册了两个 CSV 数据集和一个 Excel 数据集:
companies: type: pandas.CSVDataset filepath: data/01_raw/companies.csv reviews: type: pandas.CSVDataset filepath: data/01_raw/reviews.csv shuttles: type: pandas.ExcelDataset filepath: data/01_raw/shuttles.xlsx load_args: engine: openpyxl # Use modern Excel engine (the default since Kedro 0.18.0)在本地文件系统读取或保存文件时,每个条目需要三个要素:
- 数据集名称(key):顶层键,作为节点输入输出引用的标识符;
type:声明使用的数据集类;filepath:文件位置。
配置数据集参数
catalog.yml中数据集配置的层次结构为:
- 顶层键是数据集名称,例如
shuttles、weather; - 下一层包含多个键:第一个必填键是
type,声明数据集类型;其余键是数据集参数,因实现而异; - 部分数据集参数可以进一步由底层库决定:例如
shuttles的load_args由 pandas 加载 CSV 的选项定义,而weather的save_args由 Snowpark 的saveAsTable方法定义。
shuttles: # Dataset name type: pandas.ExcelDataset # Dataset type filepath: data/01_raw/shuttles.xlsx # pandas.ExcelDataset parameter load_args: # pandas.ExcelDataset parameter engine: openpyxl # Pandas option for loading CSV files weather: # Dataset name type: snowflake.SnowparkTableDataset # Dataset type table_name: "weather_data" database: "meteorology" schema: "observations" credentials: snowflake_client save_args: # snowflake.SnowparkTableDataset parameter mode: overwrite # Snowpark saveAsTable input option column_order: name table_type: ''注意:Kedro 数据集会把
load_args/save_args直接委托给底层实现。完整参数列表请查阅kedro-datasets文档中对应数据集类的__init__方法,其中会给出对底层库 API(如pandas.read_excel)的参数引用。
数据集type与filepath:连接任意数据存储
type:支持的数据集类型
Kedro 支持连接 CSV、Excel、Parquet、Feather、HDF5、JSON、pickle 对象、SQL 表、SQL 查询等多种数据形态,底层依赖 pandas、PySpark、NetworkX、Matplotlib 等库。完整列表见kedro-datasets文档。
filepath:基于 fsspec 的协议寻址
Kedro 依赖fsspec从各类数据存储读取和写入数据。filepath应使用protocol://path/to/data的通用形式;若不写协议,默认按本地文件系统处理(等价于file://)。可用协议包括:
| 协议 | 说明 |
|---|---|
file:// | 本地或网络文件系统,默认协议,允许相对路径 |
hdfs://user@server:port/path/to/data | Hadoop 分布式文件系统(HDFS),面向集群内高可靠、副本文件 |
s3://my-bucket-name/path/to/data | Amazon S3 远程二进制存储(常与 EC2 搭配),基于 s3fs 库 |
s3://my-bucket-name/path/to/data | S3 兼容存储(如 MinIO),同样基于 s3fs |
gcs:// | Google Cloud Storage,基于 gcsfs |
abfs:// | Azure Blob Storage / Azure Data Lake Storage Gen2 |
http:///https:// | 直接从 HTTP Web 服务器读取数据 |
fsspec还提供 SSH、FTP、WebHDFS 等其他文件系统实现。
深入配置:load_args / save_args / fs_args / validator
除type和filepath外,catalog 还接受几组影响数据加载、保存和访问方式的设置:
load_args和save_args:控制底层第三方库如何加载/保存数据。例如pandas.CSVDataset的load_args会作为关键字参数传给pd.read_csv,save_args传给pd.DataFrame.to_csv:
cars: type: pandas.CSVDataset filepath: data/01_raw/company/cars.csv load_args: sep: ',' save_args: index: False date_format: '%Y-%m-%d %H:%M' decimal: .validator:在每次加载或保存时,按 schema 或自定义规则校验数据,详见 数据集校验。fs_args:控制 Kedro 与文件系统本身的交互。顶层键传给底层文件系统类(例如 GCS 的GCSFileSystem),而open_args_load和open_args_save传给文件系统的open方法,控制文件在加载/保存时如何被打开:
test_dataset: type: ... fs_args: project: test_project # 传给 GCSFileSystem open_args_load: mode: "r" encoding: "utf-8" open_args_save: mode: "a" # 追加模式保存每个数据集实现的默认加载、保存和文件系统参数,定义在实现类中的
DEFAULT_LOAD_ARGS、DEFAULT_SAVE_ARGS、DEFAULT_FS_ARGS中,可在kedro-datasets文档中查询。
凭据(credentials):安全地访问远程存储
Kedro 在实例化DataCatalog前,会先从项目配置中读取credentials.yml(通常位于conf/local/)中的凭据,将结果字典通过credentials参数传入DataCatalog.from_config()。catalog.yml条目通过顶层的credentials:键按名称引用凭据块。
假设conf/local/credentials.yml中包含:
dev_s3: client_kwargs: aws_access_key_id: key aws_secret_access_key: secretcatalog 条目即可引用:
motorbikes: type: pandas.CSVDataset filepath: s3://your_bucket/data/02_intermediate/company/motorbikes.csv credentials: dev_s3 load_args: sep: ','Catalog 会在凭据字典中查找dev_s3,将其值作为credentials参数传入数据集构造函数。
数据集版本化(versioning)
在 catalog 条目中添加versioned: True即可启用版本化:
cars: type: pandas.CSVDataset filepath: data/01_raw/company/cars.csv versioned: True启用后,filepath成为存储各版本的目录基础。每次流水线运行产生新版本时,会存储在<filepath>/<version>/<filename>中,其中<version>是格式为YYYY-MM-DDThh.mm.ss.sssZ的时间戳。
- 默认情况下,
kedro run加载最新版本; - 要加载指定版本,使用
--load-versions参数,以数据集名:版本时间戳的格式传入:
kedro run --load-versions=cars:YYYY-MM-DDThh.mm.ss.sssZ- 版本化数据集提供
list_versions()方法列出所有可用版本:
# In a Kedro session or notebook dataset = catalog._datasets["cars"] versions = dataset.list_versions(full_path=False) print(versions) # ['2024-01-15T10.30.00.000Z', '2024-01-14T09.15.00.000Z', ...]full_path=True(默认)返回各版本的完整文件路径,full_path=False返回版本字符串(时间戳);版本按时间倒序返回(最新的在前)。
版本化支持的前提:数据集必须继承kedro.io.AbstractVersionedDataset类以接受构造参数version,并在_save/_load方法中通过_get_save_path和_get_load_path使用版本化路径。要验证数据集是否支持版本化,可检查其类继承关系,例如CSVDataset(AbstractVersionedDataset[pd.DataFrame, pd.DataFrame])即支持版本化。
HTTP(S) 是数据集实现支持的文件系统,但不能与版本化组合使用。
多环境配置:conf/base与conf/local
Kedro 通过配置加载器扫描conf文件夹下的配置文件:先扫描conf/base,再扫描conf/local(指定的覆盖环境),合并后返回配置字典。因此可以通过在不同环境放置不同版本的catalog.yml来区分开发、测试与生产配置。
例如conf/base/catalog.yml中定义生产环境的 S3 位置:
cars: filepath: s3://my_bucket/cars.csv type: pandas.CSVDataset在conf/local/catalog.yml中覆盖为本地文件:
cars: filepath: data/01_raw/cars.csv type: pandas.CSVDataset当流水线代码引用cars数据集时,本地运行使用conf/local的条目,而conf/local无覆盖时使用生产条目。完整的合并规则见 配置基础文档。
数据集工厂(dataset factories)
数据集工厂(Kedro0.18.12引入)让你用模式(pattern)批量注册配置相似的数据集,大幅减少 catalog 条目。例如以下两个条目:
factory_data: type: pandas.CSVDataset filepath: data/01_raw/factory_data.csv process_data: type: pandas.CSVDataset filepath: data/01_raw/process_data.csv可以重写为一个工厂模式:
"{name}_data": type: pandas.CSVDataset filepath: data/01_raw/{name}_data.csv运行时,模式会与流水线节点inputs/outputs中定义的数据集名称进行匹配。工厂模式的行为类似正则表达式,可理解为反向的f-string:输入数据集factory_data匹配模式{name}_data时,name被解析为factory;输出数据集process_data则解析为process。
注意:工厂模式必须用引号包裹,否则会触发 YAML 解析错误。
三类模式
- 数据集模式(dataset patterns):在
catalog.yml中用{name}_data这类占位符显式定义,任何符合命名规则的数据集都会被动态解析; - 用户 catch-all 模式:当没有数据集模式匹配时作为兜底,使用
{default_dataset}占位符;每个 catalog 只允许一个,指定多个会抛出DatasetError; - 默认运行时模式(default runtime patterns):Kedro 内置的模式,当数据集未在 catalog 中定义时(通常是流水线运行中产生的中间数据集)自动使用,例如
DataCatalog的{"{default}": {"type": "kedro.io.MemoryDataset"}}和SharedMemoryDataCatalog的{"{default}": {"type": "kedro.io.SharedMemoryDataset"}}。
模式解析顺序
- 数据集模式:最明确,优先匹配;
- 用户 catch-all 模式:无数据集模式匹配时的回退;
- 默认运行时模式:上述均不匹配时由 Kedro 在运行时自动创建
MemoryDataset或SharedMemoryDataset。
默认情况下catalog.get()不启用运行时模式,除非显式设置fallback_to_runtime_pattern=True;kedro run执行时则自动启用。
Pipeline 感知的 catalog 命令
DataCatalog通过CatalogCommandsMixin暴露一组检查模式解析结果的命令,Kedro 在初始化 session 时自动装配:
kedro catalog describe-datasets:描述流水线中用到的数据集,按解析方式(显式条目、工厂匹配、默认兜底)分组;kedro catalog list-patterns:按优先级列出 catalog 中所有工厂模式;kedro catalog resolve-patterns:将流水线数据集对全部模式进行解析,返回完整配置。
懒加载(Lazy loading)
从 Kedro0.19.10起,DataCatalog引入_LazyDataset辅助类以优化性能。它先存储数据集的配置与版本化信息,而不立即实例化数据集对象,将真正的创建(materialisation)推迟到数据集被首次访问时。
从catalog.yml实例化DataCatalog时,Kedro 不会一次性创建所有底层数据集对象,而是把每个数据集包装成_LazyDataset注册进 catalog;首次访问时(直接访问或流水线执行期间)再自动实例化:
In [1]: catalog Out[1]: { 'shuttles': kedro_datasets.pandas.excel_dataset.ExcelDataset } # 此时 'shuttles' 尚未完全实例化——只注册了其配置 In [2]: catalog["shuttles"] Out[2]: kedro_datasets.pandas.excel_dataset.ExcelDataset( filepath=PurePosixPath('/Projects/default/data/01_raw/shuttles.xlsx'), protocol='file', load_args={'engine': 'openpyxl'}, save_args={'index': False}, writer_args={'engine': 'openpyxl'} ) # 访问数据集即触发实例化 In [3]: catalog Out[3]: { 'shuttles': kedro_datasets.pandas.excel_dataset.ExcelDataset( filepath=PurePosixPath('/Projects/default/data/01_raw/shuttles.xlsx'), ... ) }该机制对大型 catalog 的启动开销有明显改善。在流水线预热阶段可以提前强制实例化所有数据集,从而尽早发现配置或导入错误、校验外部依赖、确保执行前所有数据集可创建。虽然_LazyDataset不对最终用户暴露,但理解它有助于调试 catalog 行为与数据集实例化问题。
在代码中以编程方式使用 DataCatalog
除 YAML 配置外,也可以通过kedro.io.DataCatalog以代码方式定义数据源,适合在catalog.py或 notebook 中构建 IO 层:
from kedro.io import DataCatalog from kedro_datasets.pandas import ( CSVDataset, SQLTableDataset, SQLQueryDataset, ParquetDataset, ) catalog = DataCatalog( { "bikes": CSVDataset(filepath="../data/01_raw/bikes.csv"), "cars": CSVDataset(filepath="../data/01_raw/cars.csv", load_args=dict(sep=",")), "cars_table": SQLTableDataset( table_name="cars", credentials=dict(con="sqlite:///kedro.db") ), "scooters_query": SQLQueryDataset( sql="select * from cars where gear=4", credentials=dict(con="sqlite:///kedro.db"), ), "ranked": ParquetDataset(filepath="ranked.parquet"), } )DataCatalog提供完整的映射式 API:keys()、values()、items()、__iter__、__getitem__、__contains__、__len__,以及load(name, version=None)、save(name, data)、exists(name)、release(name)、confirm(name)等方法,详见 如何在代码中访问 Data Catalog。
实用 YAML 配方精选
Data Catalog YAML 示例页 收录了大量可直接复用的catalog.yml配方,以下为几个典型场景:
读取带压缩的 CSV:
boats: type: pandas.CSVDataset filepath: data/01_raw/company/boats.csv.gz load_args: sep: ',' compression: 'gzip' fs_args: open_args_load: mode: 'rb'从 S3 加载 CSV(带凭据与加载参数):
motorbikes: type: pandas.CSVDataset filepath: s3://your_bucket/data/02_intermediate/company/motorbikes.csv credentials: dev_s3 load_args: sep: ',' skiprows: 5 skipfooter: 1 na_values: ['#NA', NA]从 GCS 加载 Excel(带文件系统参数):
rockets: type: pandas.ExcelDataset filepath: gcs://your_bucket/data/02_intermediate/company/motorbikes.xlsx fs_args: project: my-project credentials: my_gcp_credentials save_args: sheet_name: Sheet1SQL 表与查询(凭据块须包含 SQLAlchemy 兼容的连接串con):
scooters: type: pandas.SQLTableDataset credentials: scooters_credentials table_name: scooters load_args: index_col: [name] columns: [name, gear] save_args: if_exists: replace scooters_query: type: pandas.SQLQueryDataset credentials: scooters_credentials sql: select * from cars where gear=4用 YAML 锚点复用公共配置:&csv命名模板块、<<: *csv插入模板内容;模板条目名必须以_开头,Kedro 才不会将其实例化为数据集:
_csv: &csv type: spark.SparkDataset file_format: csv load_args: sep: ',' na_values: ['#NA', NA] header: True inferSchema: False cars: <<: *csv filepath: s3a://data/01_raw/cars.csv bikes: <<: *csv filepath: s3a://data/01_raw/bikes.csv load_args: header: False # 局部声明的键会覆盖插入的键源码印证:DataCatalog 如何从配置构建
DataCatalog的构建入口是类方法 from_config(),它接收配置字典、凭据字典、加载版本与保存版本等参数。从源码结构看,其完整流程为:
- 实例化时通过
CatalogConfigResolver提取并排序 catalog 中的工厂模式(见 catalog_config_resolver.py 中的_extract_patterns、_sort_patterns、match_dataset_pattern),同时从credentials.yml读取凭据并注入条目; - 每个配置条目经
parse_dataset_definition(位于 core.py)解析出数据集类与构造参数; - 数据集以
_LazyDataset占位(data_catalog.py 第 42 行)注册,首次访问时才调用materialize()完成实例化; - 版本化相关的路径解析、时间戳生成(
generate_timestamp)与版本列表查询(list_versions)由core.py中的版本辅助逻辑实现。
这种"配置解析 → 凭据注入 → 模式匹配 → 懒加载实例化"的分层设计,正是 Kedro 能把 YAML 配置、工厂模式和运行时流水线无缝衔接起来的原因。
进阶主题与相关文档
- 分区与增量数据集:处理跨多文件的数据,见 分区与增量数据集概念 与 使用方法;
- 数据集校验:加载/保存时按 schema 校验,见 dataset_validation.md;
- 自定义数据集:创建自己的数据集实现,见 自定义数据集教程;
- 凭据与参数:见 parameters_and_credentials.md。
【免费下载链接】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),仅供参考