news 2026/9/15 13:14:09

Kedro HTTP 服务器实战指南:通过 REST API 触发管道运行与项目检查

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Kedro HTTP 服务器实战指南:通过 REST API 触发管道运行与项目检查

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 /healthGET /snapshotPOST /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也在运行时才导入fastapiuvicorn,见 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-H127.0.0.1服务器绑定主机
--port-p8000服务器绑定端口
--reloadFalse代码变更时自动重载,仅限开发环境,禁止用于生产
--env-eKedro 配置环境
--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_serverenvconf_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数据类,包含metadatapipelinesdatasetsparameters四个属性,详见 Inspect a Kedro project。

失败行为:如果快照无法构建(例如 Catalog 配置出错),响应仍返回 HTTP 200,但status变为"failure",并附error字段(含异常类型与消息),此时metadatapipelinesdatasetsparameters等数据字段缺席:

{ "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),不接受按请求传入envconf_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_inputslist[str]从这些数据集名称开始运行管道
to_outputslist[str]在这些数据集名称处结束管道
from_nodeslist[str]从这些节点名称开始运行管道
to_nodeslist[str]在这些节点名称处结束管道
node_nameslist[str]仅运行指定节点
runnerstrRunner 类名或完整点分路径,必须是kedro.runner.AbstractRunner子类(默认SequentialRunner
is_asyncbool使用线程异步加载/保存节点输入输出(默认false
tagslist[str]仅运行带这些标签的节点
load_versionsdict[str, str]固定加载的数据集版本,格式{"dataset_name": "version"}
pipeline_nameslist[str]要运行的管道(省略时运行默认管道)
namespaceslist[str]仅运行这些命名空间内的节点
paramsdict传给上下文的运行时参数
only_missing_outputsbool跳过输出已存在且已持久化的节点

成功响应包含run_idstatusduration_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()写操作改变共享状态,从而避免竞态;
  • envconf_source不接受按请求传入,只能在服务器启动时通过--env/--conf-source指定(RunRequest模型注释中明确说明这一点)。
Runner 安全

runner字段是攻击面,Kedro 做了两层防护(实现于 http_server.py):

  1. 格式校验RunRequest用正则^[A-Za-z_][A-Za-z0-9_]*(\.[A-Za-z_][A-Za-z0-9_]*)*$拒绝任何非法点分字符串(tests/server/test_run_endpoint.pyos; import sys../../etc/passwd__import__('os')等恶意输入均被拒绝);
  2. 模块白名单:短名称(如SequentialRunner)总是解析到kedro.runner;全限定名称(如mypackage.runners.MyRunner)的模块前缀必须属于kedro.runner、项目自身包名,或在settings.pyRUNNER_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验证了传参、环境变量、默认值三者间的优先级关系);
  • envconf_source可在create_http_server参数中显式传入,也可通过KEDRO_SERVER_ENVKEDRO_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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/15 13:13:50

基于YOLOv11的智能抽烟行为监测系统开发实践

1. 项目概述&#xff1a;基于YOLOv11的智能抽烟行为监测系统这个项目实现了一套完整的端到端抽烟行为识别解决方案&#xff0c;从数据采集到GUI界面部署的全流程覆盖。核心采用YOLOv11目标检测算法&#xff0c;针对抽烟这一特定行为进行优化&#xff0c;最终封装成可视化管理界…

作者头像 李华
网站建设 2026/9/15 13:12:34

工业级OpenCV形状检测:从光照噪声到PLC可用的鲁棒实现

1. 这不是“画个圈就识别”的玩具功能&#xff0c;而是工业视觉的底层呼吸OpenCV形状检测——这五个字在新手教程里常被简化成“用cv2.findContours()找轮廓&#xff0c;再用cv2.approxPolyDP()拟合多边形”&#xff0c;然后贴出一张带红框的硬币、三角板和矩形纸片截图。但我在…

作者头像 李华
网站建设 2026/9/15 13:12:32

用fairseq从零训练中英NMT模型:数据清洗到参数调优全流程

从数据集清洗、BPE切分、环境配置到训练参数调优&#xff0c;完整走一遍用fairseq训练中英NMT模型的流程&#xff0c;我把过程中踩过的坑和最终跑通的配置都放在下面了。如果你正准备复现一篇翻译论文&#xff0c;或者想自己训一个离线可部署的中英翻译基线&#xff0c;这篇应该…

作者头像 李华
网站建设 2026/9/15 13:11:45

高危端口自查与加固:从80到6379的端口安全实践指南

几年前的一次应急响应&#xff0c;让我对“高危端口”这四个字有了非常直观的认知。客户反馈一台业务服务器CPU被打满、对外连接异常&#xff0c;登录上去一看&#xff0c;一个陌生的进程占了大半资源&#xff0c;顺着网络连接排查才发现&#xff0c;入口竟然是Redis的6379端口…

作者头像 李华
网站建设 2026/9/15 13:11:19

1万本金期货量化:双均线+ATR趋势跟踪策略全解析

先说结论&#xff1a;1万块本金做期货量化&#xff0c;年化17%绝对不是一个激进的目标&#xff0c;但它比大多数人想象得要难。难的不是找到能赚钱的策略&#xff0c;而是让策略在实盘中活下来。这篇文章我会完整拆解一套我实际跑过、基于双均线ATR吊灯止损的趋势跟踪策略&…

作者头像 李华
网站建设 2026/9/15 13:10:59

API越权漏洞自动化检测:Hadrian+Vespasian+crAPI本地部署实战

API 越权漏洞自动化检测是我最近反复折腾的一个方向。越权漏洞说起来简单&#xff0c;但真要在几十个接口里找出“哪个接口能看别人数据、哪个接口能调管理员功能”&#xff0c;手工点一天也未必能覆盖完整。我最后搭了一套本地组合&#xff1a;Hadrian 负责扫描编排&#xff0…

作者头像 李华