FEATURED · 精选文章

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

发布时间 / 2026/9/8 23:55:42
来源 / 创域科博编辑部
栏目 / 资讯中心
Pathway 在 Azure Container Instances 上的部署示例解析:从 launch.py 到云上运行 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 InstancesACI。通过阅读本文你将掌握该示例仓库的文件结构与职责划分、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 MarketplaceBYOL 容器与 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 eastusAZURE_SUBSCRIPTION_IDAzure 订阅 IDUUID4 格式可通过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_LOCATIONAzure 数据中心区域默认eastus。第二组Docker 镜像仓库凭据DOCKER_REGISTRY_USER YOUR_DOCKER_REGISTRY_USER DOCKER_REGISTRY_TOKEN YOUR_DOCKER_REGISTRY_TOKEN DOCKER_IMAGE_NAME pathwaycom/pathway:latestACI 从 Docker Hub 拉取pathwaycom/pathway:latest时需要认证因此ImageRegistryCredential使用serverindex.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_TOKENPATHWAY_LICENSE_KEYPathway 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. 容器资源与环境变量定义业务容器本体是ContainerCPU 与内存被刻意设得很小1 vCPU / 1.5 GB因为示例只是简单 ETLcontainer Container( nameAZURE_CONTAINER_NAME, imageDOCKER_IMAGE_NAME, resourcesResourceRequirements( requestsResourceRequests(cpu1, memory_in_gb1.5) ), ports[ContainerPort(port80)], environment_variablesget_environment_variable_overrides(), )字段含义对应文档说明name标识容器image决定运行的应用镜像resources.requests声明容器请求的最小资源量ports暴露对外通信端口environment_variables在运行时注入配置。3. 注入 PATHWAY_SPAWN_ARGS接通 Pathway 官方镜像get_environment_variable_overrides()中关键的一项也是把 Pathway 官方镜像与你的业务代码接通的一环是EnvironmentVariable( namePATHWAY_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_TOKENYOUR_GITHUB_PERSONAL_ACCESS_TOKEN \ PATHWAY_LICENSE_KEYYOUR_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( locationAZURE_LOCATION, containers[container], os_typeOperatingSystemTypes.linux, ip_addressIpAddress(ports[Port(protocolTCP, port80)], typePublic), restart_policyContainerGroupRestartPolicy.never, image_registry_credentials[ ImageRegistryCredential( serverindex.docker.io, usernameDOCKER_REGISTRY_USER, passwordDOCKER_REGISTRY_TOKEN, ) ], )要点restart_policy取never与本示例任务形态一致——airbyte-to-deltalake默认以static模式扫描一次提交后退出属于一次性批处理任务而非常驻服务若在 Azure Marketplace 的 BYOL 部署中容器则按持续运行、退出自动重启的模式设计参见 教程 对INPUT_CONNECTOR_MODEstreaming的讨论。5. 幂等部署先删旧组再创建主流程先尝试删除同名旧容器组再创建保证脚本可重复执行try: print(fDeleting existing container group {AZURE_CONTAINER_GROUP_NAME}...) client.container_groups.begin_delete( AZURE_RESOURCE_GROUP, AZURE_CONTAINER_GROUP_NAME ).wait() ... except Exception: print(fContainer group {AZURE_CONTAINER_GROUP_NAME} does not exist, skipping deletion.) client.container_groups.begin_create_or_update( resource_group_nameAZURE_RESOURCE_GROUP, container_group_nameAZURE_CONTAINER_GROUP_NAME, container_groupcontainer_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(fContainer 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_optionsstorage_options) pd_table_from_delta delta_table.to_pandas() print(Entries read and parsed: , pd_table_from_delta.shape[0])AWS_S3_ALLOW_UNSAFE_RENAMETrue的注释给出了关键工程考量单写者、无并发写入时无需 Delta Lake 的事务锁协调故关闭 DynamoDB 同步以简化部署。最后to_pandas()统计行数并打印作为云端 ETL 结果的直观验证。常遇问题与注意事项Access Token 时效az account get-access-token获取的令牌约 1 小时过期长任务启动前务必刷新资源组配额教程提醒部署失败的一个常见原因是Insufficient regional vCPU quota left区域 vCPU 配额不足需到 Azure 门户提升资源组所在区域的配额一次性 vs 常驻本示例 ACI 采用restart_policynever适合批处理若要让 Pipeline 持续消费增量事件应参考 官方部署教程 中对 Marketplace BYOL 容器与streaming模式的说明这类容器退出后会自动以相同参数重启镜像认证从 Docker Hub 拉取私有/限流镜像必须提供ImageRegistryCredential请确保 Token 仍有权限且未过期。总结一条从本地脚手架到云上 Pipeline 的完整路径整个示例呈现了清晰的职责分层launch.pylauncher负责与 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),仅供参考
RELATED — 相关阅读

相关资讯

LATEST — 最新资讯

最新发布

TODAY — 本日精选

新闻

WEEKLY — 本周精选

新闻

MONTHLY — 本月精选

新闻