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.txt | launch.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前三个(boto3、deltalake、pandas)服务于"运行结束后读取 S3 中 Delta Lake 结果并转成 pandas DataFrame"这一校验步骤;后两个(azure-identity、azure-mgmt-containerinstance)则是调用 Azure 容器实例管理 API 的核心。教程建议分两步安装核心 Azure 依赖:
pip install azure-identity pip install azure-mgmt-containerinstanceDockerfile 负责把 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_PATH、AWS_S3_ACCESS_KEY、AWS_S3_SECRET_ACCESS_KEY、AWS_BUCKET_NAME、AWS_REGION、PATHWAY_LICENSE_KEY、GITHUB_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_policy取never,与本示例任务形态一致——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),仅供参考