在Lakeflow上启动工作
此工具是在云端生成计算作业的一种固执己见的方式。通过“计算 作业”,我的意思是大规模并行数据处理作业,比如训练深度网络, 分析S3存储桶中的大型文本语料库,或1000个 对某物的并行模拟。为了让你做这些事情,这个包 要求您将代码编写为Python包,并强制您指定 a中的包依赖关系 pyproject.toml。然后,它上传该包(作为 python wheel)来执行它。
这比Databrick内置的笔记本编辑方法更重 Python脚本在他们的web UI中。作为回报,它允许您捕获大型包裹 通过git子模块跨repos的依赖关系,并导入第三方包 通过紫外线。它比大多数其他工作提交系统重量轻,因为它 不需要构建docker容器。Docker容器占用大量资源 您的系统快照,足以构建一个完整的unix环境。这些 快照的大小约为千兆字节,很难从家里上传 电脑。对于我们的大部分工作,轮子提供了我们所需的所有集装箱化 (一个轮子是几千字节)。
它还有一个观点: uv 是捕获这些Python的好方法 依赖关系,具有a pyproject.toml我们也在探索 裤子 作为管理更复杂包的一种方式。裤子也可以出口轮子,所以这种设计中没有什么能阻止我们采用裤子。
您可以使用此工具构建您的轮子,将其上传到ViewModel,生成副本 每个都有不同的命令行参数,并跟踪作业的状态。 您还可以使用ConnectionUI检查作业的状态。工具 提供了几个接口:
- 一个MCP服务器,这样你就可以让AI为你生成作业。
- 可以从shell使用的CLI。
- 一个可以从Python程序调用的程序化Python接口。
获取对copula的访问权限
通过访问检查您是否可以访问Databrick 此url.如果 你陷入了一个无限循环,在这个循环中,copula向你发送了一个不需要的代码 工作,这意味着你没有账户。在#help数据平台中请求一个。
您的包裹结构
此包假定您要在集群上运行的包具有 像这样的结构,可以运行 uv run:
my_project/
├── pyproject.toml
├── src/
└── my_package/
├── __init__.py
└── my_package_py.py它还假设您已为您的 pyproject.toml 被称为 “湖流任务”。如果您的包裹被调用 my_package,它有一个司机 脚本已调用 my_package_py.py,并且此脚本中的main函数被调用 main,您可以这样定义“lakeflow任务”入口点:
[project.scripts]
lakeflow-task = "my_package.my_package_py:main"包裹 lakeflow_demo 在此目录下,您将看到 如何设置包的具体示例。
使用CLI构建和启动包
要在集群上运行包,首先构建轮子,然后上传 然后告诉copula运行它。
为了更容易跟踪您的工件和跑步的沿袭,构建 step将当前的git提交哈希嵌入到wheel版本中(例如。 0.1.0.devabcdef1234...).这需要对工作树进行所有更改,以 在建造之前必须先完成。否则,构建将失败并出现错误 要求你承诺或藏匿。
- 从源创建作业:
你可以使用 create-job-from-source 构建、上传和创建作业。
如果你没有通过a --cluster-id,会自动创建一个新集群:
uv run lakeflow.py create-job-from-source \
"my-lakeflow-job" \
"my-package" \
--pyproject-dir-path ~/my_project \
--max-workers 4这将返回作业ID,我们将在下一步中使用它。这还没有 运行任何作业。它只是启动一个可以运行它们的集群。这 --max-workers 参数设置自动缩放的最大工作人数 在新集群上。
要使用现有集群,请传递 --cluster-id:
uv run lakeflow.py create-job-from-source \
"my-lakeflow-job" \
"my-package" \
--pyproject-dir-path ~/my_project \
--cluster-id 0202-235755-w37hoxe8如果集群未运行,它将自动启动。
您还可以显式创建集群,并在多个作业中重用它:
uv run lakeflow.py create-cluster --max-workers 4这将返回一个可以传递给的集群ID create-job-from-source 或 create-job 通过 --cluster-id.
- 开始作业:
uv run lakeflow.py trigger-run 123456 arg11 arg12
uv run lakeflow.py trigger-run 123456 arg21 arg22
uv run lakeflow.py trigger-run 123456 arg31 arg32这将使用三组不同的作业启动三个作业实例 论据。您可以让参数引用不同的数据分片,以及 开始做尽可能多的平行工作。您的工作可以检索这些 通过argv进行参数设置。它可以从环境中检索其作业id 变量 DATABRICKS_RUN_ID.
您还可以将环境变量传递给远程作业,而不会泄漏 机密(如API密钥)通过您的命令行:
uv run lakeflow.py trigger-run 123456 arg1 arg2 \
--secret-env-var MY_SECRET_KEY --secret-env-var MY_OTHER_SECRET_KEY该工具从您的本地环境中读取值,并将其上传到 Rancher的秘密和通行证 --lakeflow-secret-scope 作为 任务的命令行参数。然后,您的任务可以使用以下命令检索机密 具有该作用域名称的Databricksdbutils API。
- 监控运行情况:
uv run lakeflow.py list-job-runs 123456这列出了给定作业ID的运行。
- 获取运行日志:
uv run lakeflow.py get-run-logs 987654321这将检索特定运行ID的日志。它接受由返回的运行 trigger-run.
使用Python编程接口
上面说明了如何使用CLI。您可能会发现更容易使用 改为使用Python编程接口访问包。看 run_lakeflow_demo.py 举个例子。
使用MCP服务器
您可以将此软件包安装为MCP服务器。为此,请将此添加到 ~/.cursor/mcp.json:
{
"mcpServers": {
"lakeflow": {
"command": "uv",
"args": [
"run",
"--quiet",
"--directory",
"/path/to/lakeflow-mcp",
"python",
"lakeflow.py"
],
"env": {
"DATABRICKS_HOST": "https://hims-machine-learning-staging-workspace.cloud.databricks.com",
"DATABRICKS_TOKEN": ""
}
},
...
}
}然后,您可以要求代理执行以下操作:
let's launch 4 copies of this job on lakeflow, and pass them the arguments "fi", "fie", "fo", and "fum" respectively.考虑的替代设计
我的目标是建立一个工作提交系统,该系统:
- Python优先:可以运行包含约100个Python文件和第三方依赖关系的Python包。
- 版本化工件:对源和输出进行版本化。
- 工作流编排:可以将其步骤分解为可以在工作流编排器下缓存、检查点、重试和恢复的任务。
- 原生能力:可以容纳少量用Rust、C++或Dafny编写的非Python代码。
- 小规模:并行运行约100名远程工作人员的作业,同时为约20名工程师运行作业。
理想的系统将使用 完美 作为工作流 编排器,在我们目前使用的现有kubernetes脚手架之上 运行staging和prod。有许多工作流编排器,但Prefect是 唯一一个提供上述所有工作流功能的。这 理想的系统是Prefect前端VM,它可以扩展kubernetes集群 按需上下。推出这个产品需要一些对话 与devops团队合作,为公司引入新的技术栈。时间 因为这很快就会到来,但这个包裹不是那样的。
与此同时,该软件包使用了数据工程团队的技术栈 已经在使用。他们已经在使用ViewModel来运行笔记本 通过使用copula的专业知识,我可以快速地使用这个解决方案。copula 笔记本是DE团队在ConnectionUI中编辑的小型python文件。这些 脚本在Git下进行版本控制。copula确实提供了一个工作流 编排器,但该团队使用Airflow来完成更大的工作。总之,技术 数据处理团队已经使用的堆栈提供了70%的功能 正在努力发展。所以我决定在它的基础上构建,而不是构建一个 该软件包升级了我们现有的技术栈,以支持更多 通过Python轮子(不是 只是笔记本电脑)。
您会注意到,此包不提供任何工作流编排。 那就要来了。copula提供了一些基本的工作流功能, 我将逐步将其纳入这个系统。
