MWAA中Airflow PythonOperator生成的PEM文件无法被SSHOperator访问如何解决
问题根因
你碰到的报错来自两个必然触发的逻辑问题,和MWAA的运行机制直接相关:
- 现有
task2的代码在DAG解析阶段就会直接执行retrieve_pem_key_from_aws_secrets_manager(),根本不会等task1运行结束再拉密钥。DAG文件每次被调度器、worker加载解析时都会跑这段初始化代码,生成的临时文件存在于解析进程所在容器,和task1实际执行时生成的文件完全是两个东西。 - MWAA基于Fargate运行worker任务,默认每个任务跑在独立隔离容器中,本地文件系统是临时存储,任务结束容器销毁后文件会被直接清除,跨任务完全无法共享本地
/tmp路径下的文件,就算用XCom传递临时文件路径,在另一个运行容器里也找不到对应文件。
零环境改动的代码级解决方案
完全不需要修改MWAA后端配置、不需要搭建共享存储,直接把密钥拉取和SSH连接的逻辑放到同一个任务执行上下文里即可,两种实现方式可按需选择:
方案1:单PythonOperator封装完整流程(最稳妥)
不拆分两个任务,直接在同一个Python可调用对象里完成拉取PEM密钥、写入临时文件、初始化SSH连接、执行命令的全流程,整个流程在同一个容器进程内跑完,临时文件在进程结束后自动销毁,不存在跨容器共享问题。
示例代码:
import tempfile import os from airflow.providers.amazon.aws.hooks.secrets_manager import SecretsManagerHook from airflow.providers.ssh.hooks.ssh import SSHHook from airflow.operators.python import PythonOperator def ssh_with_pem_from_sm(**kwargs): # 从Secrets Manager拉取PEM密钥 sm_hook = SecretsManagerHook() sm_client = sm_hook.get_conn() secret = sm_client.get_secret_value(SecretId='<YOUR_SECRET_ID>') pem_content = secret["SecretString"] # 写入临时文件,设置仅当前用户可读权限(SSH强制要求私钥权限不能过宽) with tempfile.NamedTemporaryFile(mode='w+', encoding='UTF-8', delete=False) as pem_file: pem_file.write(pem_content) pem_path = pem_file.name os.chmod(pem_path, 0o600) try: # 初始化SSH连接执行命令 ssh_hook = SSHHook( remote_host='<YOUR_EC2_HOST>', username='ec2-user', key_file=pem_path, timeout=10 ) exit_code, output = ssh_hook.run( command='echo Hello', combine_stderr=True ) print(f"SSH执行输出: {output}") if exit_code != 0: raise Exception(f"SSH命令执行失败,退出码: {exit_code}") finally: # 用完即删临时密钥,避免残留 if os.path.exists(pem_path): os.remove(pem_path) # 任务定义 ssh_task = PythonOperator( task_id='run_ssh_command', python_callable=ssh_with_pem_from_sm )
方案2:自定义SSHOperator保留任务拆分逻辑
如果必须保留两个任务的拆分结构,不要在DAG初始化阶段硬编码key_file参数,通过自定义算子的生命周期钩子,在SSHOperator执行前的当前worker容器内拉取密钥、写入临时文件,执行完成后自动清理:
import tempfile import os from airflow.providers.amazon.aws.hooks.secrets_manager import SecretsManagerHook from airflow.providers.ssh.hooks.ssh import SSHHook from airflow.providers.ssh.operators.ssh import SSHOperator class PemFromSMSSHOperator(SSHOperator): def pre_execute(self, context): # 任务执行前,在当前worker容器内拉取密钥生成临时文件 sm_hook = SecretsManagerHook() sm_client = sm_hook.get_conn() secret = sm_client.get_secret_value(SecretId='<YOUR_SECRET_ID>') pem_content = secret["SecretString"] with tempfile.NamedTemporaryFile(mode='w+', encoding='UTF-8', delete=False) as pem_file: pem_file.write(pem_content) self.pem_path = pem_file.name os.chmod(self.pem_path, 0o600) # 动态初始化SSH Hook,绑定刚生成的临时密钥 self.ssh_hook = SSHHook( remote_host='<YOUR_EC2_HOST>', username='ec2-user', key_file=self.pem_path, timeout=10 ) super().pre_execute(context) def post_execute(self, context, result=None): # 执行完成后删除临时密钥 if hasattr(self, 'pem_path') and os.path.exists(self.pem_path): os.remove(self.pem_path) super().post_execute(context, result) # 任务定义(不需要提前定义task1,算子内部会自行拉取密钥) task2 = PemFromSMSSHOperator( task_id='test_ssh_connectivity', command='echo Hello' )
不推荐的方案说明
- 不要尝试用XCom传递文件路径:跨Fargate容器的本地文件系统完全隔离,传递路径后在另一个容器内找不到对应文件。
- 不需要专门将Secrets Manager配置为Airflow连接后端:当前场景只需要拉取单个PEM密钥,用现有SecretsManagerHook即可满足需求,修改后端配置属于多余操作,还需要调整MWAA环境配置,不符合尽量少改环境的要求。
- 不要用S3做中间存储传递PEM密钥:会额外增加S3权限配置、文件中转逻辑,反而提升密钥泄露风险,安全性不如同进程直接使用临时文件的方案。
内容的提问来源于stack exchange,提问作者rk92
相关产品推荐
相关产品推荐

