首页
/ Pathway 程序在 Azure 部署:基于 launch.py 与 Azure Container Instances 的完整实践指南

Pathway 程序在 Azure 部署:基于 launch.py 与 Azure Container Instances 的完整实践指南

2026-09-07 23:00:05作者:柯茵沙

本指南围绕 docs/2.developers/7.templates/ETL/_readmes/azure-aci-deploy.md 所对应的 ACI 部署示例展开:它演示了如何用一个名为 launch.py 的 Python 脚本,把 Pathway Live Data Framework(即本仓库 Pathway)的 Docker 镜像部署到 Azure Container Instances(ACI)上并托管远端 GitHub 仓库中的数据处理程序。读完本文,你将掌握示例仓库的完整文件结构、所有需要填写的常量与云服务凭证的获取方式、使用 Docker 或 virtualenv 两种方式运行该脚本的具体命令,以及脚本底层基于 Azure Python SDK 构建容器与容器组的调用原理,并能结合 examples/projects/azure-aci-deploy/ 中的源码独立复现整套部署。

阅读范围说明:本文的“主文档”是 ETL 模板目录下的示例说明,而示例的实际代码位于 examples/projects/azure-aci-deploy/;同时参考仓库内对同一主题的更完整讲解 docs/2.developers/4.user-guide/60.deployment/25.azure-aci-deploy.md 进行印证与细节补充。

示例仓库结构:三个文件承担的分工

示例目录(examples/projects/azure-aci-deploy/)只包含三个关键文件,各自职责非常清晰:

文件 作用
launch.py 核心 Python 脚本:负责把 Docker 镜像部署到 Azure Container Instances,并轮询容器状态、最后从 S3 上的 Delta Lake 读取并打印处理结果
requirements.txt launch.py 运行所需的 Python 依赖清单(见 requirements.txt
Dockerfile 定义示例自身的镜像配置,使整个示例可以在隔离环境中运行(见 Dockerfile

requirements.txt 的实际内容如下:

boto3
deltalake
pandas
azure-identity
azure-mgmt-containerinstance

其中 azure-mgmt-containerinstance 提供 ACI 管理客户端,azure-identityazure.core 提供认证支持,deltalakepandas 用于回读并展示结果,boto3 提供 S3 相关依赖。

Dockerfile 非常简单:基于 python:3.10 基础镜像,把 launch.pyrequirements.txt 拷贝进镜像,安装依赖后以 python launch.py 作为容器入口命令:

FROM python:3.10

COPY ./launch.py launch.py
COPY ./requirements.txt requirements.txt

RUN pip install -r requirements.txt

CMD ["python", "launch.py"]

使用前必读:需要更新的常量清单

运行示例之前,必须先更新 launch.py 顶部的常量。按功能可以分成四组:

1. Azure 相关(决定部署目标)

AZURE_SUBSCRIPTION_ID = "YOUR_AZURE_SUBSCRIPTION_ID"
AZURE_TOKEN_CREDENTIAL = "YOUR_AZURE_TOKEN_CREDENTIAL"
AZURE_RESOURCE_GROUP = "YOUR_AZURE_RESOURCE_GROUP"
AZURE_CONTAINER_GROUP_NAME = "pathway-test-container-group"
AZURE_CONTAINER_NAME = "pathway-test-container"
AZURE_LOCATION = "eastus"

2. Docker Hub 相关(决定拉取的镜像及其认证)

DOCKER_REGISTRY_USER = "YOUR_DOCKER_REGISTRY_USER"
DOCKER_REGISTRY_TOKEN = "YOUR_DOCKER_REGISTRY_TOKEN"
DOCKER_IMAGE_NAME = "pathwaycom/pathway:latest"

3. S3 存储相关(决定 Delta Lake 结果的落点与访问凭据)

AWS_S3_OUTPUT_PATH = "YOUR_AWS_S3_OUTPUT_PATH"
AWS_S3_ACCESS_KEY = "YOUR_AWS_S3_ACCESS_KEY"
AWS_S3_SECRET_ACCESS_KEY = "YOUR_AWS_S3_SECRET_ACCESS_KEY"
AWS_BUCKET_NAME = "YOUR_AWS_BUCKET_NAME"
AWS_REGION = "YOUR_AWS_REGION"

4. Pathway 与 GitHub 相关(决定被托管程序与功能授权)

PATHWAY_LICENSE_KEY = "YOUR_PATHWAY_LICENSE_KEY"
GITHUB_PERSONAL_ACCESS_TOKEN = "YOUR_GITHUB_PERSONAL_ACCESS_TOKEN"

launch.py 源码可见,AZURE_CONTAINER_GROUP_NAMEAZURE_CONTAINER_NAMEAZURE_LOCATION 已预置了可用默认值,属于“不必强制修改”的参数;其余常量则是实际值占位符,必须替换为真实凭证。

Azure 凭证如何获取

以下操作需要本机安装 Azure CLI。准备顺序与官方教程 docs/2.developers/4.user-guide/60.deployment/25.azure-aci-deploy.md 一致:

  1. 登录 Azure 并定位订阅 ID:

    az login
    az account get-access-token
    

    登录完成后浏览器会完成认证;从订阅列表的 Subscription ID 列(UUID4 格式)取订阅 ID 填入 AZURE_SUBSCRIPTION_ID;再从 az account get-access-token 返回的 JSON 中取 accessToken 字段填入 AZURE_TOKEN_CREDENTIAL。该令牌约每小时过期,启动容器前请确认它是新鲜的。

  2. 准备资源组:

    az group list --query "[].name" --output tsv
    az group create --name myResourceGroup --location eastus
    

    前者列出已有资源组,后者在资源组不存在时创建,把资源组名称填入 AZURE_RESOURCE_GROUP

为什么结果必须落在 S3 而不是容器本地

这个示例托管的是一个 ETL 程序(从 GitHub 拉取提交历史并写入 Delta Lake)。在云端容器里,本地磁盘是临时的——容器结束即销毁。因此示例把输出写到基于 S3 的 Delta Lake(对应 AWS_S3_OUTPUT_PATH = "s3://your-bucket/output-path" 这种格式),这样容器结束后结果仍然可访问。这也是 S3 相关常量必不可少的原因。

脚本工作原理:从常量到 ACI 上的运行容器

launch.py 的流程可以拆成六步,对应源码中的关键片段。理解这一步是把示例改造成自己的部署脚本的基础。

1. 用自定义凭证类完成 Azure SDK 认证

在 Docker 这类隔离环境中,最简做法是写一个实现 get_token 的包装类来替代交互式登录。源码 launch.py 中它把上文获得的 AZURE_TOKEN_CREDENTIAL 包装为有效期 3600 秒的访问令牌:

class TokenCredential:
    def __init__(self, token: str):
        self.token = token

    def get_token(self, *args, **kwargs):
        return AccessToken(self.token, 3600)

随后创建 ACI 管理客户端:

client = ContainerInstanceManagementClient(
    TokenCredential(AZURE_TOKEN_CREDENTIAL), AZURE_SUBSCRIPTION_ID
)

2. 把配置翻译为容器环境变量

函数 get_environment_variable_overrides()(见 launch.py)把上文的常量逐个包装成 ACI 的 EnvironmentVariable。完整列表与含义如下:

环境变量 作用
AWS_S3_OUTPUT_PATH Delta Lake 输出在 S3 上的完整路径,如 s3://your-bucket/output-path
AWS_S3_ACCESS_KEY S3 访问密钥 ID
AWS_S3_SECRET_ACCESS_KEY S3 密钥
AWS_BUCKET_NAME S3 桶名
AWS_REGION S3 桶所在区域
PATHWAY_LICENSE_KEY Pathway 许可证密钥;启用 Delta Lake 相关能力时需要
GITHUB_PERSONAL_ACCESS_TOKEN GitHub 个人访问令牌,用于拉取仓库、读取提交
PATHWAY_SPAWN_ARGS 传给 Pathway CLI spawn 的参数;本示例固定为 --repository-url https://github.com/pathway-labs/airbyte-to-deltalake python main.py

其中 PATHWAY_SPAWN_ARGS 是连接“通用 Pathway 镜像”与“你的具体程序”的桥梁。容器镜像 pathwaycom/pathway:latest 本身只装好运行环境,启动时执行 spawn-from-env,从环境变量 PATHWAY_SPAWN_ARGS 里解析出要运行的仓库与入口文件,从而实现“一个镜像跑任意 GitHub 仓库里的程序”。

3. 构建 Container 对象

源码 launch.py 中,容器配置为:

container = Container(
    name=AZURE_CONTAINER_NAME,
    image=DOCKER_IMAGE_NAME,
    resources=ResourceRequirements(
        requests=ResourceRequests(cpu=1, memory_in_gb=1.5)
    ),
    ports=[ContainerPort(port=80)],
    environment_variables=get_environment_variable_overrides(),
)

各字段含义:

  • name:容器在 Azure 内的标识名;
  • image:使用的 Docker 镜像,即 pathwaycom/pathway:latest
  • resources.requests:容器申请的 CPU(1 核)与内存(1.5 GB)。示例是轻量 ETL 任务,资源申请刻意设得很小;
  • ports:对外暴露的端口,这里开放 TCP 80;
  • environment_variables:上一步生成的运行时环境变量列表。

4. 构建 ContainerGroup 对象

ACI 中**容器组(Container Group)**是容器的最小调度单位:一组容器共享同一生命周期、网络与存储资源,同组容器运行在同一台宿主机上、通过本地网络互通、共享外部 IP 与端口。源码 launch.py 中容器组配置为:

container_group = ContainerGroup(
    location=AZURE_LOCATION,
    containers=[container],
    os_type=OperatingSystemTypes.linux,
    ip_address=IpAddress(ports=[Port(protocol="TCP", port=80)], type="Public"),
    restart_policy=ContainerGroupRestartPolicy.never,
    image_registry_credentials=[
        ImageRegistryCredential(
            server="index.docker.io",
            username=DOCKER_REGISTRY_USER,
            password=DOCKER_REGISTRY_TOKEN,
        )
    ],
)

各字段含义:

  • location:部署的 Azure 区域,与 AZURE_LOCATION 一致;
  • containers:组内容器列表,本例仅一个容器;
  • os_typeLinux 操作系统;
  • ip_address:容器组的对外 IP 配置(本例公开 TCP 80);
  • restart_policynever,表示容器终止后不自动重启;
  • image_registry_credentials:拉取私有镜像所需的 Docker Hub 凭据,serverindex.docker.io

5. 幂等化部署:先清旧组再创建

主流程(launch.py)在正式创建前,先尝试删除同名容器组:

client.container_groups.begin_delete(
    AZURE_RESOURCE_GROUP, AZURE_CONTAINER_GROUP_NAME
).wait()

删除失败(组不存在)则捕获异常跳过。随后调用 begin_create_or_update 真正发起部署。这种“先删后建”的做法保证了脚本可重复执行,不会因资源重名而中断。

6. 等待完成并从 S3 回读结果

begin_create_or_update 是异步的,因此脚本通过 wait_for_container_completionlaunch.py)每 10 秒轮询一次容器组状态,直到容器进入 Terminated 状态;随后依据退出码 exit_code 判断成功(0)或失败(非 0)。容器正常结束并写出 Delta Lake 后,脚本使用 deltalakeDeltaTable 连接 S3 并打印处理过的记录条数:

storage_options = {
    "AWS_ACCESS_KEY_ID": AWS_S3_ACCESS_KEY,
    "AWS_SECRET_ACCESS_KEY": AWS_S3_SECRET_ACCESS_KEY,
    "AWS_REGION": AWS_REGION,
    "AWS_BUCKET_NAME": AWS_BUCKET_NAME,
    # 关闭 DynamoDB 同步:本 Delta Lake 没有并发写入
    "AWS_S3_ALLOW_UNSAFE_RENAME": "True",
}
delta_table = DeltaTable(AWS_S3_OUTPUT_PATH, storage_options=storage_options)
pd_table_from_delta = delta_table.to_pandas()
print("Entries read and parsed: ", pd_table_from_delta.shape[0])

值得注意 AWS_S3_ALLOW_UNSAFE_RENAME: "True" 这个选项:示例场景是单写者写入 Delta Lake,没有并行写,因此可以关闭基于 DynamoDB 的提交协调,直接回读验证行数。

运行方式一:使用 Docker

与示例自带的 Dockerfile 配合,无需在宿主机安装 Python 依赖,两步即可运行:

# 1. 构建镜像
docker build . -t pathway-azure-container-instances-example

# 2. 运行容器
docker run -t pathway-azure-container-instances-example

构建基于 python:3.10RUN pip install -r requirements.txt 一次性装好五个依赖,CMD 直接执行 python launch.py。镜像内的脚本会去执行上面描述的 ACI 部署与轮询流程。注意:launch.py 中的常量是在构建时被拷贝进镜像的,因此修改常量后需要重新执行 docker build

运行方式二:使用 virtualenv

不引入 Docker 时,也可以在本地虚拟环境内直接跑脚本:

# 1. 创建虚拟环境
virtualenv venv

# 2. 激活虚拟环境
. venv/bin/activate

# 3. 安装依赖
pip install -r requirements.txt

# 4. 运行脚本
python launch.py

这两种方式在本示例的 README.md 中均有记载。无论哪种方式,都要求宿主机能够访问 Azure、Docker Hub、S3 与 GitHub 四个外部服务,并具备相应凭据。

底层支撑:PATHWAY_SPAWN_ARGS 与 CLI 的 spawn 机制

容器内部真正执行托管程序的是 Pathway 的 CLI。在仓库的 Python 侧实现 python/pathway/cli.py 中可以找到对应支撑:

  • spawn 命令接受 --repository-url(帮助文本为 “github repository path if the program is spawned from a repository”)与入口程序参数;spawn_program 内部会调用 checkout_repository 克隆该仓库,再在隔离环境中安装依赖并运行指定文件;
  • spawn_from_env 则是从环境变量 PATHWAY_SPAWN_ARGS 中读出参数再转发给 spawn,这正是 ACI 容器无需登录即可执行远端程序的关键(见 cli.py)。

因此在本地模拟容器内行为,可以等价地写出:

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

在 ACI 场景中,这组环境变量恰好就是 launch.pyenvironment_variables 传入容器的内容——本地可复现、云端可运行,两者共用同一套参数语义。

结果验证与常见问题

结果验证

部署成功的判断链有三环:

  1. launch.py 轮询打印的容器状态最终为 Terminated 且退出码为 0,说明托管程序正常跑完;
  2. 容器组在 Azure Portal 中可见,可查看执行阶段与 CPU、内存等资源指标;
  3. 脚本结尾打印 Entries read and parsed: <N>,N 即为写进 S3 Delta Lake 并被成功解析的记录行数。N 的具体数值取决于拉取时刻 GitHub 仓库中提交的数量(官方教程成稿时对应仓库含约 862 个提交),会随时间变化,应以实际运行结果为准。

常见问题与注意点

  • Token 过期AZURE_TOKEN_CREDENTIAL 约每小时失效,运行前重新执行 az account get-access-token 获取新值。
  • 资源配额不足:若 Azure 侧提示资源配额限制,需要到资源组中调高硬件上限后再部署。
  • 重启策略的选择:本示例托管的是“扫一遍即退出”的静态 ETL,因此 restart_policy 设为 never;如果你的程序需要持续运行(例如以 streaming 模式等待新数据),则应考虑让容器在结束后按相同参数自动重启,这通常要求把程序运行模式与重启策略配合调整。
  • 常量修改后需重新构建镜像:常量是构建期拷贝进镜像的,走 Docker 方式时必须 docker build 一次。

小结

这个示例演示的是一条完整的“云端托管 Pathway 程序”路径:在本地用 launch.py 收集好 Azure、Docker Hub、S3、GitHub 与 Pathway License 五类凭证,将其编码为 ACI 容器的环境变量,通过 Container + ContainerGroup 两个 Azure SDK 对象在 ACI 上拉起预装 Pathway 的通用镜像,再用 PATHWAY_SPAWN_ARGS 让容器内的 spawn-from-env 自动克隆并运行 GitHub 仓库中的程序,最后回到 S3 回读 Delta Lake 结果。参考示例代码 examples/projects/azure-aci-deploy/ 与更完整的教程文档 docs/2.developers/4.user-guide/60.deployment/25.azure-aci-deploy.md,把 launch.py 中的常量与 PATHWAY_SPAWN_ARGS 替换为你自己的仓库和程序入口,即可复用到任意自定义数据处理管线的云端部署上。

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

项目优选

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