You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.27 20:01:21