Airflow GCSToS3Operator keep_directory_structure参数未渲染生效问题
问题根因
该问题为Apache Airflow Amazon Provider早期版本GCSToS3Operator的已知实现bug,触发场景和Composer 2.0.18内置的provider包版本直接相关:
- Composer 2.0.18预装的
apache-airflow-providers-amazon为3.0.0之前的旧版本,对应算子的__init__方法未将replace、keep_directory_structure两个参数正确透传给底层执行逻辑,也未将两个参数加入模板渲染列表,最终导致配置的参数完全不生效,算子始终按照硬编码的默认逻辑拼接路径。 - 旧版本硬编码逻辑会强制将GCS端的完整对象key(包含prefix部分)拼接在目标S3路径后,这也是首次配置时目标路径自动带上
legacy/action/20220629层级的原因;手动修改keep_directory_structure=False时,底层逻辑错误重复拼接两次prefix路径,才会出现路径中包含两层相同prefix的异常。
修复方案
按实际环境约束选择以下任意一种方案即可:
方案1:升级Amazon Provider包(推荐)
在Composer环境中将apache-airflow-providers-amazon升级到4.0.0及以上稳定版本即可,新版本已经修复参数透传、路径拼接逻辑的相关bug,升级后原有算子配置可正常生效:
- 设置
keep_directory_structure=False时,匹配到的csv文件会直接平铺写入dest_s3_key指定的s3a://action/daily/路径,不会额外拼接GCS端的prefix层级,符合预期。 - 升级可直接在Composer控制台的PyPI包管理页面操作,注意选择和Airflow 2.2.5兼容的版本即可,4.1.0版本兼容性最优,不会出现依赖冲突。
方案2:无升级条件的临时规避
如果暂时无法调整环境内置包版本,可自定义子类重写算子逻辑,绕开参数透传bug,参考代码如下:
from airflow.providers.amazon.aws.transfers.gcs_to_s3 import GCSToS3Operator class FixedGCSToS3Operator(GCSToS3Operator): template_fields = ('dest_s3_key', 'prefix', 'bucket') def __init__(self, *, replace=False, keep_directory_structure=True, **kwargs): super().__init__(**kwargs) self.replace = replace self.keep_directory_structure = keep_directory_structure def execute(self, context): from airflow.providers.google.cloud.hooks.gcs import GCSHook s3_hook = self.get_hook() gcs_hook = GCSHook(gcp_conn_id=self.gcp_conn_id, impersonation_chain=self.impersonation_chain) bucket = gcs_hook.get_bucket(self.bucket) prefix = self.prefix or '' blobs = bucket.list_blobs(prefix=prefix, delimiter=self.delimiter) for blob in blobs: if not blob.name.endswith(self.delimiter): continue if self.keep_directory_structure: dest_key = f"{self.dest_s3_key}{blob.name}" else: dest_key = f"{self.dest_s3_key}{blob.name.split('/')[-1]}" s3_hook.load_fileobj( fileobj=blob.open('rb'), key=dest_key, replace=self.replace ) # 替换原有GCSToS3Operator调用即可 gcs_to_s3 = FixedGCSToS3Operator( task_id="gcs_to_s3", bucket="gcs_outbound", prefix="legacy/action/20220629", delimiter=".csv", dest_aws_conn_id="S3-action-outbound", dest_s3_key="s3a://action/daily/", replace=False, keep_directory_structure=False, )
校验方式
配置完成后可在Airflow UI任务实例详情页查看Rendered Template,确认replace、keep_directory_structure参数已正常传入,再触发测试任务即可,不会再出现路径多拼接、重复拼接的问题。
注意:
dest_s3_key配置时末尾需保留斜杠,否则最后一级路径会被识别为文件名前缀,导致写入路径不符合预期。
内容的提问来源于stack exchange,提问作者urvish patel
相关产品推荐
相关产品推荐

