Kedro Hooks 机制完全指南:解读 kedro.framework.hooks 的 manager、markers 与 specs 全量 API
【免费下载链接】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 官方 API 文档 kedro.framework.hooks 为核心骨架,系统讲解 Kedro 的 Hooks 扩展机制:manager(全局 hook 管理器)、markers(声明式装饰器)与specs(全部可调用 Hook 的规格定义)三个子模块的职责与用法,并结合 hooks 规范源码、会话与运行器中的调用点 以及 Hook 实践指南,帮助你掌握如何编写、注册、调试 Hook,实现在kedro run执行时间线的任意关键节点注入自定义行为。
一、Hook 机制与模块总览
Kedro 的 Hooks 机制基于 pluggy(即 pytest 所用的插件系统)构建:一个 Hook 由Hook specification(规格)与Hook implementation(实现)两部分组成。Kedro 在源码中预定义了一批规格,声明"在哪个执行节点可以注入行为";用户或插件则提供实现,声明"在该节点做什么"。
在kedro.framework.hooks包(见 kedro/framework/hooks/init.py)中,公开导出了_create_hook_manager与hook_impl,并包含三个子模块,其分工如下表:
| 模块 | 描述 |
|---|---|
kedro.framework.hooks.manager | 提供工具函数,用于在 Kedro 执行进程中创建全局hook_manager单例,并完成 Hook 的注册与插件入口点加载 |
kedro.framework.hooks.markers | 提供声明式 markers(hook_spec/hook_impl),用于标记 Kedro 的 Hook 规格与实现 |
kedro.framework.hooks.specs | 包含 Kedro 执行时间线上全部可调用 Hook 的规格定义(5 个规范类,共 12 个 Hook) |
此外,CLI 层面还有一组独立的 CLI Hooks(before_command_run/after_command_run),定义于 kedro/framework/cli/hooks/specs.py,命名空间为kedro_cli。
二、markers:Hook 的声明式标记
markers.py 是整个机制的基石,全文只有三个关键定义:
import pluggy HOOK_NAMESPACE = "kedro" hook_spec = pluggy.HookspecMarker(HOOK_NAMESPACE) hook_impl = pluggy.HookimplMarker(HOOK_NAMESPACE)HOOK_NAMESPACE = "kedro":所有 Kedro Hook(非 CLI 类)共享的命名空间。pluggy 依赖该命名空间将"规格"与"实现"进行匹配,因此实现方法名必须与规格方法名完全一致。hook_spec:HookspecMarker实例,用作@hook_spec装饰器,声明"这是一个 Hook 规格"。hook_impl:HookimplMarker实例,用作@hook_impl装饰器,声明"这是一个 Hook 实现"。
与之对应,CLI Hooks 使用独立的命名空间与 markers,见 kedro/framework/cli/hooks/markers.py:
CLI_HOOK_NAMESPACE = "kedro_cli" cli_hook_spec = pluggy.HookspecMarker(CLI_HOOK_NAMESPACE) cli_hook_impl = pluggy.HookimplMarker(CLI_HOOK_NAMESPACE)这意味着kedro.hooks与kedro.cli_hooks是两套相互独立的插件体系,分别由各自的 manager 管理。
三、specs:Kedro 执行时间线上的全部 Hook 规格
specs.py 定义了 5 个规范类(namespace),每个类用@hook_spec标记其方法,构成 Kedro 运行生命周期的 12 个核心 Hook。
3.1 命名约定
非错误类 Hook 遵循<before/after>_<noun>_<past_participle>约定:
<before/after>与<past_participle>表示执行时机,例如before <something> was run、after <something> was created;<noun>表示被注入行为的组件,例如catalog、node、pipeline、dataset、context。
错误类 Hook 遵循on_<noun>_error约定,<noun>表示抛出错误的组件。
3.2 KedroContextSpecs:上下文生命周期
after_context_created是一次 Kedro 运行中最早触发的 Hook,在KedroContext创建完成后立即调用(见 session.py 中的调用)。规格签名:
@hook_spec def after_context_created(self, context: KedroContext) -> None:其中context是刚创建的KedroContext实例,携带credentials、config_loader、env等有用信息,适合在此完成全局初始化(如注入外部服务客户端)。
3.3 DataCatalogSpecs:数据目录生命周期
after_catalog_created在数据目录创建后触发,接收catalog以及KedroContext._create_catalog的全部入参(调用点见 context.py):
@hook_spec def after_catalog_created( self, catalog: CatalogProtocol, conf_catalog: dict[str, Any], conf_creds: dict[str, Any], parameters: dict[str, Any], save_version: str, load_versions: dict[str, str], ) -> None:catalog:已创建的目录实例;conf_catalog:用于创建目录的配置;conf_creds:用于创建目录的凭据配置;parameters:目录创建后注入的参数;save_version:目录中所有数据集save操作使用的版本号;load_versions:目录中各数据集load操作使用的版本号。
典型用途:在目录就绪后统一添加监控、记录数据集清单,或按版本号做审计。
3.4 NodeSpecs:节点生命周期
节点级 Hook 围绕单个节点的执行展开(全部调用点位于 kedro/runner/task.py)。
before_node_run在节点执行前触发,且允许通过返回值改写节点输入:
@hook_spec def before_node_run( self, node: Node, catalog: CatalogProtocol, inputs: dict[str, Any], is_async: bool, run_id: str, ) -> dict[str, Any] | None:inputs键是数据集名称,值是已加载的实际数据而非数据集实例;is_async表示节点是否以异步模式运行;run_id是本次运行的 ID;- 返回值:
None或"数据集名 -> 新值"的字典;若返回字典,Kedro 将用它更新节点输入,从而实现输入覆写(例如注入测试数据或 mock)。
after_node_run在节点执行成功后触发,额外提供outputs(输出数据字典),可用于数据质量检查、指标采集:
@hook_spec def after_node_run( self, node: Node, catalog: CatalogProtocol, inputs: dict[str, Any], outputs: dict[str, Any], is_async: bool, run_id: str, ) -> None:on_node_error在节点抛出未捕获异常时触发,其签名与before_node_run一致并额外携带error:
@hook_spec def on_node_error( self, error: Exception, node: Node, catalog: CatalogProtocol, inputs: dict[str, Any], is_async: bool, run_id: str, ) -> None:常用于失败告警、错误上报或失败现场转储。
3.5 DatasetSpecs:数据集加载/保存生命周期
数据集级 Hook 在目录对单个数据集执行load/save前后触发,签名均在 specs.py 的 DatasetSpecs 类 中定义:
@hook_spec def before_dataset_loaded(self, dataset_name: str, node: Node) -> None: @hook_spec def after_dataset_loaded(self, dataset_name: str, data: Any, node: Node) -> None: @hook_spec def before_dataset_saved(self, dataset_name: str, data: Any, node: Node) -> None: @hook_spec def after_dataset_saved(self, dataset_name: str, data: Any, node: Node) -> None:其中dataset_name为数据集名称,data为实际加载/保存的数据,node为触发该操作(或刚刚运行完)的节点。这类 Hook 适合做细粒度的数据血缘记录、敏感数据脱敏或 I/O 审计。
3.6 PipelineSpecs:流水线生命周期
流水线级 Hook 在整条流水线运行前后触发,调用点位于 session.py(ServiceSession同构,见 service_session.py)。三者共享同一份run_params字典,其完整 schema 在规格 docstring 内联定义:
{ "run_id": str, "project_path": str, "env": str, "kedro_version": str, "tags": Optional[List[str]], "from_nodes": Optional[List[str]], "to_nodes": Optional[List[str]], "node_names": Optional[List[str]], "from_inputs": Optional[List[str]], "to_outputs": Optional[List[str]], "load_versions": Optional[List[str]], "runtime_params": Optional[Dict[str, Any]], "pipeline_names": Optional[List[str]], "namespaces": Optional[List[str]], "runner": str, "only_missing_outputs": bool, }before_pipeline_run在流水线运行前触发,run_params携带本次运行的全部参数,可用于权限校验、参数审计:
@hook_spec def before_pipeline_run( self, run_params: dict[str, Any], pipeline: Pipeline, catalog: CatalogProtocol ) -> None:after_pipeline_run在流水线运行成功后触发,额外提供run_result(流水线运行输出),可用于结果落库、指标聚合:
@hook_spec def after_pipeline_run( self, run_params: dict[str, Any], run_result: dict[str, Any], pipeline: Pipeline, catalog: CatalogProtocol, ) -> None:on_pipeline_error在流水线抛出未捕获异常时触发,签名与before_pipeline_run一致并携带error:
@hook_spec def on_pipeline_error( self, error: Exception, run_params: dict[str, Any], pipeline: Pipeline, catalog: CatalogProtocol, ) -> None:3.7 CLI Hooks:命令生命周期
除上述运行期 Hook 外,Kedro 还定义了 CLI 级 Hook,在 CLI 命令执行前后触发(调用点见 kedro/framework/cli/cli.py),其规格在 kedro/framework/cli/hooks/specs.py:
@cli_hook_spec def before_command_run(self, project_metadata: ProjectMetadata, command_args: list[str]) -> None: @cli_hook_spec def after_command_run(self, project_metadata: ProjectMetadata, command_args: list[str], exit_code: int) -> None:project_metadata:Kedro 项目的元数据;command_args:本次使用的全部命令行参数(包含命令与子命令本身);exit_code:Click 应用完成后的退出码(仅after_command_run有)。
官方文档明确指出,kedro-telemetry插件正是依赖这两组 CLI Hooks 来收集 CLI 使用统计的。
四、manager:全局 Hook 管理器
manager.py 是 Hooks 的"调度中枢",提供四个核心工具函数/类。
4.1_create_hook_manager():创建全局管理器
def _create_hook_manager() -> PluginManager: manager = PluginManager(HOOK_NAMESPACE) manager.trace.root.setwriter( logger.debug if logger.getEffectiveLevel() == logging.DEBUG else None ) manager.enable_tracing() manager.add_hookspecs(NodeSpecs) manager.add_hookspecs(PipelineSpecs) manager.add_hookspecs(DataCatalogSpecs) manager.add_hookspecs(DatasetSpecs) manager.add_hookspecs(KedroContextSpecs) return manager关键点:
- 基于 pluggy 的
PluginManager,命名空间为HOOK_NAMESPACE; - 一次性注册全部 5 个规范类;
- 当项目日志级别为
DEBUG时启用 pluggy 的 tracing,将每次 Hook 的执行过程写入调试日志(便于排查,但会显著增加日志噪音、拖慢流水线;文档建议生产环境保持INFO及以上)。
KedroSession在创建时即调用该函数持有全局hook_manager(见 session.py)。
4.2_register_hooks():注册项目自定义 Hook
def _register_hooks(hook_manager: PluginManager, hooks: Iterable[Any]) -> None: for hooks_collection in hooks: if not hook_manager.is_registered(hooks_collection): if isclass(hooks_collection): raise TypeError( "KedroSession expects hooks to be registered as instances. " "Have you forgotten the `()` when registering a hook class ?" ) hook_manager.register(hooks_collection)两个值得注意的行为:
- 重复注册保护:若 Hook 已被注册则直接跳过,避免用户多次调用注册逻辑导致崩溃;
- 实例校验:Hook 必须以实例(而非类)形式注册,忘记加
()会抛出上述TypeError——这一点被 tests/framework/hooks/test_manager.py 中的test_register_hooks参数化测试明确覆盖([ExampleHook]报错、[ExampleHook()]通过)。
4.3_register_hooks_entry_points():加载插件 Hook
_PLUGIN_HOOKS = "kedro.hooks" # entry-point to load hooks from for installed plugins def _register_hooks_entry_points(hook_manager, disabled_plugins) -> None: already_registered = hook_manager.get_plugins() hook_manager.load_setuptools_entrypoints(_PLUGIN_HOOKS) disabled_plugins = set(disabled_plugins) plugininfo = hook_manager.list_plugin_distinfo() # ...根据 DISABLE_HOOKS_FOR_PLUGINS 逐一 unregister 指定插件- 通过
kedro.hooks入口点自动加载已安装插件声明的 Hook,这是 Kedro默认开启的自动发现机制; - 对
disabled_plugins中列出的插件(按dist.project_name匹配,而非入口点名),调用unregister()禁用其 Hook,并记录插件名-版本到调试日志。
4.4_NullPluginManager:空实现兜底
class _NullPluginManager: def __init__(self, *args, **kwargs): ... def __getattr__(self, name): return self def __call__(self, *args, **kwargs): ...这是一个"吞掉一切调用"的空管理器:当没有实例化任何hook_manager时,它让 runner 仍可正常运行,所有 Hook 调用被静默忽略。测试test_null_plugin_manager_returns_none_when_called验证了其调用返回None的行为。
五、编写并注册一个 Hook 实现
5.1 最小实现示例
在项目src/<package_name>/hooks.py中,用@hook_impl标记实现,方法名必须与规格同名,参数可取规格参数的子集(得益于 pluggy 的 opt-in arguments 机制,未声明参数会被省略):
# src/<package_name>/hooks.py import logging from kedro.framework.hooks import hook_impl from kedro.io import DataCatalog class DataCatalogHooks: @property def _logger(self): return logging.getLogger(__name__) @hook_impl def after_catalog_created(self, catalog: DataCatalog) -> None: self._logger.info(catalog.list())5.2 在 settings.py 中注册
在src/<package_name>/settings.py的HOOKS键下注册实现(见 kedro/framework/project/init.py 中_HOOKS校验器的定义):
# src/<package_name>/settings.py from <package_name>.hooks import ProjectHooks, DataCatalogHooks HOOKS = (ProjectHooks(), DataCatalogHooks())同一规格可以注册多个实现,按LIFO(后进先出)顺序调用:即元组中靠后的实现先执行。KedroSession创建时依次调用_register_hooks(hook_manager, settings.HOOKS)与_register_hooks_entry_points(...)(见 session.py),因此自动发现的插件 Hook 先执行,settings.py中指定的 Hook 后执行。
5.3 禁用插件的自动注册 Hook
通过settings.py的DISABLE_HOOKS_FOR_PLUGINS键(对应 project/init.py 中_DISABLE_HOOKS_FOR_PLUGINS校验器)禁用指定插件的自动 Hook:
# src/<package_name>/settings.py DISABLE_HOOKS_FOR_PLUGINS = ("<plugin_name>",)5.4 两个重要注意事项
- 不要为 Hook 参数设置默认值:由于 pluggy 传参机制,带默认值的参数会收到默认值而非 Kedro 传入的真实值。例如
def before_pipeline_run(self, run_params: dict = {})中run_params恒为空字典,必须写成无默认值的形式。 - ParallelRunner 下的行为差异:使用
ParallelRunner时,catalog、context、pipeline级 Hook 会在主进程中正常执行,但dataset与node级 Hook不会在并行 worker 进程中执行。若项目依赖这些 Hook,请改用SequentialRunner或ThreadRunner。
六、执行时间线与调用链验证
Kedro 在kedro run生命周期中按下述顺序触发 Hook,右侧为仓库中的真实调用点:
| 顺序 | Hook | 调用位置 |
|---|---|---|
| 1 | after_context_created | session.py |
| 2 | after_catalog_created | context.py |
| 3 | before_pipeline_run | session.py |
| 4 | before_dataset_loaded/after_dataset_loaded | runner/task.py |
| 5 | before_node_run | runner/task.py |
| 6 | before_dataset_saved/after_dataset_saved | runner/task.py |
| 7 | after_node_run(成功)/on_node_error(失败) | runner/task.py |
| 8 | after_pipeline_run(成功)/on_pipeline_error(失败) | session.py |
整个生命周期可通过文档配套的示意图直观理解:
关于参数匹配的完整性,仓库测试 test_manager.py 中的test_hook_manager_can_call_hooks_defined_in_specs会对全部 12 个 Hook 逐一断言其参数集合与规格一致;test_hook_args_doc_table_matches_specs更是直接解析 docs/extend/hooks/introduction.md 中的参数表格,确保文档与源码签名不漂移(issue #4564)。这说明上表所列参数与规格签名即为可依赖的"地面真相"。
七、调试与最佳实践小结
- 开启 tracing:将项目日志级别调至
DEBUG即可看到每个 Hook 的执行记录;排查完毕后恢复INFO及以上以避免性能损耗。 - 保持实现无副作用顺序依赖:除
tryfirst/trylast参数外,Hook 之间的执行顺序不保证,业务逻辑不应依赖特定次序。 - 聚合相关实现:建议用类把相关 Hook 实现分组,放入项目中的
hooks.py(模块名任意,不强制)。 - 利用返回值覆写:
before_node_run的返回值可用于替换节点输入,是实现测试替身、数据注入的官方入口。 - 插件化复用:将 Hook 打包为插件并在
pyproject.toml中声明kedro.hooks(CLI 类为kedro.cli_hooks)入口点即可自动注册;若要在settings.py之外彻底理解入口点加载与禁用逻辑,可继续阅读 manager.py 与 kedro/framework/cli/hooks/manager.py。
综上,kedro.framework.hooks通过 pluggy 将"规格声明"(specs + markers)与"注册调度"(manager)解耦:规格定义好 12 个运行期 Hook 与 2 个 CLI Hook 的完整签名与参数语义,manager 负责创建单例、校验注册、加载插件入口点并按settings.py的配置组织执行顺序,从而为生产级数据流水线提供了稳定、可插拔的扩展点。
【免费下载链接】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),仅供参考