首页
/ Pathway 上云实战:用 AWS Fargate + BYOL 容器一键部署 Pathway 数据管道

Pathway 上云实战:用 AWS Fargate + BYOL 容器一键部署 Pathway 数据管道

2026-09-04 21:24:48作者:段琳惟

本文以 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.pycheckout_repository 实现看,当提供 repository_url 时,CLI 会:

  1. 检查 GitPython 是否可用(不可用则报错退出);
  2. 创建一个临时目录,用 git.Repo.clone_from 将仓库 clone 进去,若指定了 branch 还会执行 checkout
  3. 在临时目录内创建一个带 pip 的独立 venv。

随后 spawn_programcli.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.pyspawn_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 及其他必需环境变量注入容器,代码即可在云端跑起来——这就是整条部署链路最"轻"的部分。

获取容器镜像

  1. 打开 Marketplace 上的 Pathway BYOL 列表页,点击 "Continue to Subscribe";
  2. 阅读分发条款(BSL 1.1 License);
  3. 确认后点击 "Continue to Configuration";
  4. 履约方式(fulfillment method)选择 "Container image",并选择软件版本(建议最新版);
  5. 点击 "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);
    • essentialTrue 表示该容器为任务必需,它停止或失败时整个任务的其他容器都会停止。单容器任务必须设为 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.pywait_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 部署 教程,思路类似,可对照选择。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.12 K
2.72 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
527
590
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
904
1.82 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
854
1.34 K
docsdocs
暂无描述
Markdown
889
5.78 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.52 K
1.01 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.33 K
1.45 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
982
502
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
540
384