资讯详情

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

📅 2026/9/15 13:14:34 | 华诺云谱 👁 阅读
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 /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.pyDEFAULT_HOST 127.0.0.1、DEFAULT_HTTP_PORT 8000。启动选项选项短选项默认值说明--host-H127.0.0.1服务器绑定主机--port-p8000服务器绑定端口--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, factoryTrue, ...)启动即以工厂模式加载应用。服务端环境变量的优先级create_http_server对env与conf_source的解析遵循编程传参 环境变量的优先级见 http_server.pyresolved_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:passhost、签名参数泄露。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_inputslist[str]从这些数据集名称开始运行管道to_outputslist[str]在这些数据集名称处结束管道from_nodeslist[str]从这些节点名称开始运行管道to_nodeslist[str]在这些节点名称处结束管道node_nameslist[str]仅运行指定节点runnerstrRunner 类名或完整点分路径必须是kedro.runner.AbstractRunner子类默认SequentialRunneris_asyncbool使用线程异步加载/保存节点输入输出默认falsetagslist[str]仅运行带这些标签的节点load_versionsdict[str, str]固定加载的数据集版本格式{dataset_name: version}pipeline_nameslist[str]要运行的管道省略时运行默认管道namespaceslist[str]仅运行这些命名空间内的节点paramsdict传给上下文的运行时参数only_missing_outputsbool跳过输出已存在且已持久化的节点成功响应包含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 的extraforbid请求中出现未知字段会直接报错而不是被静默忽略。会话生命周期与并发模型第一个/run请求会创建KedroServiceSession后续请求复用该会话见 http_server.py会话创建受threading.Lock保护避免竞态端点运行在线程池中因此并发的/run请求共享同一个会话管道运行之间不互相隔离服务器以serving_modeTrue创建会话。在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.pyserving_mode下构造配置加载器时强制设置restrict_runtime_params_type_selectionTrue且该参数不受项目自身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, envprod, ) # Serve with uvicorn import uvicorn uvicorn.run(app, host127.0.0.1, port8000)关键约定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),仅供参考
📝

华诺云谱内容团队

资深建站顾问 · 行业研究员

10年+企业数字化服务经验,专注智能建站、SEO优化与品牌营销,持续输出建站技巧、行业洞察与营销干货,已帮助5000+企业实现数字化增长。

你可能需要的服务

订阅华诺云谱资讯周报

每周一封,精选建站技巧、SEO与营销干货,直达邮箱。已有 8,000+ 企业主订阅,助你少走弯路。