多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

在 Floci 中模拟 Amazon MWAA:双容器架构、启动脚本与 CLI 代理实战指南

在 Floci 中模拟 Amazon MWAA:双容器架构、启动脚本与 CLI 代理实战指南 在 Floci 中模拟 Amazon MWAA双容器架构、启动脚本与 CLI 代理实战指南【免费下载链接】flociLight, fluffy, and always free - The AWS Local Emulator alternative项目地址: https://gitcode.com/gh_mirrors/fl/flociFloci 的 MWAAAmazon Managed Workflows for Apache Airflow实现通过真实运行 Apache AirflowLocalExecutor与独立 Postgres 元数据库的方式为本地开发与 CI 提供与 AWS 托管服务行为一致的 Airflow 编排环境。本文基于 docs/services/mwaa.md 展开结合src/main/java/io/github/hectorvent/floci/services/mwaa/下的源码与测试系统讲解其协议入口、双容器架构、环境生命周期、DAG/依赖同步机制、配置项与已知边界帮助你正确地在本地复现 MWAA 的 API 形状与运行时行为。协议与路由REST-JSON、双重签名注册MWAA 服务遵循 REST-JSON 协议所有请求都经由 Floci 的主端点http://localhost:4566/通过 JAX-RS 路径路由分发。与 EKS 类似MWAA 使用标准 HTTP 动词承载 JSON 请求体而非 JSON 1.1X-Amz-Target或 Query 协议——路径与动词严格对齐真实 MWAA APICreateEnvironment是PUTUpdateEnvironment是PATCH。这一映射体现在 MwaaController.java操作方法与路径返回CreateEnvironmentPUT /environments/{name}{Arn: ...}GetEnvironmentGET /environments/{name}{Environment: {...}}ListEnvironmentsGET /environments{Environments: [...]}UpdateEnvironmentPATCH /environments/{name}{Arn: ...}DeleteEnvironmentDELETE /environments/{name}{}CreateWebLoginTokenPOST /webtoken/{name}见下文CreateCliTokenPOST /clitoken/{name}见下文一个容易被忽略的细节是签名名signing name与端点 id 的分离MWAA 请求以airflow作为 SigV4 签名名却注册在mwaa端点 id 下。在 ResolvedServiceCatalog.java 中两者被同时注册Set.of(airflow, mwaa)与 Bedrock Runtime 的签名名/端点 id 拆分采用同一套双重注册技术确保无论按哪种 scope 查找凭证都能正确解析。两种运行模式Mock 模式mock: trueMock 模式下环境元数据仅保存在进程内存储后端为mwaa-environments.json见 MwaaService.java不启动任何 Docker 容器。CreateEnvironment后环境状态直接转为AVAILABLE。适合 CI 场景或仅需 MWAA API 形状而非真实 Airflow 实例的契约测试。真实模式mock: false默认每个环境启动一对容器由 MwaaEnvironmentManager.java 负责全生命周期管理floci-mwaa-account.region.name-db私有postgres元数据库。从不发布宿主机端口仅由同网络内的 Airflow 兄弟容器通过 Docker 网络访问Floci 直接从 Docker daemon 解析其容器 IP见 resolveContainerIp避免依赖默认 bridge 网络上的容器名 DNS。floci-mwaa-account.region.name-airflow真实apache/airflow容器以LocalExecutor运行webserver 与 scheduler 在同一进程树内并挂载dags、logs两个命名卷。Airflow 容器需要被 Floci 自身就绪轮询器、Web 代理后端可达因此采用与 Neptune/RDS 后端容器一致的宿主机端口模式Floci 原生运行时分配动态宿主机端口Floci 自身也在容器内时则仅暴露端口见 startAirflowContainer。Airflow 版本是真实生效的而非回显占位。请求的AirflowVersion直接决定镜像 tagapache/airflow:version-python3.112.10.x 及以前与apache/airflow:version-python3.122.11.0 起严格对齐真实 Amazon MWAA 各版本对应的 Airflow/Python 配对。其判定逻辑见 pythonTagFor主版本大于 2或主版本等于 2 且次版本 ≥ 11 时使用 python3.12。版本值会经supported-versions白名单校验见 resolveAirflowVersion不合法的请求在创建前即被拒绝。Airflow 容器覆盖了官方镜像的 entrypoint改为运行一段引导脚本airflowBootstrapScript先执行可选的/startup.sh再airflow db migrate随后幂等地创建 admin 用户复用官方_AIRFLOW_WWW_USER_*变量名后台启动 scheduler最后exec airflow webserver使其成为 PID 1。scheduler 与迁移步骤之间通过保证顺序避免 webserver 与迁移竞争启动。AirflowConfigurationOptions是真实生效的每个section.key条目都会被翻译为 Airflow 容器上的AIRFLOW__SECTION__KEY环境变量——这正是真实 Amazon MWAA 采用的机制因此 DAG 内的conf.get(section, key)能读到配置值。转换逻辑见 airflowConfigurationOptionsEnv缺少section.key点号分隔点前或点后为空的畸形键会被跳过并记录日志会与 Floci 自身管理的变量executor、数据库连接、安全密钥、AWS-SDK 重定向等冲突的条目同样被跳过并记录日志避免破坏环境必需的接线保护清单见 PROTECTED_ENV_VARS转换后的变量名遵循 Airflow 自身的 section/key 模型core.dags_are_paused_at_creation→AIRFLOW__CORE__DAGS_ARE_PAUSED_AT_CREATION。同时CreateEnvironment在服务端强制校验键/值形状与 botocore 对该字段的约束一致键为 1–64 字符、匹配[a-z]([a-z0-9._]*[a-z0-9_])?值为 1–65536 字符的可打印 ASCII见 validateAirflowConfigurationOptions。由于 Floci 可通过裸 HTTP 或 CLI 的--cli-input-json访问绕过 SDK 客户端侧校验这一服务端校验是必要的兜底。环境生命周期与就绪判定创建流程createEnvironment校验名称非空、环境中不存在重名botocore 模型对 CreateEnvironment 未声明 “already exists” 形状因此重复创建以ValidationException返回解析并校验AirflowVersion、校验AirflowConfigurationOptions构造 ARN 与环境实体初始状态为CREATING未提供的字段取默认值EnvironmentClass默认mw1.small、WebserverAccessMode默认PUBLIC_ONLY、MaxWorkers10、MinWorkers1、Schedulers2真实模式下按顺序完成从 S3 拉取启动脚本 → 启动容器 → 从配置区间分配代理端口 → 启动 Web/CLI 代理任一步骤失败则环境标记为CREATE_FAILED并在finally中回滚停止代理、停止已启动的容器、释放预留端口rollbackFailedCreate绝不泄漏端口或遗留孤儿容器。就绪判定后台就绪轮询器每 3 秒检查一次首个检查延迟 2 秒见 startReadinessPoller。只有当 Airflow 免认证的/health端点同时报告metadatabase与scheduler为healthy时环境才转为AVAILABLEisReady 及其无依赖的 JSON 解析 healthySection。如果 Postgres 或 Airflow 容器在就绪前退出——比如启动脚本或airflow db migrate中途失败、Postgres 本身挂掉——同一轮询器检测到已停止的容器后会将环境转为CREATE_FAILEDcheckReadiness而不是永远轮询死容器。区域语义遗留环境保留其 ARN 中记录的区域。另一区域的请求不会采纳该环境需在请求区域重建。Web/CLI 代理Floci 为每个环境运行一个 HTTP 代理发布在配置端口区间的宿主机端口上前端对接真实 Airflow webserver。除POST /aws_mwaa/cli外所有请求原样转发至 Airflow 内部host:8080并流式返回响应实现基于 Vert.xHttpClient/HttpServer并剥离 hop-by-hop 头见 MwaaWebProxy.java。POST /aws_mwaa/cli由代理自行拦截handleCli校验Authorization: Bearer CliToken头是否匹配该环境由CreateCliToken签发的令牌令牌保存在进程内映射中不持久化无效令牌返回 401将请求体原始 Airflow CLI 命令如dags list通过docker exec在容器内以sh -c airflow cmd执行runAirflowCli以 AWS 文档规定的{stdout: base64, stderr: base64}形状返回。Environment.WebserverUrl与CreateWebLoginToken返回的 URL 都指向该代理因此 Airflow UI 浏览与 CLI-token 流程走同一入口。令牌生成使用SecureRandom的 24 字节 URL-safe Base64generateToken。代理的启停由 MwaaProxyManager.java 统一管理环境删除或 Floci 关闭时随之回收代理端口在环境删除时显式释放deleteEnvironment避免反复创建/删除耗尽端口区间。启动脚本StartupScriptS3Path若配置了启动脚本Floci 会在环境创建时从 S3 拉取脚本并在 Airflow 容器启动前早于airflow db migrate/scheduler/webserver注入容器顺序与真实 MWAA 一致。三个从 API 表面看不出的关键行为免密 sudo容器内airflowOS 用户被授予无密码sudo写入/etc/sudoers.d/floci-airflow-nopasswd见 startAirflowContainer专为支持 AWS 官方启动脚本文档中sudo apt-get ...这类写法——官方apache/airflow镜像自带 sudo 但默认要求密码若无此授权任何 sudo 行都会挂起。该授权是无条件授予的不依赖是否配置了脚本因为它属于环境的基础能力而非脚本相关能力。受保护环境变量快照恢复Floci 管理的一组环境变量Fernet 密钥、数据库连接串、executor 设置、AWS-SDK 重定向变量会在脚本运行前快照、运行后立即恢复脚本无论有意无意改动其中任何一个都不会破坏环境。这对应真实 MWAA 的 “reserved environment variables” 行为。在引导脚本中体现为先以_FLOCI_ORIG_*快照、if [ -f /startup.sh ]; then . /startup.sh || exit 1; fi以 source 方式执行导出变量对后续步骤可见、随后恢复并继续迁移airflowBootstrapScript。脚本失败即创建失败脚本返回非零、或注入失败都直接使CreateEnvironment失败与真实 MWAA 以启动脚本成功为环境创建门槛的行为一致。注入失败会被当作硬错误抛出见 startAirflowContainer 中对启动脚本与 sudoers 授权的区别处理启动脚本对象缺失/不可读也视为真实配置错误并传播失败fetchStartupScript。DAG 内 AWS SDK 调用的重定向Airflow 容器通过LaunchedContainerAwsEnv.sdkBaselineEnv被指回 Floci 自身AWS_ENDPOINT_URL、占位凭证、区域与 Lambda/ECS 容器访问模拟器的方式一致见 startAirflowContainer。因此 DAG 内自己的boto3.client(s3)、boto3.client(secretsmanager)等调用会解析到 Floci 而非真实 AWS——这使本地端到端编排Airflow 调度 DAG 内 AWS 操作成为可能。前置条件Docker socket真实模式会启动特权 Docker 容器。需要挂载 Docker socket 并设置 Docker 网络使容器间可互相访问services: floci: image: floci/floci:latest volumes: - /var/run/docker.sock:/var/run/docker.sock ports: - 4566:4566 environment: FLOCI_SERVICES_MWAA_DOCKER_NETWORK: my_project_default配置项环境变量默认值说明FLOCI_SERVICES_MWAA_ENABLEDtrue启用 MWAA 服务FLOCI_SERVICES_MWAA_MOCKfalse仅元数据模式不启动 DockerFLOCI_SERVICES_MWAA_DEFAULT_POSTGRES_IMAGEpostgres:16-alpine元数据库 Docker 镜像独立于 RDS 的配置旋钮FLOCI_SERVICES_MWAA_SUPPORTED_VERSIONS2.10.5,2.9.3,2.8.4逗号分隔的AirflowVersion白名单CreateEnvironment拒绝其他值FLOCI_SERVICES_MWAA_DEFAULT_VERSION2.10.5请求未指定时采用的AirflowVersionFLOCI_SERVICES_MWAA_PROXY_BASE_PORT8700Web/CLI 代理端口区间起始值第一个环境占用FLOCI_SERVICES_MWAA_PROXY_MAX_PORT8799代理端口区间上界含FLOCI_SERVICES_MWAA_DATA_PATH./data/mwaa宿主机数据目录根FLOCI_SERVICES_MWAA_DOCKER_NETWORK未设置Postgres/Airflow 容器所在 Docker 网络未设置时回退到全局FLOCI_SERVICES_DOCKER_NETWORK再回退到 Floci 自身网络FLOCI_SERVICES_MWAA_KEEP_RUNNING_ON_SHUTDOWNfalseFloci 停止后保留 Postgres/Airflow 容器运行FLOCI_SERVICES_MWAA_DAG_SYNC_INTERVAL_SECONDS30轮询 S3 以同步 DAG 与requirements.txt变更的间隔FLOCI_SERVICES_MWAA_INSTALL_REQUIREMENTStrueRequirementsS3Path内容变化时执行pip install -r这些配置定义于 EmulatorConfig.java 的MwaaServiceConfig接口。另外容器与卷命名支持通过FLOCI_DOCKER_RESOURCE_NAMESPACE实现命名空间感知与所有 Docker 后端服务一致经ContainerStorageHelper.dockerName处理见 dbContainerName/airflowContainerName。DAG 与依赖同步每个dag-sync-interval-seconds周期Floci 会列出环境DagS3Path前缀下、SourceBucketArn所指 bucket 中的对象把内容变化的文件复制进容器/opt/airflow/dags删除已不在 S3 的文件——完全复用 Floci 进程内的 S3 模拟不新增任何面向 AWS 的表面syncDags。同步状态按 “S3 key 相对路径 → ETag” 记录仅同步 ETag 变化的文件。RequirementsS3Path若已设置则在每次 DAG 同步通过时检查其 ETag内容变化即pip install -rsyncRequirementsIfNeeded写入/tmp/mwaa-requirements.txt后以pip install --no-cache-dir -r执行installRequirements。单个 DAG 文件损坏或某次pip install失败只会记日志绝不会导致环境失败——与真实 MWAA 行为一致。DAG 同步轮询器只对AVAILABLE环境生效且单个环境的同步异常不会波及其他环境startDagSyncPoller。ARN 格式arn:aws:airflow:region:accountId:environment/nameARN 解析与 tag 相关操作TagResource/UntagResource/ListTagsForResource由 findByArn 与MwaaService实现的TagHandler接口共同完成tags 直接存储于环境实体。快速开始CLI 示例export AWS_ENDPOINT_URLhttp://localhost:4566 export AWS_DEFAULT_REGIONus-east-1 export AWS_ACCESS_KEY_IDtest export AWS_SECRET_ACCESS_KEYtest # 创建环境 aws mwaa create-environment \ --name my-environment \ --execution-role-arn arn:aws:iam::000000000000:role/mwaa-role \ --source-bucket-arn arn:aws:s3:::my-mwaa-bucket \ --dag-s3-path dags \ --network-configuration SubnetIdssubnet-1,subnet-2,SecurityGroupIdssg-1 \ --airflow-version 2.10.5 # 查询环境轮询 Status 直到 AVAILABLE aws mwaa get-environment --name my-environment # 列出环境 aws mwaa list-environments # 获取 Web 登录令牌 aws mwaa create-web-login-token --name my-environment # 获取 CLI 令牌然后直接调用 bridge aws mwaa create-cli-token --name my-environment curl -s -X POST http://localhost:proxyPort/aws_mwaa/cli \ -H Authorization: Bearer CliToken \ -H Content-Type: text/plain \ --data-raw dags list # 打标签 aws mwaa tag-resource \ --resource-arn arn:aws:airflow:us-east-1:000000000000:environment/my-environment \ --tags envdev,teamplatform # 删除环境 aws mwaa delete-environment --name my-environmentget-environment返回的环境对象中Status会经历CREATING→AVAILABLE或失败时的CREATE_FAILEDWebserverUrl指向代理地址max-workers、min-workers、schedulers等字段属于 SDK 兼容性回显字段——在 LocalExecutor 下无实际扩缩容语义但会被原样往返保存避免读取它们的 SDK/CLI 报错见 Environment.java 的注释。Java SDK 示例MwaaClient mwaa MwaaClient.builder() .endpointOverride(URI.create(http://localhost:4566)) .region(Region.US_EAST_1) .credentialsProvider(StaticCredentialsProvider.create( AwsBasicCredentials.create(test, test))) .build(); // 创建环境 CreateEnvironmentResponse created mwaa.createEnvironment(r - r .name(my-environment) .executionRoleArn(arn:aws:iam::000000000000:role/mwaa-role) .sourceBucketArn(arn:aws:s3:::my-mwaa-bucket) .dagS3Path(dags) .networkConfiguration(n - n .subnetIds(subnet-1, subnet-2) .securityGroupIds(sg-1)) .airflowVersion(2.10.5)); // 查询环境 GetEnvironmentResponse described mwaa.getEnvironment(r - r .name(my-environment)); System.out.println(described.environment().statusAsString()); // AVAILABLE // 列出环境 ListString names mwaa.listEnvironments(r - {}).environments(); // 打标签 mwaa.tagResource(r - r .resourceArn(created.arn()) .tags(Map.of(team, platform))); // 删除环境 mwaa.deleteEnvironment(r - r.name(my-environment));仓库的 SDK 兼容性测试如 sdk-test-java 目录体系进一步印证了多语言客户端对该 API 的覆盖MWAA 自身的单元与集成测试位于 src/test/java/io/github/hectorvent/floci/services/mwaa/例如 MwaaEnvironmentManagerTest.java 验证了镜像 tag 替换apache/airflow:2.10.5-python3.11、section.key→ 双下划线格式转换、受保护变量跳过等行为可作实现细节的进一步参考。已知边界Not Implemented当前实现存在以下明确缺口PluginsS3Path会被接受并存储于环境但不会拉取或应用——真实 MWAA 通过plugins.zip分发自定义二进制文件的机制尚无对应实现。UpdateEnvironment仅支持元数据字段tags、description、logging-config 开关。修改AirflowVersion、AirflowConfigurationOptions、RequirementsS3Path路径变更或StartupScriptS3Path会抛出ValidationException需要删除并重建环境——因为这些变更本质要求重建运行中的容器超出了原地更新的范围updateEnvironment。注意RequirementsS3Path指向的文件内容变更仍会通过 DAG 同步自动生效只是不能更换路径。CreateWebLoginToken返回桩令牌stub token不执行真实的 Airflow FAB/JWT SSO 握手返回结构包含WebToken、WebServerHostname、IamIdentityassumed-role/floci-local/user与AirflowIdentityadmincreateWebLoginToken。无按组件区分的MWAA_AIRFLOW_COMPONENT环境变量——真实 MWAA 会在 webserver、scheduler、worker 上分别运行启动脚本而 Floci 将 webserver 与 scheduler 同置一个容器无组件之分。无triggerer进程——使用 deferrable operator 的 DAG 只会记录警告因为只有 scheduler 与 webserver 在运行。Docker 容器与卷名称经由FLOCI_DOCKER_RESOURCE_NAMESPACE命名空间感知与所有 Docker 后端服务一致但不按 AWS 账户 id 隔离——两个账户创建同名环境可能冲突。这是所有 Docker 后端服务EKS、RDS 等共有的既有特性并非 MWAA 独有。Floci 重启后从持久存储加载的环境不会自动重连其 Docker 容器——同样是当前所有 Docker 后端服务共有的缺口非 MWAA 特有。小结Floci 的 MWAA 服务在本地完整复刻了 Amazon MWAA 的 API 表面与关键运行时语义REST-JSON 协议与双重签名注册、双容器Postgres LocalExecutor Airflow架构、真实生效的版本与配置选项、S3 驱动的 DAG/依赖同步、启动脚本的门槛语义与受保护变量、以及统一入口的 Web/CLI 代理。理解这些机制后你可以用它进行 Airflow DAG 的端到端本地开发、CI 中的契约测试以及在 mock 模式下仅验证 API 形状的轻量测试。【免费下载链接】flociLight, fluffy, and always free - The AWS Local Emulator alternative项目地址: https://gitcode.com/gh_mirrors/fl/floci创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表