Pathway 在 Azure Container Instances 上的部署示例解析:从 launch.py 到云上运行
2026/9/8 23:55:34 网站建设 项目流程

Pathway 在 Azure Container Instances 上的部署示例解析:从 launch.py 到云上运行

【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway

本指南以 examples/projects/azure-aci-deploy 目录中的README.md与配套脚本为核心,深入解析如何将 Pathway 数据处理程序部署到 Azure Container Instances(ACI)。通过阅读本文,你将掌握该示例仓库的文件结构与职责划分、launch.py中各类 Azure/Docker/AWS 常量的配置方法,以及使用 Docker 与 virtualenv 两种方式运行示例的完整操作流程,并理解底层pathway spawn-from-env机制如何驱动容器内的代码执行。

示例定位:本仓库中可直接运行的 ACI 部署脚手架

该目录(examples/projects/azure-aci-deploy/)是 Pathway 官方教程"在 Azure 中使用 Azure Container Instances 运行 Pathway 程序"的配套代码。它并不重新实现一条业务 Pipeline,而是提供一个**"启动器"(launcher)脚手架**:在你本地机器上运行一个 Python 脚本,由该脚本代表你在 Azure 上创建容器组、拉取官方 Pathway 镜像、注入运行所需的环境变量,从而把一段远程 GitHub 仓库中的 Pathway 代码搬到云端运行。

对比仓库中的 官方部署教程(docs/2.developers/4.user-guide/60.deployment/25.azure-aci-deploy.md),示例中的launch.py即教程 Step 5/6 所述 Azure Python SDK 配置过程的完整可运行实现,教程同时提供了 Azure Marketplace(BYOL 容器)与 ACI 两条部署路径,而本示例专注的是纯 ACI 编程式部署这一条路。

仓库结构(共 4 个文件)

文件职责
launch.py核心 Python 脚本:在 Azure Container Instances 中部署 Docker 镜像并轮询容器状态、校验运行结果
requirements.txtlaunch.py运行所需的 Python 依赖清单
Dockerfile定义示例自身(launcher)的 Docker 镜像,便于在隔离环境中运行启动器
README.md使用说明

注意这里的双层镜像关系:启动器launch.py及其 Dockerfile)负责"调 Azure API",而真正承载 Pathway 业务程序的是 Docker Hub 上的pathwaycom/pathway官方镜像,该镜像由启动器在 ACI 中拉起。

launch.py 全貌:一个"常量先行"的部署脚本

launch.py的整体逻辑可划分为三块:顶部常量区(需要你手工填写)、get_environment_variable_overrides()环境变量构造函数、以及__main__中的创建容器组 → 删除旧组 → 部署 → 等待完成 → 读取 S3 上的 Delta Lake 结果。运行前,README 明确要求:launch.py中的常量更新为真实值。这些常量按用途分为四组。

第一组: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"
  • AZURE_SUBSCRIPTION_ID:Azure 订阅 ID(UUID4 格式),可通过az login后从订阅列表中获取;
  • AZURE_TOKEN_CREDENTIAL:访问令牌。脚本内自定义的TokenCredential类(见下)会包装它并通过AccessToken(token, 3600)暴露给 Azure SDK。该令牌约每小时过期,长时运行前需重新获取并更新;
  • AZURE_RESOURCE_GROUP:目标资源组名称,可执行az group create --name myResourceGroup --location eastus新建;
  • AZURE_CONTAINER_GROUP_NAME/AZURE_CONTAINER_NAME:容器组与容器名称,本例默认值pathway-test-container-group/pathway-test-container可保留;
  • AZURE_LOCATION:Azure 数据中心区域,默认eastus

第二组:Docker 镜像仓库凭据

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

ACI 从 Docker Hub 拉取pathwaycom/pathway:latest时需要认证,因此ImageRegistryCredential使用server="index.docker.io"+ 用户名 + Personal Access Token 的组合。你需要在 Docker Hub 账户中生成一个访问令牌(Personal Access Token)填入DOCKER_REGISTRY_TOKEN

第三组: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"

为什么业务输出要放在 S3?教程明确指出:容器是有状态易失的,一旦运行结束文件即被销毁,本地 Delta Lake 对用户不可达;因此本示例(对应 ETL 教程 所述的数据准备管线)把结果写到 S3 中的 Delta Lake,运行完成后仍可在本地读取。

第四组:业务运行凭据

PATHWAY_LICENSE_KEY = "YOUR_PATHWAY_LICENSE_KEY" GITHUB_PERSONAL_ACCESS_TOKEN = "YOUR_GITHUB_PERSONAL_ACCESS_TOKEN"
  • PATHWAY_LICENSE_KEY:Pathway Live Data Framework 许可证,启用 Delta Lake 等高级能力所需(可申请免费许可证);
  • GITHUB_PERSONAL_ACCESS_TOKEN:用于让pathway spawn解析 GitHub 仓库中的提交历史。

依赖清单与镜像构建方式

requirements.txt 内容如下:

boto3 deltalake pandas azure-identity azure-mgmt-containerinstance

前三个(boto3deltalakepandas)服务于"运行结束后读取 S3 中 Delta Lake 结果并转成 pandas DataFrame"这一校验步骤;后两个(azure-identityazure-mgmt-containerinstance)则是调用 Azure 容器实例管理 API 的核心。教程建议分两步安装核心 Azure 依赖:

pip install azure-identity pip install azure-mgmt-containerinstance

Dockerfile 负责把 launcher 本身容器化:

FROM python:3.10 COPY ./launch.py launch.py COPY ./requirements.txt requirements.txt RUN pip install -r requirements.txt CMD ["python", "launch.py"]

即以python:3.10为基座,将启动器脚本及其依赖打进镜像,容器启动即执行python launch.py

运行方式一:本地 Docker 运行启动器

README 给出最直接的执行路径——先构建镜像再运行容器:

docker build . -t pathway-azure-container-instances-example docker run -t pathway-azure-container-instances-example

该方式把 launcher 装进隔离环境,避免污染本机 Python 环境,适合"只运行一次云端部署任务"的场景。

运行方式二:virtualenv 运行启动器

若希望在本机直接调试脚本,可改用虚拟环境:

virtualenv venv . venv/bin/activate pip install -r requirements.txt python launch.py

两种方式最终都执行同一段launch.py主流程。

launch.py 主流程源码级拆解

1. 访问令牌封装

脚本自定义了极简凭据类以适配 Azure SDK 的TokenCredential协议:

class TokenCredential: def __init__(self, token: str): self.token = token def get_token(self, *args, **kwargs): return AccessToken(self.token, 3600)

get_token返回有效期 3600 秒(1 小时)的AccessToken,与管理令牌的过期节奏一致。随后由此构造管理客户端:

client = ContainerInstanceManagementClient( TokenCredential(AZURE_TOKEN_CREDENTIAL), AZURE_SUBSCRIPTION_ID )

2. 容器资源与环境变量定义

业务容器本体是Container,CPU 与内存被刻意设得很小(1 vCPU / 1.5 GB),因为示例只是简单 ETL:

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标识容器;image决定运行的应用镜像;resources.requests声明容器请求的最小资源量;ports暴露对外通信端口;environment_variables在运行时注入配置。

3. 注入 PATHWAY_SPAWN_ARGS,接通 Pathway 官方镜像

get_environment_variable_overrides()中关键的一项(也是把 Pathway 官方镜像与你的业务代码"接通"的一环)是:

EnvironmentVariable( name="PATHWAY_SPAWN_ARGS", value="--repository-url https://github.com/pathway-labs/airbyte-to-deltalake python main.py", ),

官方镜像默认入口命令为pathway spawn-from-env。查看 python/pathway/cli.py 的实现可知其语义:若环境变量PATHWAY_SPAWN_ARGS存在,就把它作为参数追加到spawn子命令后重新执行;若未设置则告警退出:

@cli.command() def spawn_from_env(): cli_spawn_arguments = os.environ.get("PATHWAY_SPAWN_ARGS") if cli_spawn_arguments is not None: args = ["spawn"] + cli_spawn_arguments.split(" ") os.execl(sys.executable, sys.executable, sys.argv[0], *args) else: logging.warning("PATHWAY_SPAWN_ARGS variable is unspecified, exiting...")

因此 ACI 容器启动后会自动执行等价于下面的本地命令——把airbyte-to-deltalake这个公共仓库中的main.py拉下来运行(从源码结构看,spawn --repository-url会处理仓库检出、依赖安装与进程拉起,具体分支逻辑也在cli.py的 spawn 实现中):

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

这种"配置即命令"的设计,正是整套 ACI 方案只需传环境变量即可运行任意公开 GitHub 仓库中 Pathway 代码的原因。示例中完整注入的环境变量还包括AWS_S3_OUTPUT_PATHAWS_S3_ACCESS_KEYAWS_S3_SECRET_ACCESS_KEYAWS_BUCKET_NAMEAWS_REGIONPATHWAY_LICENSE_KEYGITHUB_PERSONAL_ACCESS_TOKEN

4. 容器组定义

单个容器之上还需要容器组ContainerGroup,它决定生命周期、网络与重启策略:

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, ) ], )

要点:restart_policynever,与本示例任务形态一致——airbyte-to-deltalake默认以static模式扫描一次提交后退出,属于一次性批处理任务而非常驻服务;若在 Azure Marketplace 的 BYOL 部署中,容器则按"持续运行、退出自动重启"的模式设计(参见 教程 对INPUT_CONNECTOR_MODE=streaming的讨论)。

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

主流程先尝试删除同名旧容器组再创建,保证脚本可重复执行:

try: print(f"Deleting existing container group '{AZURE_CONTAINER_GROUP_NAME}'...") client.container_groups.begin_delete( AZURE_RESOURCE_GROUP, AZURE_CONTAINER_GROUP_NAME ).wait() ... except Exception: print(f"Container group '{AZURE_CONTAINER_GROUP_NAME}' does not exist, skipping deletion.") client.container_groups.begin_create_or_update( resource_group_name=AZURE_RESOURCE_GROUP, container_group_name=AZURE_CONTAINER_GROUP_NAME, container_group=container_group, )

若容器组不存在则忽略删除异常,继续执行begin_create_or_update,随后进入轮询阶段。

6. 轮询容器状态直到退出

wait_for_container_completion每 10 秒通过client.container_groups.get拉取一次容器状态,等待instance_view.current_state可用后判断:

container_state = instance_view.current_state.state if container_state == "Terminated": exit_code = container.instance_view.current_state.exit_code if exit_code == 0: print("Container completed successfully.") else: print(f"Container failed with exit code {exit_code}.")

初次创建后instance_view可能尚未就绪,脚本会打印提示并继续等待;状态为Terminated时依据退出码判定成败并跳出循环;拉取状态本身出错(HttpResponseError)也会中断轮询。

7. 结果校验:读取 S3 中的 Delta Lake

容器退出后,脚本用deltalake库直接读取 S3 上的输出,验证 Pipeline 是否写入了预期数据:

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, # Disabling DynamoDB sync since there are no parallel writes into this 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 同步以简化部署。最后to_pandas()统计行数并打印,作为云端 ETL 结果的直观验证。

常遇问题与注意事项

  • Access Token 时效az account get-access-token获取的令牌约 1 小时过期,长任务启动前务必刷新;
  • 资源组配额:教程提醒部署失败的一个常见原因是Insufficient regional vCPU quota left(区域 vCPU 配额不足),需到 Azure 门户提升资源组所在区域的配额;
  • 一次性 vs 常驻:本示例 ACI 采用restart_policy=never,适合批处理;若要让 Pipeline 持续消费增量事件,应参考 官方部署教程 中对 Marketplace BYOL 容器与streaming模式的说明,这类容器退出后会自动以相同参数重启;
  • 镜像认证:从 Docker Hub 拉取私有/限流镜像必须提供ImageRegistryCredential,请确保 Token 仍有权限且未过期。

总结:一条从本地脚手架到云上 Pipeline 的完整路径

整个示例呈现了清晰的职责分层:launch.py(launcher)负责与 Azure 控制面交互,pathwaycom/pathway:latest官方镜像负责运行环境,PATHWAY_SPAWN_ARGS则是把"要跑哪个仓库里的哪段代码"以纯环境变量方式传递的粘合剂,S3 + Delta Lake 承担持久化输出,轮询与to_pandas完成闭环验证。

若你的场景与之类似——业务代码托管在公共 GitHub 仓库、结果需落到 S3/Delta Lake、又不想维护 Kubernetes 集群,那么 examples/projects/azure-aci-deploy 是一个开箱即用的参照实现:替换launch.py顶部常量、按 README 的 Docker 或 virtualenv 方式启动即可。更完整的 Azure Marketplace 图形化部署流程与每一步的 Azure CLI 命令,可继续阅读 docs/2.developers/4.user-guide/60.deployment/25.azure-aci-deploy.md(该目录下还有同系列的 AWS Fargate 与 Nebius 部署教程可供横向对照)。

【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询