Kedro 会话生命周期管理:KedroSession与KedroServiceSession实战指南
【免费下载链接】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 的会话(Session)机制负责管理一次或多次 Kedro 运行的完整生命周期,将库组件(由KedroContext管理)与运行时的静态、动态数据解耦。本文以 docs/extend/session.md 为骨架,结合kedro/framework/session/源码与测试,讲解KedroSession(单次运行)与KedroServiceSession(服务化多运行)的创建参数、调用方式,以及bootstrap_project/configure_project的项目初始化流程。读完本文,你将能在脚本、Notebook 或 Web 服务中正确创建会话并按需执行管道。
概览:会话为什么存在
一次 Kedro 运行涉及两类截然不同的对象:
- 库组件:由
KedroContext统一管理,负责配置加载、Catalog 构建、钩子调度等; - 会话数据:包括静态数据(项目路径、环境名、运行时参数)与动态数据(CLI 上下文、Git 信息、用户名、异常信息等)。
会话(Session)将二者解耦,使 Kedro 组件与插件可以无需导入KedroContext对象即可访问会话数据。Kedro 提供两个基于抽象基类AbstractSession(定义于 kedro/framework/session/abstract_session.py)实现的会话类:
| 特性 | KedroSession | KedroServiceSession |
|---|---|---|
| 定位 | 单次运行(single-run) | 多运行(multi-run)服务场景 |
| 生命周期 | 一次会话对应一次管道运行 | 一个会话内可多次运行管道,每次可传不同运行时参数 |
| 典型场景 | CLIkedro run、脚本化单次执行 | Web 服务 / API 中保持会话存活、按需触发运行 |
| 状态 | 持久化运行时参数与对应 session ID | 不落盘 session 数据,close()仅记录日志 |
KedroSession会捕获会话元数据,包括CLI 上下文、环境信息、运行时参数、用户名以及项目 Git 状态;而KedroServiceSession处于活跃开发中,可能偶发破坏性变更,官方鼓励试用并反馈意见(该提示同时出现在 service_session.py 的类 docstring 中)。
两个类共有的核心方法:
create():携带会话数据创建会话实例;load_context():实例化KedroContext。对于KedroServiceSession,此方法接受可选参数runtime_params,用于更新该次运行的KedroContext参数;close():关闭当前会话。KedroSession在save_on_close=True时还会把会话数据保存到磁盘;run():以给定参数运行管道,详见 Running pipelines。
在 abstract_session.py 中,AbstractSession实现了上下文管理器协议:__enter__返回自身,__exit__调用close(),这就是两个会话类都能直接配合with语句使用的原因。
创建KedroSession:单次运行
以下代码以上下文管理器方式创建KedroSession并在上下文内运行管道。脚本可以从 Kedro 项目中的任意位置调用——find_kedro_project会从当前目录向上逐级寻找含pyproject.toml且带[tool.kedro]配置的项目根目录,因此无需硬编码路径;退出with块时会话自动关闭:
from pathlib import Path from kedro.framework.session import KedroSession from kedro.framework.startup import bootstrap_project from kedro.utils import find_kedro_project # Get project root current_dir = Path(__file__).resolve().parent project_root = find_kedro_project(current_dir) bootstrap_project(Path(project_root)) # Create and use the session with KedroSession.create(project_path=project_root) as session: session.run()KedroSession.create()的可选参数如下:
| 参数 | 类型/默认值 | 说明 |
|---|---|---|
project_path | Path \| str \| None | 项目根目录路径;默认取Path.cwd()(源码见 session.py) |
save_on_close | bool,默认True(create()内默认) | 会话关闭时是否将 session 数据保存到磁盘 |
env | str \| None | KedroContext使用的环境名;不传时读取环境变量KEDRO_ENV |
runtime_params | dict \| None | 传递给底层KedroContext的运行时项目参数;指定后会更新并优先于项目配置中读取的参数 |
conf_source | str \| None | KedroContext的配置来源目录;默认取settings.CONF_SOURCE(即项目的conf/) |
从源码看,create()内部流程如下(session.py):
- 调用
validate_settings()校验项目配置就绪; - 以
generate_timestamp()生成时间戳作为 session ID,构造会话并初始化 session store; - 组装
session_data:project_path、session_id、CLI 上下文(_jsonify_cli_context)、env、runtime_params、当前用户名(getpass.getuser())、项目 Git 状态(_describe_git,含commit_sha与dirty标志); - 将全部数据写入
session._store。
_describe_git与_jsonify_cli_context的实现可分别参见 session.py 与 session.py,它们是“捕获会话元数据”这一能力的落地。
会话数据持久化与 session store
KedroSession的close()在save_on_close=True时调用self._store.save()(session.py)。store 由_init_store()根据settings.SESSION_STORE_CLASS实例化(默认是 store.py 中的BaseSessionStore),默认存储路径为<project_root>/sessions。注意BaseSessionStore本身是不落盘的临时实现——read()返回空 dict、save()仅记录调试日志,真正持久化由自定义 store 子类完成,可通过项目settings.py中的SESSION_STORE_CLASS与SESSION_STORE_ARGS替换(对应 settings 定义见 kedro/framework/project/init.py)。
一次会话只允许一次运行
KedroSession.run()内部通过self._run_called标志强制“会话与运行 1:1 映射”——若同一会话被调用第二次,会抛出KedroSessionError(session.py)。这正是需要KedroServiceSession的原因:服务化场景必须能在同一会话中多次运行。
创建KedroServiceSession:服务化多运行
以下代码创建KedroServiceSession,在同一会话中先后以不同运行时参数运行两次管道,最后手动关闭会话:
from pathlib import Path from kedro.framework.session import KedroServiceSession from kedro.framework.startup import bootstrap_project from kedro.utils import find_kedro_project # Get project root current_dir = Path(__file__).resolve().parent project_root = find_kedro_project(current_dir) bootstrap_project(Path(project_root)) # Create and use the session session = KedroServiceSession.create(project_path=project_root) # first run session.run(runtime_params={"param1": "value1"}) # second run with different runtime parameters session.run(runtime_params={"param1": "value2"}) # close the session when done session.close()KedroServiceSession.create()的可选参数:
| 参数 | 类型/默认值 | 说明 |
|---|---|---|
session_id | str \| None | 会话标识符;不传时自动生成 UUID(str(uuid.uuid4()),见 service_session.py) |
project_path | Path \| str \| None | 项目根目录路径 |
env | str \| None | KedroContext使用的环境名;不传时读取KEDRO_ENV |
conf_source | str \| None | 配置来源目录,默认conf/ |
serving_mode | bool,默认False | 会话是否处理并发run()调用;为True时在create()阶段急切预加载全部管道,以避免并发竞争条件 |
create()与KedroSession的差异
KedroServiceSession没有save_on_close参数——服务会话不落盘,其close()仅输出关闭日志(service_session.py);KedroServiceSession没有create()级runtime_params参数——运行时参数改由每次run()传入,从而为每一次具体运行独立更新KedroContext参数。与之对应,run()也接受独立的run_id参数(默认用generate_timestamp()生成),每次运行都有独立标识。
KedroServiceSession.load_context(runtime_params=None)会构造一个全新的 config loader 并注入该次运行的参数(service_session.py);测试 tests/framework/session/test_service_session.py 验证了load_context(runtime_params={"param1": "value1"})后context.config_loader.runtime_params即为该值。
serving_mode:并发安全的服务化运行
serving_mode=True的语义在源码中有精细实现:
_enable_serving_mode()先调用_preload_pipelines()成功之后才把_serving_mode置位,保证预加载失败时会话仍停留在普通 CLI 模式而非不一致状态(service_session.py);- 预加载时通过
pipelines.set_requested(None)+list(pipelines)将pipeline_registry.py中注册的全部管道填入共享单例pipelines._content(service_session.py); - 运行阶段,serving 模式跳过
set_requested(),避免并发请求间互相改写共享管道注册表;而普通模式下set_requested()会告诉懒加载器只导入需要的管道,加快 CLI 启动速度(service_session.py); - 安全加固:serving 模式下 config loader 被强制设置
restrict_runtime_params_type_selection=True——因为运行时参数可能来自不可信的 HTTP 请求体,项目自身的CONFIG_LOADER_ARGS也无法削弱该限制(service_session.py,对应 issue kedro-org/kedro#5706)。
这些行为均有测试覆盖:test_serving_mode_preloads_all_pipelines、test_serving_mode_flag_stays_false_when_preload_fails 以及用 20 线程ThreadPoolExecutor压测并发run()的 test_serving_mode_concurrent_runs_do_not_raise。
bootstrap_project与configure_project:运行前的项目初始化
Kedro 的 CLI 在kedro run启动时已自动执行这两个函数,因此日常命令行使用无需手动调用;只有以编程方式(脚本、服务、Notebook)与 Kedro 项目交互时才需要。若想在 Notebook 等交互式环境中加载 Kedro 项目,也可以直接使用%reload_kedroline magic(参见 Kedro and Notebooks)。
整体启动流程如下:
两个函数都负责 Kedro 项目的初始化,区别在于:
bootstrap_project(项目模式):内部调用configure_project,并额外读取pyproject.toml的[tool.kedro]段获得package_name、project_name、kedro_init_version等元数据,把项目源码目录(默认src/)加入sys.path并写入PYTHONPATH,使项目可作为 Python 包被导入。实现见 startup.py 与_add_src_to_path(startup.py)。适合直接操作项目源码的开发场景。configure_project(打包模式):读取项目的settings.py与pipeline_registry.py,在 Kedro 运行前注册配置与管道,并记录全局PACKAGE_NAME(kedro/framework/project/init.py)。如果你的 Kedro 项目是打包安装的,执行管道前需调用configure_project。
ValueError: package name not found
ValueError: Package name not found. Make sure you have configured the project using `bootstrap_project`. This should happen automatically if you are using Kedro command line interface.该错误由validate_settings()抛出(kedro/framework/project/init.py):当全局PACKAGE_NAME为None(即项目尚未被bootstrap_project/configure_project配置)时即触发。使用 CLI 时无需担心,但使用multiprocessing时需特别注意。
取决于操作系统,Python 有不同的多进程启动方式(fork/spawn/forkserver)。若进程以spawn方式启动,Python 会在每个子进程中重新导入所有模块,此时必须在新进程启动处再次调用configure_project。Kedro 在ParallelRunner中正是这样处理的,例如 kedro/runner/task.py:
if ( multiprocessing.get_start_method() in ("spawn", "forkserver") and package_name ): Task._bootstrap_subprocess(package_name, logging_config)_bootstrap_subprocess在子进程中重新执行configure_project(package_name)并(可选)重放日志配置(task.py),从而保证子进程也能正确解析项目包。关于各启动方式的取舍,可参考 Run a pipeline 中对fork/forkserver/spawn的对比说明。
实战选择建议
- CLI / 批处理 / 单次脚本:默认使用
kedro run,由 CLI 自动完成初始化与KedroSession生命周期管理;需要编程式控制时用KedroSession.create()+with上下文管理器,并注意一次会话仅一次run()。 - Web 服务 / API / 按需触发:使用
KedroServiceSession,在同一会话中多次run()并每次注入不同runtime_params;若服务会并发处理请求,务必设置serving_mode=True让管道在创建期一次性预加载,同时获得运行时参数驱动 Catalog 类型选择的强制安全限制。 - 多进程/分布式:凡使用
spawn启动方式,子进程中必须重新调用configure_project(或bootstrap_project),否则会遇到 “Package name not found” 错误。
【免费下载链接】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),仅供参考