首页
/ 将 Pathway ETL 流水线部署到 AWS Fargate:CLI、BYOL 容器与任务定义的完整实战

将 Pathway ETL 流水线部署到 AWS Fargate:CLI、BYOL 容器与任务定义的完整实战

2026-09-04 18:29:41作者:毕习沙Eudora

本文以 Pathway 官方开发者文档 AWS Fargate 部署指南 为主体,讲解如何将本地开发完成的 Pathway 数据流水线部署到 AWS Fargate 无服务器计算环境:从 Pathway CLI 的 spawn / spawn-from-env 机制、AWS Marketplace 上的 BYOL(Bring Your Own License)容器,到 ECS 任务定义、集群创建、run_task 启动与环境变量注入的完整链路,并逐字段解读每个配置项。读完本文,你将掌握一套“仓库地址 + 环境变量即可上云”的 Pathway 生产化部署方案,并理解其背后的 CLI 源码实现。

前置条件

在开始之前,文档明确了项目需要满足两项基本要求:

  1. 项目托管在公开的 GitHub 仓库中——这是 --repository-url 参数能直接拉取并运行代码的前提;
  2. 根目录存在 requirements.txt,并列出项目所有 Python 依赖——CLI 会在隔离环境中自动安装这些依赖。

示例 ETL 流水线:airbyte-to-deltalake

教程以“为 Spark 分析做数据准备”这一 ETL 场景为例:流水线追踪 GitHub 仓库的 commit 历史、剔除敏感数据,并把结果写入 Delta Lake。该示例项目(airbyte-to-deltalake)相对原版教程做了两处简化,正好契合云容器场景:

  • GitHub PAT(Personal Access Token)改为从环境变量读取;
  • 移除了 Spark 计算部分,因为云容器内并不必要。

有一个关键点必须注意:输出目的地。原始教程有两种输出模式——本地 Delta Lake 或 S3 上的 Delta Lake。在云端部署时,本地 Delta Lake 只存在于远端容器内部,容器销毁后数据即丢失,用户也无法访问,因此本教程采用 S3 上的 Delta Lake 作为输出。相应地,容器需要额外的一组 S3 访问环境变量(下文 Step 4 详细展开)。

Pathway 仓库的 examples/projects 目录下保留了与本文档对应的完整示例工程 aws-fargate-deploy,包含本地启动器 launch.pyDockerfilerequirements.txt。同目录下还有 spark-data-preparationrealtime-log-monitoringkafka-ETL 等可参考的 ETL 项目模板。

两大部署工具:Pathway CLI 与 BYOL 容器

Pathway CLI:从仓库直接运行代码

安装 Pathway Live Data Framework 时会附带命令行工具。pathway spawn 命令用于以多线程/多进程方式启动应用,例如 pathway spawn python main.py 运行本地的 main.py

本教程重点使用它的另一个能力:直接运行 GitHub 仓库中的代码,无需本地克隆。以 airbyte-to-deltalake 为例,设置两个环境变量——GITHUB_PERSONAL_ACCESS_TOKEN(GitHub PAT)与 PATHWAY_LICENSE_KEY(Pathway 许可证密钥)——然后调用 pathway spawn --repository-url 指定仓库:

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 后,CLI 会自动完成三件事:检出仓库、在隔离环境中安装 requirements.txt 中声明的依赖、运行指定入口文件。

源码级印证:--repository-urlspawn-from-env 的实现

这一行为可以在 python/pathway/cli.py 中得到验证。spawn 命令通过 click 声明了 --repository-url(帮助文本为 "github repository path if the program is spawned from a repository")以及 --branch(用于检出非默认分支)两个选项,随后交给 spawn_program 处理:

  • checkout_repository:若传入 repository_url,会检查系统已安装 git(否则报 "To run the code from Git repository please have git installed"),然后在临时目录中 git.Repo.clone_from 克隆仓库,并按需 checkout 指定分支;
  • spawn_program:在临时根目录下划分 repositoryvenv 两个路径,读取仓库根目录的 requirements.txt 安装依赖,最后 os.chdir(repository_path) 进入仓库目录再启动程序。

这正是“仓库托管于 GitHub + requirements.txt 声明依赖”这两条前置要求的底层原因。

此外还有 PATHWAY_SPAWN_ARGS 环境变量作为 pathway spawn 的快捷方式:

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

从源码 spawn_from_env 可以看到其机制非常简洁:读取 PATHWAY_SPAWN_ARGS,按空格切分后拼接到 spawn 命令之后,再用 os.execl 以当前 Python 解释器重新执行;若变量未设置则输出警告 "PATHWAY_SPAWN_ARGS variable is unspecified, exiting..." 并退出。CHANGELOG.md 中也记录了该命令的引入:“A new CLI command, spawn-from-env, has been added. This command runs the Pathway CLI spawn command using arguments provided in the PATHWAY_SPAWN_ARGS environment variable.” 这个“环境变量携带启动参数”的设计正是云容器部署的关键——容器镜像可以保持通用不变,只需在运行时注入不同的 PATHWAY_SPAWN_ARGS 即可指向任意仓库。

BYOL 容器:AWS Marketplace 上的现成镜像

Pathway 提供 BYOL(Bring Your Own License)容器,作为 AWS Marketplace 上架的 Docker 镜像,内部已预装 Pathway Live Data Framework 及其全部依赖。该 listing 免费,可在 Marketplace 直接获取;不输入许可证密钥也能运行,输入密钥则解锁框架全部功能。容器以 BSL 1.1 许可证分发,对应仓库根目录的 LICENSE.txt

容器启动时执行的就是 pathway spawn-from-env 命令:把 PATHWAY_SPAWN_ARGS 以及其他所需环境变量传入容器,代码即可在云上跑起来。

在 AWS Fargate 中运行示例

获取 BYOL 容器镜像

  1. 打开 AWS Marketplace 上的 Pathway BYOL 容器 listing,点击 "Continue to Subscribe";
  2. 阅读容器要约及其分发条款(BSL 1.1 License),确认后再次点击 "Continue to Configuration";
  3. 选择履约方式 Container image(本教程采用此方式),并选择软件版本(建议使用最新版);
  4. 点击 "Continue to Launch"。

完成后,你就获得了一个托管在 AWS 基础设施上的 Pathway BYOL Docker 镜像,可用于 Amazon ECS 和 Amazon EKS 服务。镜像路径显示在页面 "Container images" 区块的 CONTAINER_IMAGE 变量中,指向 Amazon ECR 仓库,后续任务定义中将用到。

本地启动脚本

后续的登录、配置 Fargate、运行实例等步骤,可以全部放在一个本地运行的启动器脚本(示例中的 launch.py)里一次性完成。仓库中的 launch.py 即为完整参考实现,其依赖在 requirements.txt 中声明(boto3deltalakepandas)。

Step 1: 登录 AWS CLI

先安装 AWS CLI 并完成命令行登录,同时准备好 Access Key 与 Secret Access Key。登录后,用 boto3 创建 ECS 与 ECR 客户端(区域名按你的环境调整):

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 的实际写法,脚本会先尝试 get_authorization_token(),若抛异常则打印 "AWS connection failed" 并 exit(1),作为登录/凭证有效性的快速校验。

Step 2: 在 AWS 中注册任务定义

ECS 的**任务定义(Task Definition)**是容器化应用的蓝图:声明容器镜像、CPU/内存、网络设置与环境变量,保证任务在不同环境中一致、可重复地运行,并可对接负载均衡与自动扩缩容。任务定义以字典形式描述:

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 任务定义族名,将同一任务定义的多个版本分组,便于追踪变更与回滚;示例值为 "pathway-container-test"
networkMode 容器网络模式;"awsvpc" 让每个任务拥有独立的弹性网络接口,隔离性与安全性更好
requiresCompatibilities 任务兼容的启动类型;"FARGATE" 表示使用 Fargate 无服务器计算引擎,无需额外底层设置
cpu 任务分配的 CPU 单位,"2048" 即 2 vCPU(1024 units = 1 vCPU)
memory 任务内存(MiB),"8192" 即 8 GiB
executionRoleArn IAM 角色 ARN,授权 ECS 拉取镜像与推送日志;将 YOUR_ACCOUNT_ID 替换为你的账户 ID 即可得到该 ARN
containerDefinitions 任务中的容器定义数组。本教程只运行一个容器
containerDefinitions[].name 容器名,任务内唯一;示例为 "pathway-container-test-definition"
containerDefinitions[].image 容器镜像地址(仓库 + 路径 + tag),使用前面获取的 CONTAINER_IMAGE
containerDefinitions[].essential 标记容器是否为任务必需;必需容器停止/失败会导致整个任务停止。本例仅一个容器,故设为 True

注册并拿到任务定义 ARN(唯一标识,务必保存):

response = ecs_client.register_task_definition(**task_definition)
task_definition_arn = response["taskDefinition"]["taskDefinitionArn"]

Step 3(可选): 创建集群

集群(Cluster) 是任务与服务在 ECS 中的逻辑分组。使用 Fargate 时,集群作为任务运行的环境,屏蔽底层基础设施细节,同时划定网络边界、隔离不同应用/环境。若已有 ECS 集群可直接复用;否则用默认设置新建一个即可满足教程需要:

cluster_name = "pathway-test-cluster"
response = ecs_client.create_cluster(clusterName=cluster_name)

Step 4: 启动容器

先配置环境变量。如前所述,输出不能落在容器本地(容器销毁即丢失),而应写入 S3 上的 Delta Lake:

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",  # You can get it at https://pathway.com/features
    },
    {
        "name": "GITHUB_PERSONAL_ACCESS_TOKEN",
        "value": "YOUR_GITHUB_PERSONAL_ACCESS_TOKEN",  # You can get it at https://github.com/settings/tokens
    },
    {
        "name": "PATHWAY_SPAWN_ARGS",
        "value": "--repository-url https://github.com/pathway-labs/airbyte-to-deltalake python main.py",  # Doesn't need to be changed
    },
]

各环境变量的用途:

  • PATHWAY_LICENSE_KEY:运行 Delta Lake 功能所必需的 Pathway 许可证密钥(可获取免费许可);
  • AWS_S3_OUTPUT_PATH:输出在 S3 中的完整路径;
  • AWS_S3_ACCESS_KEY / AWS_S3_SECRET_ACCESS_KEY:S3 访问密钥与秘密密钥;
  • AWS_BUCKET_NAME:目标 S3 桶名称;
  • AWS_REGION:S3 桶所在区域;
  • GITHUB_PERSONAL_ACCESS_TOKEN:GitHub 个人访问令牌,用于拉取私有/受保护的仓库代码;
  • PATHWAY_SPAWN_ARGS:Pathway CLI 参数。本例指定从 pathway-labs/airbyte-to-deltalake 仓库运行 main.py,该值通常无需修改。

注意 containerOverrides 的机制:环境变量不是写死在任务定义里,而是在启动时通过 overrides 注入容器——这与 BYOL 容器“镜像通用、参数注入”的设计完全一致。

随后用 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(或 family 名);
  • cluster:任务运行的 ECS 集群,即 Step 3 创建(或已有)的集群;
  • launchType"FARGATE" 表示在 AWS Fargate 上运行;
  • networkConfiguration / awsvpcConfiguration:任务在 VPC 中的网络配置——
    • subnets:子网 ID 列表,至少指定一个,可从 AWS 控制台的 Subnets 页面任选;
    • assignPublicIp"ENABLED" 表示任务启动时分配公网 IP,使其可被互联网访问;
  • count:任务实例数,本例只需 1 个;
  • overrides / containerOverrides:为本次云启动覆盖注入的环境变量。

对照仓库中 launch.py 的完整实现,run_task 返回后脚本会取出任务 ARN,并进入一个 wait_for_completion 轮询循环:每 10 秒调用一次 describe_tasks 检查 lastStatus,直到状态变为 STOPPED 后打印 "Task ... is completed."。运行 python launch.py 后,你也可以在 ECS 控制台的集群 overview 页面观察任务状态直至执行结束。

验证执行结果

任务完成后,可用 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

其中 AWS_S3_ALLOW_UNSAFE_RENAME 置为 "True" 是因为本场景不存在对同一 Delta Lake 的并发写入,可以安全地禁用 DynamoDB 锁表同步。文档给出的验证结果是 664 行——写作时点 pathwaycom/pathway 仓库恰好有 664 个 commit,与流水线处理数量吻合。仓库示例 launch.py 的最后一段同样完成了这段校验,并打印 "Entries read and parsed: ..."。

总结

云部署是数据工程项目走向生产的关键一步:代码在远端可靠、可预测地运行,资源可灵活管理,还能自主选择可用区;但对初学者而言,容器、云服务、虚拟机等组件的组合确实容易让人望而却步。

本文展示的 Pathway 方案把这件事简化成了三步:从 AWS Marketplace 获取预装 Pathway CLI 的 BYOL 容器,设置 PATHWAY_SPAWN_ARGS(指向 GitHub 仓库与入口脚本)及其余环境变量,然后在 Fargate 上启动任务。CLI 源码层面看,spawn-from-env 只是把环境变量切成参数重新 exec 到 spawn,而 spawn 负责克隆仓库、安装 requirements.txt 依赖并运行入口文件——整条链路上没有任何“把代码打进镜像”的环节,这正是该部署模式可复用性强的原因。你可以克隆仓库中的示例工程 examples/projects/aws-fargate-deploy 作为起点,开发自己的云端数据抽取方案;同类部署场景还可以参考仓库中的 Azure ACI 部署文档

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