Kedro HTTP 服务器实战指南:通过 REST 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 内置 HTTP 服务器的完整用法:从安装kedro[server]可选依赖、通过kedro server start启动服务,到调用GET /health、GET /snapshot、POST /run三大端点触发管道运行与读取项目快照,并深入剖析其基于KedroServiceSession的会话生命周期与 Runner、数据集类型等安全约束。读完本文,你将掌握如何把 Kedro 项目以最小成本暴露为 REST 服务,并通过create_http_server编程式扩展自定义端点。
服务器定位与设计边界
Kedro 内置的 HTTP 服务器让外部系统能够以 REST 方式与 Kedro 项目交互——触发管道运行、查看项目元数据等。它的核心支撑是KedroServiceSession,该会话可以在多次请求之间保持存活,避免每次请求都重新初始化项目上下文。
需要特别强调的是,这个服务器是刻意保持极简的:
- 不包含认证(authentication)与授权(authorisation);
- 不提供请求排队、异步任务执行、运行历史记录;
- 不提供按请求隔离的会话(per-request session isolation)。
因此,切勿在未添加适当安全防护的情况下将其公开暴露到公网。它本质上是"一个可以按需扩展的接口基座",而非开箱即用的生产级网关。从源码看,kedro/server/http_server.py中 FastAPI 应用在启动(lifespan)阶段调用bootstrap_project完成项目引导,并在发现settings.SESSION_CLASS不是KedroServiceSession时记录警告——说明该服务器专为服务型会话设计。
安装可选依赖
HTTP 服务器依赖 FastAPI、Pydantic、Uvicorn 等第三方包,需要额外安装:
pip install 'kedro[server]'在 入口模块 中,create_http_server采用懒加载方式导入kedro.server.http_server:如果环境中缺少fastapi,会抛出ModuleNotFoundError并提示安装kedro[server]。CLI 命令kedro server start也在运行时才导入fastapi与uvicorn,见 CLI 实现。
启动服务器
在 Kedro 项目根目录内执行:
kedro server start默认监听http://127.0.0.1:8000。默认值与相关环境变量的定义集中在 kedro/server/utils.py:DEFAULT_HOST = "127.0.0.1"、DEFAULT_HTTP_PORT = 8000。
启动选项
| 选项 | 短选项 | 默认值 | 说明 |
|---|---|---|---|
--host | -H | 127.0.0.1 | 服务器绑定主机 |
--port | -p | 8000 | 服务器绑定端口 |
--reload | — | False | 代码变更时自动重载,仅限开发环境,禁止用于生产 |
--env | -e | — | Kedro 配置环境 |
--conf-source | — | — | 自定义配置目录路径 |
示例:
# 绑定到 localhost 的 8080 端口 kedro server start --host 127.0.0.1 --port 8080 # 使用 staging 环境并开启自动重载 kedro server start --env staging --reload源码层面的细节值得注意(见 CLI 实现):
- 命令会读取
metadata.project_path并写入KEDRO_PROJECT_PATH环境变量; - 若传入了
--env/--conf-source,分别写入KEDRO_SERVER_ENV/KEDRO_SERVER_CONF_SOURCE环境变量; - 未传入时,会先清空这两个环境变量,避免复用上次运行残留的值;
--reload会触发UserWarning提醒仅限开发使用,同时将reload_dirs指向项目根目录,使 Uvicorn 只监听项目内代码变更;- 最终以
uvicorn.run("kedro.server.http_server:create_http_server", factory=True, ...)启动,即以工厂模式加载应用。
服务端环境变量的优先级
create_http_server对env与conf_source的解析遵循"编程传参 > 环境变量"的优先级(见 http_server.py):
resolved_env = env if env is not None else os.environ.get(KEDRO_SERVER_ENV) resolved_conf_source = ( conf_source if conf_source is not None else os.environ.get(KEDRO_SERVER_CONF_SOURCE) )tests/server/test_run_endpoint.py中的test_run_endpoint_factory_defaults_override_env_vars等用例明确验证了"工厂参数覆盖环境变量"这一行为。项目路径同样支持KEDRO_PROJECT_PATH环境变量,_resolve_project_path(见 utils.py)会校验路径存在性,否则抛出ServerSettingsError。
HTTP 端点详解
服务器共暴露三个内置端点,路由定义与响应模型分别位于 http_server.py 与 models.py。
GET /health
返回服务器状态与所用 Kedro 版本:
curl http://127.0.0.1:8000/health{ "status": "healthy", "kedro_version": "<installed-kedro-version>" }两点事实性说明:
kedro_version是正在运行服务器的 Kedro 包版本(源码中取自kedro.__version__),而非项目pyproject.toml中声明的版本;- 响应模型
HealthResponse严格限定为{"status", "kedro_version"}两个字段(见 models.py),不会泄露project_path等内部信息——tests/server/test_http_server.py中的test_health_endpoint_response_model_validation对此有断言。
GET /snapshot
返回项目的只读结构快照:元数据、已注册管道、Catalog 数据集与参数键名。
curl http://127.0.0.1:8000/snapshot{ "status": "success", "metadata": { "project_name": "My Project", "package_name": "my_project", "kedro_version": "1.0.0" }, "pipelines": [ { "name": "__default__", "nodes": [ { "name": "split_data_node", "func_name": "split_data", "inputs": ["example_iris_data"], "outputs": ["X_train", "X_test"], "tags": [], "namespace": null, "source": { "filepath": "src/my_project/pipelines/data_science/nodes.py", "line_start": 12, "line_end": 25 } } ], "inputs": ["example_iris_data"], "outputs": ["example_predictions"] } ], "datasets": { "example_iris_data": { "name": "example_iris_data", "type": "pandas.CSVDataset", "filepath": "data/01_raw/iris.csv" } }, "parameters": ["example_learning_rate", "example_num_train_iter"] }快照结构的程序化对应物是kedro.inspection.get_project_snapshot返回的ProjectSnapshot数据类,包含metadata、pipelines、datasets、parameters四个属性,详见 Inspect a Kedro project。
失败行为:如果快照无法构建(例如 Catalog 配置出错),响应仍返回 HTTP 200,但status变为"failure",并附error字段(含异常类型与消息),此时metadata、pipelines、datasets、parameters等数据字段缺席:
{ "status": "failure", "error": { "type": "MissingConfigException", "message": "No config files found matching the pattern(s) 'catalog*'" } }从源码看,该端点内部用try/except Exception捕获所有异常,并调用_redact_url_credentials对异常消息与堆栈做脱敏处理后再返回和记录日志,防止数据集 URL 中的凭据(如user:pass@host、签名参数)泄露。tests/server/test_run_endpoint.py中的test_execute_pipeline_failure_redacts_credentials_from_exception验证了响应与日志中均不出现敏感内容。
注意:
/snapshot使用的是服务器启动时配置的环境与配置源(--env/KEDRO_SERVER_ENV、--conf-source/KEDRO_SERVER_CONF_SOURCE),不接受按请求传入env或conf_source参数。若需在同一进程内对比多个环境的快照,应使用程序化 API(见 Inspect a Kedro project)。
POST /run
触发管道运行。所有字段均为可选;发送空 JSON 对象{}即使用默认设置运行默认管道。
运行默认管道:
curl -X POST http://127.0.0.1:8000/run \ -H "Content-Type: application/json" \ -d '{}'运行指定管道并携带运行时参数:
curl -X POST http://127.0.0.1:8000/run \ -H "Content-Type: application/json" \ -d '{"pipeline_names": ["training"], "params": {"n_splits": 5}}'请求字段一览(对应RunRequestPydantic 模型,见 models.py):
| 字段 | 类型 | 说明 |
|---|---|---|
from_inputs | list[str] | 从这些数据集名称开始运行管道 |
to_outputs | list[str] | 在这些数据集名称处结束管道 |
from_nodes | list[str] | 从这些节点名称开始运行管道 |
to_nodes | list[str] | 在这些节点名称处结束管道 |
node_names | list[str] | 仅运行指定节点 |
runner | str | Runner 类名或完整点分路径,必须是kedro.runner.AbstractRunner子类(默认SequentialRunner) |
is_async | bool | 使用线程异步加载/保存节点输入输出(默认false) |
tags | list[str] | 仅运行带这些标签的节点 |
load_versions | dict[str, str] | 固定加载的数据集版本,格式{"dataset_name": "version"} |
pipeline_names | list[str] | 要运行的管道(省略时运行默认管道) |
namespaces | list[str] | 仅运行这些命名空间内的节点 |
params | dict | 传给上下文的运行时参数 |
only_missing_outputs | bool | 跳过输出已存在且已持久化的节点 |
成功响应包含run_id、status、duration_ms:
{ "status": "success", "run_id": "2024-01-01T00.00.00.000Z", "duration_ms": 142.3 }失败响应额外包含error对象(异常类型与消息):
{ "status": "failure", "run_id": "2024-01-01T00.00.00.000Z", "duration_ms": 12.1, "error": { "type": "DatasetError", "message": "Failed to load dataset 'raw_data'" } }注意:
RunRequest使用严格校验(Pydantic 的extra="forbid"),请求中出现未知字段会直接报错,而不是被静默忽略。
会话生命周期与并发模型
- 第一个
/run请求会创建KedroServiceSession,后续请求复用该会话(见 http_server.py,会话创建受threading.Lock保护,避免竞态); - 端点运行在线程池中,因此并发的
/run请求共享同一个会话,管道运行之间不互相隔离; - 服务器以
serving_mode=True创建会话。在KedroServiceSession中,serving_mode会在会话创建时通过_preload_pipelines()预先加载全部已注册管道(见 service_session.py),这样并发请求查找管道时只是读取已填充的共享单例,不会通过set_requested()写操作改变共享状态,从而避免竞态; env和conf_source不接受按请求传入,只能在服务器启动时通过--env/--conf-source指定(RunRequest模型注释中明确说明这一点)。
Runner 安全
runner字段是攻击面,Kedro 做了两层防护(实现于 http_server.py):
- 格式校验:
RunRequest用正则^[A-Za-z_][A-Za-z0-9_]*(\.[A-Za-z_][A-Za-z0-9_]*)*$拒绝任何非法点分字符串(tests/server/test_run_endpoint.py中os; import sys、../../etc/passwd、__import__('os')等恶意输入均被拒绝); - 模块白名单:短名称(如
SequentialRunner)总是解析到kedro.runner;全限定名称(如mypackage.runners.MyRunner)的模块前缀必须属于kedro.runner、项目自身包名,或在settings.py的RUNNER_MODULE_ALLOWLIST中,否则根本不会执行导入:
# settings.py RUNNER_MODULE_ALLOWLIST = ["external_lib.runners"]白名单默认值为空元组(见 kedro/framework/project/init.py)。此外,通过load_obj加载后还会校验目标必须是AbstractRunner的子类(tests中_NotARunner、函数等非类对象均被拒绝)。
数据集类型安全
Catalog 条目中的type字段决定 Kedro 导入并实例化哪个数据集类,因此HTTP 请求的params绝不能有能力选中它。服务器会拒绝任何通过runtime_params解析(${runtime_params:...})得到请求方提供值的 catalogtype,无论嵌套多深,并返回InterpolationResolutionError。
# 通过 HTTP 会被拒绝:`type` 由 runtime_params 解析 companies: type: "${runtime_params:dataset.type}" filepath: data/01_raw/companies.csv # 可接受:`type` 固定,仅其他字段(如 filepath)使用 runtime_params companies: type: pandas.CSVDataset filepath: "${runtime_params:folder, 'data/01_raw'}/companies.csv"该限制仅作用于 HTTP 请求。在受信任的非服务器场景(如kedro run --params)中,仍支持用runtime_params选择数据集type。底层实现位于 service_session.py:serving_mode下构造配置加载器时强制设置restrict_runtime_params_type_selection=True,且该参数不受项目自身CONFIG_LOADER_ARGS配置削弱——这是对不可信调用方(HTTP 请求体)的强制安全约束;相关逻辑见 omegaconf_config.py 与 templating 文档。
交互式 API 文档
服务器运行时,FastAPI 会自动生成交互式 API 文档,访问http://127.0.0.1:8000/docs即可查看所有端点、请求与响应 schema,并直接在浏览器中试调。对于调试阶段排查请求体格式非常实用。
编程方式使用 create_http_server
除了 CLI,也可以直接创建 FastAPI 应用并自行托管:
from kedro.server import create_http_server app = create_http_server( project_path="/path/to/project", env="prod", ) # Serve with uvicorn import uvicorn uvicorn.run(app, host="127.0.0.1", port=8000)关键约定:
project_path缺省时,从KEDRO_PROJECT_PATH环境变量解析(tests/server/test_run_endpoint.py验证了传参、环境变量、默认值三者间的优先级关系);env与conf_source可在create_http_server参数中显式传入,也可通过KEDRO_SERVER_ENV、KEDRO_SERVER_CONF_SOURCE环境变量提供,函数参数优先于环境变量;- 返回的
app是标准 FastAPI 应用,启动时(lifespan)自动执行bootstrap_project,关闭时若已创建会话则调用session.close()(见 tests/server/test_http_server.py 的test_lifespan_closes_session_on_shutdown)。
扩展服务器:自定义端点
由于create_http_server返回标准 FastAPI 应用,可以直接在它之上挂载额外路由或中间件,并自动继承同一套会话生命周期。
例如,暴露已注册管道列表:
from kedro.framework.project import pipelines from kedro.server import create_http_server app = create_http_server(project_path="/path/to/project") @app.get("/pipelines") def list_pipelines() -> dict: return {"pipelines": list(pipelines.keys())}新增的/pipelines端点与内置的/health、/snapshot、/run路由并存,并共享相同的KedroServiceSession生命周期。类似的扩展模式还适用于:在请求前后追加认证中间件、为敏感端点补充权限校验、或增加指标采集路由——这正好呼应了文档开头"接口可以按需扩展"的设计初衷。相关端点的路由定义与响应模型可在 http_server.py 与 models.py 中查看,测试用例集中在 tests/server 目录,可作为二次开发时的参考基准。
【免费下载链接】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),仅供参考