【问题标题】:Templating `bucket_key` in S3KeySensor in Apache Airflow在 Apache Airflow 的 S3KeySensor 中模板化 `bucket_key`
【发布时间】:2018-05-24 16:28:01
【问题描述】:

气流版本:1.9.0

在气流 dag 文件中,我有一个名为 run_query 的 PythonOperator 任务,它在其 python_callable 函数中设置以下 xcom 变量:

kwargs['ti'].xcom_push(key='query_result_loc', value=query_result_loc)

在同一天,我有一个 S3KeySensor 任务,它使用上述位置作为它的 bucket_key 参数:

S3KeySensor(task_id = 'check_file_in_s3',
                 bucket_key = '{{  ti.xcom_pull(task_ids="run_query",key="query_result_loc")  }}' ,
                 bucket_name = None,
                 wildcard_match = False,
                 poke_interval=60,
                 timeout=1200,
                 dag = dag
                 )

现在,当我运行 dag(在测试模式或 trigger_dag 模式下)时,S3KeySensor 抱怨缺少 bucket_name,它来自 S3KeySensor definition 中的这段代码:

    class S3KeySensor(BaseSensorOperator):
    """
    Waits for a key (a file-like instance on S3) to be present in a S3 bucket.
    S3 being a key/value it does not support folders. The path is just a key
    a resource.

    :param bucket_key: The key being waited on. Supports full s3:// style url
        or relative path from root level.
    :type bucket_key: str
    :param bucket_name: Name of the S3 bucket
    :type bucket_name: str
    :param wildcard_match: whether the bucket_key should be interpreted as a
        Unix wildcard pattern
    :type wildcard_match: bool
    :param aws_conn_id: a reference to the s3 connection
    :type aws_conn_id: str
    """
    template_fields = ('bucket_key', 'bucket_name')

    @apply_defaults
    def __init__(
                self, bucket_key,
                bucket_name=None,
                wildcard_match=False,
                aws_conn_id='aws_default',
                *args, **kwargs):
            super(S3KeySensor, self).__init__(*args, **kwargs)
            # Parse
            if bucket_name is None:
                parsed_url = urlparse(bucket_key)
                if parsed_url.netloc == '':
                    raise AirflowException('Please provide a bucket_name')
                else:
                    bucket_name = parsed_url.netloc
                    if parsed_url.path[0] == '/':
                        bucket_key = parsed_url.path[1:]
                    else:
                        bucket_key = parsed_url.path
            self.bucket_name = bucket_name
            self.bucket_key = bucket_key

看起来模板在这个阶段没有被渲染。

如果我注释掉 if 块,它工作正常。这是错误还是模板字段的错误使用?

更新,基于@kaxil 的评论:

  • 如果没有提供 bucket_name 并且“if”块未注释,气流甚至无法检测到 dag。在 UI 上,我看到了这个错误:Broken DAG: [/XXXX/YYYY/project_airflow.py] Please provide provide a bucket_name
  • 没有提供bucket_name,但对'if'块进行了以下修改(参见删除了if parsed_url.netloc == ''检查),它工作正常:

    if bucket_name is None:
        parsed_url = urlparse(bucket_key)
        bucket_name = parsed_url.netloc
        if parsed_url.path[0] == '/':
            bucket_key = parsed_url.path[1:]
        else:
            bucket_key = parsed_url.path
    
  • 如果提供了 bucket_name,它可以在 Rendered Template 选项卡下使用 bucket_key 和 bucket_name 的呈现值正常工作。

【问题讨论】:

  • 在 Airflow Web UI 中,您是否在页面 Rendered Template 中看到任何用于这些任务的任何渲染模板?可能是检查模板渲染的好方法。如果该字段被模板化并且变量被填充,它应该可以工作。
  • @tobi6 我确实在 dag run 的 UI 上的“渲染模板”选项卡下看到了正确的替换。
  • @crackjack 你能用 Airflow Web UI 的截图更新你的问题吗? -> 两种情况下的渲染模板 i) 使用存储桶名称,ii) 没有存储桶名称。
  • @kaxil 我不能在不屏蔽渲染值的情况下发布屏幕截图,这会违背你的提问目的。相反,我试图在我的问题中解决您的评论。希望这会有所帮助。
  • @crackjack 你能提供bucket_key的值吗?如它是s3://BUCKET_NAME/OBJECT_NAME.csv 格式还是s3:///BUCKET_NAME/OBJECT_NAME.csv ???不要提供实际值,只是让我知道格式。

标签: jinja2 airflow


【解决方案1】:

您现在可能已经解决了这个问题,但这是因为您没有为 task_ids 提供数组。格式应为task_ids=["run_query"]。将其更改为此解决了我的问题。

【讨论】:

    【解决方案2】:

    正如您在更新报价时正确注意到的那样,存储桶密钥验证发生在 jinja 模板之前。

    我发现的唯一解决方案是不使用仅提供存储桶密钥的功能。如果要使用模板,请提供存储桶名称和存储桶密钥。在这种情况下,验证将不会发生,字段将被模板化。

    【讨论】:

      猜你喜欢
      • 2018-07-19
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-02-07
      • 2023-01-27
      • 2018-07-01
      • 2018-03-14
      相关资源
      最近更新 更多