Pathway 上云实战:用 AWS Fargate + BYOL 容器一键部署 Pathway 数据管道
本文以 Pathway 官方开发者文档 AWS Fargate 部署教程 为主体,完整讲解如何将本地开发验证过的 Pathway 实时数据管道部署到 AWS 云端:使用 Pathway CLI 的 spawn --repository-url / spawn-from-env 从 GitHub 仓库直接拉起代码,配合 AWS Marketplace 上免费的 Pathway BYOL 容器,通过一段本地 launch.py 脚本完成 ECS 任务定义注册、Fargate 集群创建、容器任务启动,最终把结果写入 S3 上的 Delta Lake 并用 deltalake 包验证。配套的可运行示例位于 examples/projects/aws-fargate-deploy 目录。
前置条件
在开始部署前,官方文档要求你的项目满足两条基础约束:
- 项目托管在一个公开的 GitHub 仓库中(因为 Fargate 容器内要通过
--repository-url直接 clone 代码); - 项目根目录的
requirements.txt中完整列出全部 Python 依赖(CLI 会在容器内的隔离虚拟环境中依据该文件安装依赖)。
示例管道:Data Preparation for Spark Analytics
教程选用的示例是 Data Preparation for Spark Analytics 模板:构建一个 ETL 流程,追踪 GitHub 仓库的 commit 历史、剔除敏感数据,再把清洗结果加载进 Delta Lake。示例基于 pathway-labs/airbyte-to-deltalake 仓库,并做了两处简化:
- GitHub PAT(Personal Access Token)改为从环境变量读取;
- 移除了 Spark 计算部分,因为在云端容器中它不是必需的。
这里有一个云端部署的关键设计点:输出目标必须是 S3 上的 Delta Lake,而不是本地文件系统上的 Delta Lake。原本该任务支持两种输出模式,但在 Fargate 容器里写本地路径时,数据只存在于远端 worker 的容器内部,容器销毁后输出即丢失,用户根本无法访问。因此教程统一采用 S3 作为输出端,代价是容器需要额外一组访问 S3 服务的环境变量(下文 Step 4 详述)。
Pathway CLI:从 GitHub 仓库直接运行代码
安装 Pathway 后会附带一个命令行工具,最基础的能力是用多线程或多进程拉起本地代码,例如 pathway spawn python main.py 运行本地的 main.py。本教程重点使用它的另一项能力:通过 --repository-url 直接从 GitHub 仓库拉取并运行代码,无需在本地 checkout。
以 airbyte-to-deltalake 示例为例,设置两个环境变量(GITHUB_PERSONAL_ACCESS_TOKEN 提供仓库访问凭证,PATHWAY_LICENSE_KEY 提供 Pathway 许可证)后执行:
GITHUB_PERSONAL_ACCESS_TOKEN=YOUR_GITHUB_PERSONAL_ACCESS_TOKEN \
PATHWAY_LICENSE_KEY=YOUR_PATHWAY_LICENSE_KEY \
pathway spawn --repository-url https://github.com/pathway-labs/airbyte-to-deltalake python main.py
源码视角:--repository-url 到底做了什么
从 python/pathway/cli.py 的 checkout_repository 实现看,当提供 repository_url 时,CLI 会:
- 检查 GitPython 是否可用(不可用则报错退出);
- 创建一个临时目录,用
git.Repo.clone_from将仓库 clone 进去,若指定了branch还会执行checkout; - 在临时目录内创建一个带 pip 的独立 venv。
随后 spawn_program(cli.py)在 venv 中执行 pip install -r requirements.txt(安装失败会抛出 RuntimeError: Failed to install dependencies),把工作目录切到 clone 出的仓库路径,最后以 venv 里的解释器启动目标脚本。这正是文档所说"CLI 自动处理 checkout、依赖安装与运行"的底层机制,也解释了为什么前置条件必须要求仓库公开、requirements.txt 完整。
另外,CLI 支持 PATHWAY_SPAWN_ARGS 环境变量作为 pathway spawn 参数的快捷方式,配合 spawn-from-env 子命令使用:
GITHUB_PERSONAL_ACCESS_TOKEN=YOUR_GITHUB_PERSONAL_ACCESS_TOKEN \
PATHWAY_LICENSE_KEY=YOUR_PATHWAY_LICENSE_KEY \
PATHWAY_SPAWN_ARGS='--repository-url https://github.com/pathway-labs/airbyte-to-deltalake python main.py' \
pathway spawn-from-env
对应的实现在 cli.py 的 spawn_from_env:它从环境变量读取 PATHWAY_SPAWN_ARGS 并交给 spawn 流程解析;若该变量未设置,则打印警告后退出。这个子命令正是 BYOL 容器的启动方式(见下节)。
BYOL 容器:AWS Marketplace 上的 Pathway 镜像
Pathway 在 AWS Marketplace 上提供 BYOL(Bring Your Own License)容器:一个预装了 Pathway 框架及其全部依赖的现成 Docker 镜像。可以不带许可证 key 使用,配置 key 后解锁完整功能;该 Marketplace 列表本身免费订阅使用(许可协议为 BSL 1.1)。
关键行为:容器启动时执行的是 pathway spawn-from-env 命令。因此你只需把 PATHWAY_SPAWN_ARGS 及其他必需环境变量注入容器,代码即可在云端跑起来——这就是整条部署链路最"轻"的部分。
获取容器镜像
- 打开 Marketplace 上的 Pathway BYOL 列表页,点击 "Continue to Subscribe";
- 阅读分发条款(BSL 1.1 License);
- 确认后点击 "Continue to Configuration";
- 履约方式(fulfillment method)选择 "Container image",并选择软件版本(建议最新版);
- 点击 "Continue to Launch"。
完成后即得到一个托管在 AWS 基础设施上的 BYOL Docker 镜像,可用于 Amazon ECS 与 EKS;镜像路径在页面 "Container images" 一节的 CONTAINER_IMAGE 变量中,后续任务定义的 image 字段就填它。
用 launch.py 脚本驱动整个 Fargate 部署
文档将全部 AWS 操作收敛到一个本地运行的启动脚本 launch.py 中。仓库中提供了完整可运行的 launch.py,其依赖(boto3、deltalake、pandas)见 requirements.txt;该目录还附带一个 Dockerfile,基于 python:3.10 安装依赖后 CMD ["python", "launch.py"],方便把 launcher 本身也容器化运行。
Step 1:登录 AWS CLI 并建立客户端
先安装并登录 AWS CLI,准备好 access key 与 secret access key。登录后用 boto3 创建 ECS 与 ECR 客户端(region 按你的环境调整,示例用 eu-central-1):
import boto3
ecs_client = boto3.client("ecs", region_name="eu-central-1")
ecr_client = boto3.client("ecr", region_name="eu-central-1")
登录成功的话,还能获取当前会话的 ECR 认证 token,用于验证凭证有效:
auth_token = ecr_client.get_authorization_token()
仓库版 launch.py 将这一步包在 try/except 中,失败时打印 "AWS connection failed" 并 exit(1),是一个值得照抄的健壮性处理。
Step 2:注册 ECS 任务定义
任务定义(task definition)是 ECS 管理容器化应用的"蓝图":描述容器镜像、CPU/内存、网络模式与环境变量,保证任务可被一致地重复部署,并能对接负载均衡、自动扩缩等其他 AWS 服务。本例用字典定义:
task_definition = {
"family": "pathway-container-test",
"networkMode": "awsvpc",
"requiresCompatibilities": ["FARGATE"],
"cpu": "2048",
"memory": "8192",
"executionRoleArn": "arn:aws:iam::YOUR_ACCOUNT_ID:role/ecsTaskExecutionRole",
"containerDefinitions": [
{
"name": "pathway-container-test-definition",
"image": "PATH_TO_YOUR_IMAGE",
"essential": True,
}
],
}
逐字段说明:
family:任务定义族名称,将同一任务的多个版本归组,便于追踪变更与回滚;networkMode:"awsvpc"为每个任务分配独立的弹性网络接口,隔离性与安全性更好(Fargate 强制要求 awsvpc 模式);requiresCompatibilities:"FARGATE"表示使用 AWS Fargate 无服务器计算引擎,省去自建节点等额外配置;cpu:CPU 单位数,"2048"= 2 vCPU(1024 单位 = 1 vCPU);memory:任务内存(MiB),"8192"= 8 GiB;executionRoleArn:ECS 拉取镜像、发布日志所用的 IAM 角色 ARN,将你的账号 ID 填入模板即可;containerDefinitions:任务内的容器列表。本教程只运行一个容器:name:容器名(任务内唯一),示例为"pathway-container-test-definition";image:容器 Docker 镜像,填入 Step 0 获取的CONTAINER_IMAGE路径(指向 Amazon ECR);essential:True表示该容器为任务必需,它停止或失败时整个任务的其他容器都会停止。单容器任务必须设为True。
提交任务定义:
response = ecs_client.register_task_definition(**task_definition)
ECS 客户端返回创建出的任务定义 ARN(其唯一标识),务必保存,供后续步骤使用:
task_definition_arn = response["taskDefinition"]["taskDefinitionArn"]
Step 3(可选):创建 ECS 集群
集群是任务与服务的一种逻辑分组,用于管理和扩缩容器化应用。使用 Fargate 时,集群是任务运行的环境,屏蔽了底层基础设施细节;它也用于组织与隔离资源、定义网络边界、关联 IAM 角色。已有集群可直接复用,否则几行代码即可创建:
cluster_name = "pathway-test-cluster"
response = ecs_client.create_cluster(clusterName=cluster_name)
以默认设置创建的集群已足够运行本教程。
Step 4:注入环境变量并启动容器
再次强调:输出必须走 S3,本地输出在容器销毁后即丢失。因此需要配置一组 S3 输出与环境凭证变量:
s3_output_path = "YOUR_S3_OUTPUT_PATH" # Example: "s3://aws-demo/runs/16.08.2024/1/"
s3_bucket_name = "YOUR_S3_BUCKET_NAME" # Example: "aws-demo"
s3_region = "YOUR_S3_REGION" # Example: "eu-central-1"
s3_access_key = "YOUR_AWS_S3_ACCESS_KEY"
s3_secret_access_key = "YOUR_AWS_S3_SECRET_ACCESS_KEY"
environment_vars = [
{
"name": "AWS_S3_OUTPUT_PATH",
"value": s3_output_path,
},
{
"name": "AWS_S3_ACCESS_KEY",
"value": s3_access_key,
},
{
"name": "AWS_S3_SECRET_ACCESS_KEY",
"value": s3_secret_access_key,
},
{
"name": "AWS_BUCKET_NAME",
"value": s3_bucket_name,
},
{
"name": "AWS_REGION",
"value": s3_region,
},
{
"name": "PATHWAY_LICENSE_KEY",
"value": "YOUR_PATHWAY_LICENSE_KEY",
},
{
"name": "GITHUB_PERSONAL_ACCESS_TOKEN",
"value": "YOUR_GITHUB_PERSONAL_ACCESS_TOKEN",
},
{
"name": "PATHWAY_SPAWN_ARGS",
"value": "--repository-url https://github.com/pathway-labs/airbyte-to-deltalake python main.py",
},
]
各环境变量要点:
| 变量 | 说明 |
|---|---|
PATHWAY_LICENSE_KEY |
Delta Lake 等功能依赖许可证 key,可免费注册获取 |
AWS_S3_OUTPUT_PATH |
输出在 S3 中的完整路径 |
AWS_S3_ACCESS_KEY / AWS_S3_SECRET_ACCESS_KEY |
访问 S3 的凭证 |
AWS_BUCKET_NAME |
S3 bucket 名称 |
AWS_REGION |
bucket 所在 region |
GITHUB_PERSONAL_ACCESS_TOKEN |
拉取 GitHub 仓库所需 PAT(在 GitHub 的 "Personal access tokens" 页面创建) |
PATHWAY_SPAWN_ARGS |
Pathway CLI 参数,本例指定从 pathway-labs/airbyte-to-deltalake 仓库运行 main.py,一般无需修改 |
然后调用 run_task 启动容器:
task = {
"taskDefinition": task_definition_arn,
"cluster": cluster_name,
"launchType": "FARGATE",
"networkConfiguration": {
"awsvpcConfiguration": {
"subnets": ["REPLACE_WITH_YOUR_SUBNET_ID"],
"assignPublicIp": "ENABLED",
}
},
"count": 1,
"overrides": {
"containerOverrides": [
{
"name": "pathway-container-test-definition",
"environment": environment_vars,
}
]
}
}
ecs_client.run_task(**task)
参数逐项解释:
taskDefinition:Step 2 注册的任务定义 ARN;cluster:Step 3 创建的 ECS 集群名;launchType:"FARGATE",在 Fargate 无服务器引擎上运行;networkConfiguration/awsvpcConfiguration:VPC 网络设置subnets:任务将启动的 VPC 子网 ID 列表,至少一个(可从 AWS 控制台的 Subnets 页面任选);assignPublicIp:"ENABLED"表示任务启动时分配公网 IP,可被互联网访问;
count:任务实例数,本例只需 1 个;overrides/containerOverrides:为云端启动注入的环境变量覆盖。
执行后任务即被创建并运行,可在 ECS 集群概览页观察执行状态,直至任务结束。仓库版 launch.py 额外提供了一个 wait_for_completion 轮询函数:循环调用 describe_tasks 检查 lastStatus,每 10 秒打印一次状态,直到变为 STOPPED 才结束——这让脚本能自动衔接后续的验证步骤。
验证执行结果
任务完成后,用 deltalake(delta-rs 的 Python 包)读取 S3 上的 Delta Lake 验证数据:
from deltalake import DeltaTable
# Create an S3 connection settings dictionary
storage_options = {
"AWS_ACCESS_KEY_ID": s3_access_key,
"AWS_SECRET_ACCESS_KEY": s3_secret_access_key,
"AWS_REGION": s3_region,
"AWS_BUCKET_NAME": s3_bucket_name,
# Disabling DynamoDB sync since there are no parallel writes into this Delta Lake
"AWS_S3_ALLOW_UNSAFE_RENAME": "True",
}
# Read a table from S3
delta_table = DeltaTable(
s3_output_path,
storage_options=storage_options,
)
pd_table_from_delta = delta_table.to_pandas()
# Print the number of commits processed
pd_table_from_delta.shape[0]
输出(文档写作时):
664
与当时 pathwaycom/pathway 仓库的 commit 总数一致,可交叉验证数据完整性。其中 AWS_S3_ALLOW_UNSAFE_RENAME 用于禁用 DynamoDB 锁同步——因为该场景不存在并行写入,单写者下可以关闭以提升简单性。launch.py 在 wait_for_completion 返回后自动执行这段验证逻辑,并打印 "Entries read and parsed: N"。
小结
云端部署是数据工程项目走向生产的关键一步:代码远端运行、免受本地环境干扰,且具备弹性资源管理与更高的稳定性,但对新手而言容器、云服务、虚拟机等组件组合起来往往相当复杂。本教程展示的 Pathway 方案把复杂度压到最低:从 AWS Marketplace 免费订阅 BYOL 容器(内含 Pathway CLI 与全部依赖)→ 在任务定义里指定镜像 → 用 PATHWAY_SPAWN_ARGS 声明要运行的 GitHub 仓库与入口脚本 → 通过 ECS/Fargate 的 run_task 一键启动。仓库内的 launch.py 覆盖了从登录校验、任务注册、集群创建、轮询等待到 Delta Lake 结果验证的完整闭环,是可直接照搬改造的部署脚手架。若你的目标平台不是 AWS,仓库还提供了 GCP 部署 与 Render 部署 教程,思路类似,可对照选择。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0622
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00