Pathway 上云实战:用 Azure Marketplace 与 Azure Container Instances 部署 Pathway 数据管道
本文基于 Pathway 官方开发者指南《Deploy to Azure》,讲解如何把本地开发调试完毕的 Pathway Live Data Framework(Python ETL 流处理框架)应用部署到 Azure 云生态:先梳理一条跟踪 GitHub 提交历史并写入 S3 版 Delta Lake 的示例 ETL 管道,再完整给出两条部署路径——通过 Azure Marketplace 的 Pathway BYOL 容器走四步向导快速建集群,以及用 Python 脚本驱动 Azure Container Instances(ACI)从零拉起公共 Docker 镜像。读完本篇,你可以直接复用 launch.py 的完整编排逻辑,掌握 pathway spawn --repository-url 远程仓库执行机制与环境变量配置的全部细节。
为什么需要云部署,以及本篇的示例场景
Pathway Live Data Framework 用于定义和运行各种数据处理管道,官方文档提供了大量教程,例如实时日志监控、基于 Kafka 的 ETL 管道、为 Spark 分析做数据准备等。这些管道在本地开发与测试完成后,下一步就是部署上云:让代码在远程运行,摆脱本地机器故障带来的中断,是进入生产环境的关键一步。云部署还支持弹性资源管理、更高的稳定性,以及自由选择应用可用的可用区。
本篇以仓库中真实存在的示例工程 examples/projects/azure-aci-deploy 为主线。它对应"Data Preparation for Spark Analytics"教程的改造版,管道逻辑为:跟踪 pathwaycom/pathway 仓库的 GitHub 提交历史,去除敏感数据,把结果加载进 Delta Lake。相比原版教程,云端示例做了两处简化:
- GitHub PAT(Personal Access Token)改为从环境变量读取;
- 移除了 Spark 计算部分,因为云端容器里并不必要。
还有一个关键的输出端设计:原版有两种输出模式——本地 Delta Lake 或 S3 版 Delta Lake。云部署时本地模式不实际,因为数据只存在于远端容器内部,用户无法访问;所以示例改用 S3 版 Delta Lake 存放结果,容器需要额外的一组 S3 访问环境变量(下文详述)。
开始部署前,项目需满足两个基本要求:
- 项目托管在公开的 GitHub 仓库中;
- 仓库根目录的
requirements.txt列出了项目全部 Python 依赖。
核心机制:Pathway CLI 直接运行 GitHub 仓库中的代码
Pathway Live Data Framework 安装后自带命令行工具。基础的 spawn 命令用于以多线程/多进程方式启动应用,例如 pathway spawn python main.py 运行本地的 main.py。本篇用到的核心能力是另一项:不要求代码在本地,直接从 GitHub 仓库拉取并执行。
以 airbyte-to-deltalake 示例为例,设置 GITHUB_PERSONAL_ACCESS_TOKEN(你的 GitHub PAT)和 PATHWAY_LICENSE_KEY(Pathway 许可证密钥)两个环境变量后,通过 --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 实现源码 可以看到这套机制的完整调用链。spawn 命令定义在 python/pathway/cli.py#L307-L410,其中 --repository-url 与 --branch 选项声明于 L362-L371。当提供 --repository-url 参数时,spawn_program 会先调用 checkout_repository(L192-L208):
- 通过 GitPython 把仓库克隆到临时目录(要求运行环境装有 git,否则直接报错退出);
- 如指定了
--branch,checkout 到对应分支; - 在临时目录内创建一个隔离的 venv 环境。
随后在 L224-L246 中,CLI 检查仓库根目录是否存在 requirements.txt,若存在则用 venv 内的 pip 安装全部依赖,失败即抛错终止;成功后 chdir 进入仓库目录,以 venv 中的解释器启动 program。这解释了为什么"公开仓库 + requirements.txt"是硬性前提:依赖安装完全依赖仓库自带的 requirements.txt。
spawn-from-env:为容器化部署设计的启动入口
除了显式传参,Pathway CLI 还提供 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 的实现非常简洁(python/pathway/cli.py#L531-L538):它读取 PATHWAY_SPAWN_ARGS,按空格切分后拼接到 spawn 命令后用 os.execl 重新执行自身;若该变量未设置,仅打印警告后退出。
这正是所有容器化部署方案(包括本指南的两条 Azure 路径)的启动约定:官方 Docker 镜像与 Azure Marketplace BYOL 容器均以 pathway spawn-from-env 作为入口,使用者只需把 spawn 参数写进 PATHWAY_SPAWN_ARGS、把凭据写进容器环境变量,代码即可在云上跑起来。
路径一:Azure Marketplace 四步向导部署(推荐)
Pathway 在 Azure Marketplace 提供 Pathway - BYOL(Bring Your Own License)容器。该条目本身免费,不会产生 Marketplace 费用(需自行申请免费层许可证密钥)。它本质上是一个预装了 Pathway 框架与全部依赖的 Docker 镜像,外加一套配置好的 Kubernetes 环境,适合在 Azure 内快速落地。
选择这条路径的核心理由是简单:跟随四步部署向导、填入项目特定配置即可。若 Marketplace 方案不适用,再考虑下文更复杂的 ACI 方案。
需要留意一个运行语义:容器和标准云部署一样,期望程序持续运行。程序若因任何原因结束,会以相同的启动参数和环境变量自动重启。这也直接影响示例管道的 INPUT_CONNECTOR_MODE 取值(见下文环境变量说明)。
第 1 步:Basics(基础配置)
点击 Marketplace 条目页的 "Get It Now" 按钮打开向导。第一步在 Basics 页:
-
通过下拉菜单选择订阅与资源组。没有订阅可先在 Azure 门户的 Subscriptions 页面创建。
-
新建资源组可点 "Resource group" 下拉框下的 "Create new",或用 Azure CLI:
az group create --name myResourceGroup --location eastus若走 CLI 路线,需先执行
az login登录。 -
选定后,指定是否需要开发集群(选 "Yes"),并从下拉框选择集群所在区域。
第 2 步:Cluster Details(集群详情)
- 在 "AKS cluster name" 字段填入部署集群的字母数字名称;
- 选择最新的 Kubernetes 版本;
- 硬件配置默认值通常够用,如需调整可点 "Change size" 链接修改 vCPU 与内存;
- 自动扩缩容默认开启,可设置多台虚拟机。
第 3 步:Application Details(应用详情)
-
为集群扩展与应用程序命名;
-
"Pathway App License" 字段填入 Pathway 许可证密钥;
-
添加
PATHWAY_SPAWN_ARGS,指定pathway spawn启动参数。本教程从airbyte-to-deltalake仓库启动管道,取值为:--repository-url https://github.com/pathway-labs/airbyte-to-deltalake python main.py -
附加环境变量:定义所有必需变量。示例管道同时用到 S3 与 GitHub,需要:
环境变量 说明 AWS_S3_OUTPUT_PATHS3 输出存储路径,如 s3://your-bucket/output-pathAWS_S3_ACCESS_KEYS3 访问密钥 AWS_S3_SECRET_ACCESS_KEYS3 私有访问密钥 AWS_BUCKET_NAMES3 桶名称 AWS_REGIONS3 桶所在区域 GITHUB_PERSONAL_ACCESS_TOKENGitHub 访问令牌(在 GitHub 设置的 tokens 页面生成) INPUT_CONNECTOR_MODE决定如何从 GitHub 轮询提交,可取 "static"或"streaming"。默认"static"只扫描全部提交、处理一遍后退出;"streaming"则持续运行,等待新提交并按到达顺序追加写入输出。由于 Azure 部署中程序结束后会自动重启,应设为"streaming",使程序常驻并把新提交追加到已有集合
第 4 步:Review + Create(确认并创建)
阅读 BYOL 容器的条款与隐私策略,核对全部输入无误后点击 "Create"。
创建完成后的检查
- 创建发起后,页面右上角出现部署进行中的通知;一切顺利时很快会看到 "Deployment succeeded"。
- "Deployment details" 区可查看集群创建的所有资源;某资源部署失败时点击它查看错误详情。常见错误是 "Insufficient regional vCPU quota left",此时需在资源组中提高硬件配额。
- 整个创建过程约需 5–10 分钟。完成后页面显示 "Your deployment is complete",在 "Next steps" 区点击 "Go to resource",进入与应用绑定的 Kubernetes 服务页面——服务此时已上线运行。
这条路径无需复杂步骤,只要有订阅和资源组即可完成,是官方推荐的云上部署方式。
路径二:Azure Container Instances(ACI)从零部署
若 Marketplace 不可用,或你需要以容器实例方式完全自定义部署,可改用 Dockerhub 上的公共 Pathway 镜像(镜像名 pathwaycom/pathway:latest)。该镜像预装了 Pathway 与全部依赖,且不绑定任何云生态;不使用许可证密钥也可以跑,填入密钥则解锁框架完整功能,条目本身免费。
与 Marketplace 方案的区别在于:镜像来自 Azure 生态之外,没有 Marketplace 获取步骤,但需要本地完成登录、配置 ACI、指定 S3 存储变量、启动框架实例等一系列操作。这些操作由一个本地运行的启动脚本(即 launch.py)一次性完成,整体流程是:先在本地配置系统并从各相关系统获取令牌,再用这些令牌在云端运行计算。
该示例工程的文件结构(见 README):
- launch.py:负责在 ACI 中部署 Docker 镜像的 Python 脚本;
- requirements.txt:
launch.py的依赖清单(boto3、deltalake、pandas、azure-identity、azure-mgmt-containerinstance); - Dockerfile:基于
python:3.10的隔离运行环境,CMD ["python", "launch.py"]。
运行方式二选一:构建 Docker 镜像后 docker run,或用 virtualenv 安装 requirements.txt 后 python launch.py。
Step 1:Azure 配置与令牌获取
需要 Azure CLI(未安装请按官方安装指南处理)。关键三步:
-
登录 Azure:
az login会打开浏览器完成认证,选择要使用的订阅租户。从
Subscription ID列复制 UUID4 格式的订阅 ID,作为AZURE_SUBSCRIPTION_ID。 -
获取访问令牌:
az account get-access-token复制返回 JSON 中的
accessToken值赋给AZURE_TOKEN_CREDENTIAL。令牌较长,建议存为环境变量;注意令牌每小时过期,启动容器前确认其为最新值。 -
资源组:
az group list --query "[].name" --output tsv az group create --name myResourceGroup --location eastus把资源组名称赋给
AZURE_RESOURCE_GROUP。
其余参数预先定义、一般无需修改,汇总如下(与 launch.py 源码 中的常量一一对应):
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_CONTAINER_GROUP_NAME:容器组名称(本例只有一个容器);AZURE_CONTAINER_NAME:容器名称;AZURE_LOCATION:Azure 数据中心区域,如eastus。
Step 2:Dockerhub 凭据配置
ACI 支持直接从 Dockerhub 拉取并运行容器,因此只需在代码中存放 Dockerhub 凭据(与 launch.py#L30-L32 对应):
DOCKER_REGISTRY_USER = "YOUR_DOCKER_REGISTRY_USER"
DOCKER_REGISTRY_TOKEN = "YOUR_DOCKER_REGISTRY_TOKEN"
DOCKER_IMAGE_NAME = "pathwaycom/pathway:latest"
其中个人访问令牌在 Dockerhub 的 "Personal access tokens" 页面生成;镜像直接使用最新的 Pathway 镜像 pathwaycom/pathway:latest。
Step 3:Delta Lake 存储后端(S3)配置
由于容器结束后内部文件会被删除,结果必须落在持久化存储中,示例使用 Amazon S3 承载 Delta Lake,需配置:
- 输出路径、桶名、区域:
AWS_S3_OUTPUT_PATH(如s3://your-bucket/output-path)、AWS_S3_BUCKET_NAME、AWS_REGION; - S3 凭据:AWS Access Key 与 Secret Access Key(获取方式见 AWS IAM 文档)。
AWS_S3_OUTPUT_PATH = "YOUR_AWS_S3_OUTPUT_PATH"
AWS_BUCKET_NAME = "YOUR_AWS_BUCKET_NAME"
AWS_REGION = "YOUR_AWS_REGION"
AWS_S3_ACCESS_KEY = "YOUR_AWS_S3_ACCESS_KEY"
AWS_S3_SECRET_ACCESS_KEY = "YOUR_AZURE_S3_SECRET_ACCESS_KEY"
Step 4:许可证密钥与 GitHub PAT
启用 Delta Lake 功能、解析 GitHub 提交还需要最后两样东西:Pathway 许可证密钥(可在 Pathway 官网申请免费层)与 GitHub Personal Access Token(在 GitHub "Personal access tokens" 页面生成):
PATHWAY_LICENSE_KEY = "YOUR_PATHWAY_LICENSE_KEY"
GITHUB_PERSONAL_ACCESS_TOKEN = "YOUR_GITHUB_PERSONAL_ACCESS_TOKEN"
Step 5:用 Azure Python SDK 配置容器
在 ACI 中,Container 是一个轻量、独立、可执行的软件包,包含运行应用所需的一切(代码、运行时、库、依赖),各容器隔离运行但共享宿主内核;ACI 让你无需管理底层基础设施即可在云上运行容器。管理这些资源使用 Azure Python SDK,需先安装:
pip install azure-identity
pip install azure-mgmt-containerinstance
导入配置过程中需要的类:
from azure.core.credentials import AccessToken
from azure.core.exceptions import HttpResponseError
from azure.mgmt.containerinstance import ContainerInstanceManagementClient
from azure.mgmt.containerinstance.models import (
Container,
ContainerGroup,
ContainerGroupRestartPolicy,
ContainerPort,
EnvironmentVariable,
ImageRegistryCredential,
IpAddress,
OperatingSystemTypes,
Port,
ResourceRequests,
ResourceRequirements,
)
配置过程分两步:先建 environment_variables 数组,再实例化 Container 类。以下代码与 launch.py 中 get_environment_variable_overrides() 的实际实现一致:
environment_variables = [
EnvironmentVariable(name="AWS_S3_OUTPUT_PATH", value=AWS_S3_OUTPUT_PATH),
EnvironmentVariable(name="AWS_S3_ACCESS_KEY", value=AWS_S3_ACCESS_KEY),
EnvironmentVariable(name="AWS_S3_SECRET_ACCESS_KEY", value=AWS_S3_SECRET_ACCESS_KEY),
EnvironmentVariable(name="AWS_BUCKET_NAME", value=AWS_BUCKET_NAME),
EnvironmentVariable(name="AWS_REGION", value=AWS_REGION),
EnvironmentVariable(name="PATHWAY_LICENSE_KEY", value=PATHWAY_LICENSE_KEY),
EnvironmentVariable(name="GITHUB_PERSONAL_ACCESS_TOKEN", value=GITHUB_PERSONAL_ACCESS_TOKEN),
EnvironmentVariable(
name="PATHWAY_SPAWN_ARGS",
value="--repository-url https://github.com/pathway-labs/airbyte-to-deltalake python main.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=environment_variables,
)
Container 各字段说明:
name:容器在 Azure 中的标识名;image:容器使用的 Docker 镜像,包含应用代码与依赖;resources:分配给容器的计算资源。requests指定容器运行所需的精确资源量(CPU 与内存)。此处设为 1 核 / 1.5GB,因为这条简单 ETL 管道不需要太多资源;ports:容器对外暴露的通信端口;environment_variables:传入容器的环境变量,用于运行时配置应用。其中PATHWAY_SPAWN_ARGS指定执行pathway-labs/airbyte-to-deltalake仓库中的main.py,容器入口pathway spawn-from-env会据此启动管道。
Step 6:创建容器组
在 ACI 中,Container Group 是一组共享同一生命周期、网络与存储资源的容器:同组容器运行在同一宿主机上,可通过本地网络互通,共享对外 IP 与端口。同样用 SDK 构建:
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 区域;containers:组内容器列表,本例只有一个;os_type:使用 Linux 操作系统;ip_address:IP 地址配置(公网类型,暴露 80 端口 TCP);restart_policy:never,即容器组不自动重启——运行一次、终止后不再拉起。这与 Marketplace 方案(自动重启)语义相反:ACI 场景下程序可以按static模式跑完即止,容器随后终止;image_registry_credentials:访问容器镜像仓库(这里是 Dockerhub,服务器地址index.docker.io)的凭据。
Step 7:启动容器
最后创建 Azure 客户端实例并调用相应方法。在隔离环境中处理 Azure SDK 认证的最简方式是自实现一个令牌包装类(与 launch.py#L44-L49 一致):
class TokenCredential:
def __init__(self, token: str):
self.token = token
def get_token(self, *args, **kwargs):
return AccessToken(self.token, 3600)
然后创建客户端并发起创建:
client = ContainerInstanceManagementClient(
TokenCredential(AZURE_TOKEN_CREDENTIAL), AZURE_SUBSCRIPTION_ID
)
client.container_groups.begin_create_or_update(
resource_group_name=AZURE_RESOURCE_GROUP,
container_group_name=AZURE_CONTAINER_GROUP_NAME,
container_group=container_group,
)
启动后可到 Azure 门户查看运行阶段与资源用量等指标,也可以在门户中停止执行。
实际 launch.py 中额外的两处工程细节
对照 launch.py 完整源码 可以发现两处教程未强调、但复用时很有价值的健壮性处理:
- 幂等清理:创建前先尝试
begin_delete删除同名容器组,若不存在则跳过(L152-L162)。这保证脚本可重复执行,不会被上一次残留的容器组阻塞; - 轮询等待容器终结:
wait_for_container_completion(L90-L118)每 10 秒查询一次容器组状态,直到容器状态变为Terminated,并依据exit_code判断成功(0)或失败——这也呼应了restart_policy=never的设计:跑完即止,由脚本确认结果。
访问并验证执行结果
服务成功启动并完成 ETL 后,可用 delta-rs Python 包(deltalake,已在 requirements.txt 中列出)验证 S3 上的 Delta Lake 内容。launch.py 结尾 的实现:
from deltalake import DeltaTable
# 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,
# 该 Delta Lake 不存在并行写入,禁用 DynamoDB 锁同步
"AWS_S3_ALLOW_UNSAFE_RENAME": "True",
}
# 从 S3 读取表
delta_table = DeltaTable(
AWS_S3_OUTPUT_PATH,
storage_options=storage_options,
)
pd_table_from_delta = delta_table.to_pandas()
# 打印处理的提交数量
pd_table_from_delta.shape[0]
其中 AWS_S3_ALLOW_UNSAFE_RENAME 用于关闭 DynamoDB 锁机制——因为本示例没有并发写入该 Delta Lake。教程撰写时验证得到的行数为 862,与当时 pathwaycom/pathway 仓库的提交总数一致,可作为端到端数据完整性的核对方式。
两种方案对比与结语
| 维度 | Azure Marketplace(BYOL) | Azure Container Instances |
|---|---|---|
| 部署方式 | 四步向导,自动创建 AKS 集群 | 本地 launch.py 脚本 + Azure Python SDK |
| 复杂度 | 低,推荐默认选择 | 较高,适合无 Marketplace 或需深度定制的场景 |
| 重启语义 | 程序结束后自动以原参数重启,应使用 streaming 模式 |
restart_policy=never,跑完即止,脚本轮询确认 Terminated |
| 凭据来源 | 向导内直接填写 | az account get-access-token(每小时过期)+ Dockerhub PAT |
| 成本 | Marketplace 条目免费,另计 Azure 资源费用 | 镜像条目免费,另计 ACI 资源费用 |
两条路径的共同底座是同一套 CLI 约定:PATHWAY_SPAWN_ARGS + pathway spawn-from-env + 一组环境凭据。理解了 python/pathway/cli.py 中 checkout_repository 的"克隆仓库 → 建 venv → 装 requirements → 启动"调用链,你就能把这一套部署模式平移到仓库中其他部署文档(如 AWS Fargate、GCP、Render 教程)——只是承载容器的那朵云换了而已。云部署是项目走向生产的关键一环:它让解决方案稳定、可预测地运行,并提供灵活的资源管理与可用区选择;而借助 Pathway CLI 与上述两种 Azure 路径,原本令人望而生畏的容器与云服务配置被压缩到了"填变量、点 Create、跑脚本"的程度。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0623
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00