
Kedro HTTP Server 完全指南kedro.server 模块与 create_http_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/kedroKedro 在kedro.server模块中提供了一个基于 FastAPI 的可选 HTTP 服务器让外部系统可以通过 REST 接口触发管道运行、查看项目元数据与目录快照。本文围绕 API 文档中公开的唯一入口create_http_server展开结合 HTTP 服务器实现、请求/响应模型、CLI 命令 与端到端测试用例完整讲解服务器的启动方式、三大端点、请求参数、会话生命周期与内置安全机制读完即可在自己的 Kedro 项目上独立部署并二次扩展。模块概览kedro.server 包含什么kedro.server是一个可选模块依赖额外的server扩展包才能使用。在 kedro/server/init.py 中模块公开的唯一 API 是create_http_server工厂函数其完整构成如下模块文件职责kedro/server/init.py对外入口懒加载真正的实现kedro/server/http_server.pyFastAPI 应用工厂、/health、/snapshot、/run路由与管道执行逻辑kedro/server/models.pyPydantic 请求/响应模型RunRequest、RunResponse、SnapshotResponse、HealthResponsekedro/server/utils.py服务器环境变量、默认主机端口与项目路径解析kedro/framework/cli/server.pykedro server startCLI 命令模块级create_http_server采用懒导入策略kedro/server/init.py#L16-L30只有在真正调用时才会导入 FastAPI 依赖。如果环境中缺少fastapi会抛出带有明确提示的ModuleNotFoundError引导用户执行pip install kedro[server]安装依赖后即可从项目内部启动服务器。核心工厂函数 create_http_servercreate_http_server是构建整个 HTTP 服务器的唯一入口签名定义如下kedro/server/http_server.py#L71-L87def create_http_server( project_path: str | Path | None None, env: str | None None, conf_source: str | None None, ) - FastAPI:三个参数全部可选其解析优先级遵循编程传参 环境变量的原则参数作用对应的环境变量默认行为project_path要服务的 Kedro 项目根目录KEDRO_PROJECT_PATH未传时从环境变量解析两者皆缺则抛出ServerSettingsErrorenvKedro 配置环境如base、local、prodKEDRO_SERVER_ENV未传时读环境变量否则为Noneconf_source配置文件目录路径KEDRO_SERVER_CONF_SOURCE未传时读环境变量否则为None在 tests/server/test_http_server.py 中可以验证这一优先级即使设置了KEDRO_SERVER_ENVfrom_env_var传入envfrom_argument后app.state.default_env仍为from_argument。create_http_server内部完成四件事解析配置合并函数参数与环境变量得到resolved_env与resolved_conf_source并存入app.statekedro/server/http_server.py#L89-L94解析项目路径调用_resolve_project_path校验路径存在且为目录kedro/server/utils.py#L22-L56结果写入app.state.project_path创建 FastAPI 应用设置标题Kedro Server、版本号取运行时 Kedro 包版本而非项目 pyproject 声明的版本、并挂载lifespan生命周期管理器kedro/server/http_server.py#L96-L105注册路由依次注册GET /health、GET /snapshot、POST /run三个内置端点。工厂函数还创建了一个threading.Lock()存入app.state.session_lock用于保护后续会话的创建与关闭。生命周期启动时引导项目FastAPI 的lifespan上下文管理器kedro/server/http_server.py#L47-L68在应用启动时调用bootstrap_project(project_path)完成项目引导并把ProjectMetadata存入app.state.metadata供/snapshot端点复用。同时它会检查settings.SESSION_CLASS如果项目在settings.py中自定义了SESSION_CLASS服务器会打印警告——HTTP 服务器专为KedroServiceSession设计不会使用自定义会话类应用关闭时在锁的保护下关闭已创建的会话。test_lifespan_calls_bootstrap_on_startup 与 test_lifespan_closes_session_on_shutdown 分别验证了引导与关停行为。三个内置 HTTP 端点GET /health健康检查最简单的端点kedro/server/http_server.py#L107-L116返回服务器状态与 Kedro 运行版本curl http://127.0.0.1:8000/health{ status: healthy, kedro_version: installed-kedro-version }注意kedro_version是正在运行服务器的 Kedro 包版本不是项目pyproject.toml中声明的版本。响应严格受 HealthResponse 模型约束仅含status与kedro_version两个字段不会泄露项目路径等内部信息见 test_health_endpoint_response_model_validation。GET /snapshot只读项目快照返回项目结构的只读快照元数据、已注册管道、目录数据集与参数键kedro/server/http_server.py#L118-L160curl 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] }关键行为快照通过get_project_snapshot(env..., conf_source..., metadata...)构建使用服务器启动时配置的环境与配置源KEDRO_SERVER_ENV/KEDRO_SERVER_CONF_SOURCE不接受按请求传入的env/conf_source快照构建失败时如目录配置缺失HTTP 状态码仍为 200但status变为failure且metadata、pipelines、datasets、parameters字段不出现只保留结构化错误{ status: failure, error: { type: MissingConfigException, message: No config files found matching the pattern(s) catalog* } }数据集filepath与异常消息在返回和写日志前都会经过_redact_url_credentials脱敏处理防止 URL 中的账号密码、签名查询参数如X-Amz-Signature和 fragment 泄露。相关行为在 tests/server/test_snapshot_endpoint.py 中有完整覆盖。POST /run触发管道运行核心执行端点kedro/server/http_server.py#L162-L185。所有字段均可选发送空 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}}成功响应{ status: success, run_id: 2024-01-01T00.00.00.000Z, duration_ms: 142.3 }失败响应多出一个error对象异常类型名与脱敏后的消息HTTP 状态码同样为 200{ status: failure, run_id: 2024-01-01T00.00.00.000Z, duration_ms: 12.1, error: { type: DatasetError, message: Failed to load dataset raw_data } }成功与失败分别由RunSuccess与RunFailure建模并通过Field(discriminatorstatus)联合为 RunResponse。测试 test_run_endpoint_returns_error_detail_on_failure 验证了失败响应中不会包含 traceback。RunRequest请求体字段全解请求体由 RunRequest 模型严格校验。下表汇总全部字段及其含义字段类型默认值说明pipeline_nameslist[str]None要运行的已注册管道名列表缺省运行默认管道paramsdictNone传入上下文的运行时参数runnerstrNone默认SequentialRunnerRunner 类名或完整点分路径必须是kedro.runner.AbstractRunner的子类is_asyncboolFalse是否用线程异步加载/保存节点输入输出tagslist[str]None只运行带这些标签的节点node_nameslist[str]None只运行指定名称的节点from_nodeslist[str]None从这些节点开始to_nodeslist[str]None到这些节点结束from_inputslist[str]None从这些数据集开始to_outputslist[str]None到这些数据集结束load_versionsdict[str, str]None指定加载的数据集版本如{dataset1: 2024-01-01}namespaceslist[str]None只运行这些命名空间内的节点only_missing_outputsboolFalse只运行输出缺失的节点跳过已持久化输出的节点该模型的校验规则值得注意禁止多余字段model_config ConfigDict(extraforbid)env与conf_source不被允许出现在请求体中因为它们在服务器启动时统一配置防止每个请求各自指定环境导致混乱runner 字段正则校验必须匹配合法 Python 点分标识符kedro/server/models.py#L19如SequentialRunner、kedro.runner.ParallelRunner、mypackage.runners.MyRunner。类似os; import sys、../../etc/passwd、含空格的字符串会被直接拒绝见 tests/server/test_run_endpoint.py。会话管理与执行流程/run的会话管理遵循首次创建、后续复用的模型第一个/run请求到达时在session_lock保护下通过KedroServiceSession.create(..., serving_modeTrue)创建会话kedro/server/http_server.py#L176-L183serving_modeTrue会预先急切加载全部已注册管道见 kedro/framework/session/service_session.py#L121-L145保证并发请求读取共享pipelines单例时不产生共享状态变更从而保证多线程查找安全后续请求直接复用同一会话_execute_pipeline负责解析 runner、调用session.run()并返回结构化结果kedro/server/http_server.py#L207-L296。执行流程会把RunRequest的各字段映射为session.run()的参数列表型参数如tags、node_names转为元组params映射为runtime_params。test_run_endpoint_reuses_service_session 验证了两次请求只创建一次会话且_execute_pipeline被调用两次test_run_endpoint_with_all_parameters 验证了全部字段的透传。同时注意服务模式下配置加载器会被强制设置restrict_runtime_params_type_selectionTruekedro/framework/session/service_session.py#L154-L179即禁止运行时参数动态决定 catalog 数据集的type字段防止不可信的 HTTP 调用方通过params让 Kedro 导入任意类。这一限制是项目配置无法解除的仅针对 HTTP 请求生效kedro run --params等受信任场景不受影响。通过 CLI 启动kedro server start在 Kedro 项目目录内执行kedro server start默认监听http://127.0.0.1:8000。完整选项定义于 kedro/framework/cli/server.py选项简写默认值说明--host-H127.0.0.1绑定主机--port-p8000绑定端口--reloadFalse代码变更自动重载仅限开发环境生产环境禁用--env-e无Kedro 配置环境--conf-source无自定义配置目录路径常用示例# 绑定到 8080 端口 kedro server start --host 127.0.0.1 --port 8080 # 使用 staging 环境并开启自动重载仅开发 kedro server start --env staging --reloadCLI 的底层实现逻辑是校验fastapi/uvicorn依赖 → 将project_path写入KEDRO_PROJECT_PATH环境变量若指定了--env/--conf-source则同步写入KEDRO_SERVER_ENV/KEDRO_SERVER_CONF_SOURCE→ 调用uvicorn.run(kedro.server.http_server:create_http_server, factoryTrue, ...)。因此kedro server start --env staging与create_http_server(envstaging)在效果上等价。服务器运行后FastAPI 会在http://127.0.0.1:8000/docs自动生成可交互的 API 文档页面列出全部端点与请求/响应 Schema可直接在浏览器中试运行。程序化使用与二次扩展create_http_server返回标准 FastAPI 应用因此既可以配合 uvicorn 直接运行也可以自由挂载自定义路由与中间件。from kedro.server import create_http_server import uvicorn app create_http_server( project_path/path/to/project, envprod, ) uvicorn.run(app, host127.0.0.1, port8000)添加自定义端点的示例——比如暴露已注册管道列表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共享同一套会话生命周期。安全边界与部署注意事项Kedro 官方对 HTTP 服务器的定位是刻意保持极简的可扩展接口它不内置认证、授权、请求排队、异步任务执行、运行历史或按请求的会话隔离。因此切勿不加防护地将其暴露在公网必须自行叠加安全控制。内置的安全机制包括三层Runner 模块白名单简短的 runner 名如SequentialRunner始终从kedro.runner解析完整点分路径则必须属于kedro.runner、项目自身的 Python 包或 settings.py 中RUNNER_MODULE_ALLOWLIST列出的模块kedro/server/http_server.py#L190-L204否则模块根本不会被导入。允许外部 runner 时需显式配置# settings.py RUNNER_MODULE_ALLOWLIST [external_lib.runners]AbstractRunner 子类校验即便模块通过白名单加载出的对象仍必须是AbstractRunner的子类且是类对象否则直接拒绝执行对应测试见 tests/server/test_run_endpoint.py凭据与签名脱敏所有异常消息、traceback 与数据集filepath在写日志和返回前都经过_redact_url_credentials清洗避免数据集 URL 中的账号密码、AWS 签名参数与 fragment 泄露见 tests/server/test_run_endpoint.py。这三层防线配合RunRequest的严格校验、serving_mode的并发安全管道预加载以及运行参数禁止驱动数据集type的硬限制共同构成了 Kedro HTTP 服务器可安全暴露管道执行能力的基础。延伸阅读HTTP 服务器核心实现工厂函数、路由与执行逻辑全文请求与响应模型RunRequest全部字段与校验规则CLI 服务器命令kedro server start参数与启动流程服务会话实现KedroServiceSession与serving_mode的管道预加载机制服务器测试用例test_http_server.py、test_snapshot_endpoint.py、test_run_endpoint.py覆盖全部端点行为与安全边界用户指南Serving Kedro pipelines over HTTP端到端的实操文档【免费下载链接】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),仅供参考