【问题标题】:How to print a unique message in Airflow operator?如何在 Airflow 运算符中打印唯一消息?
【发布时间】:2018-09-03 06:18:19
【问题描述】:

在气流中,可以使用简单的print() 或使用建议的记录器here 写入日志

但是,当尝试在运算符内部打印时,这将不起作用。

我有以下代码:

for i in range(5, 0, -1):
    gcs_export_uri_template = ["adstest/{{ macros.ds_format(macros.ds_add(ds, -params.i), '%Y-%m-%d', '%Y/%m/%d') }}/*"]
    update_bigquery = GoogleCloudStorageToBigQueryOperator(
        dag=dag,
        task_id='load_ads_to_BigQuery-{}'.format(i),
        bucket=GCS_BUCKET_ID,
        destination_project_dataset_table=table_name_template,
        source_format='CSV',
        source_objects=gcs_export_uri_template,
        schema_fields=dc(),
        params={'i': i},
        create_disposition='CREATE_IF_NEEDED',
        write_disposition='WRITE_APPEND',
        skip_leading_rows=1,
        google_cloud_storage_conn_id=CONNECTION_ID,
        bigquery_conn_id=CONNECTION_ID
    )

现在说我要打印"My name is load_ads_to_BigQuery-{}".format{i) 如您所见,此打印是每个操作员唯一的。

如果我这样做:

for i in range(5, 0, -1):
    print("My name is load_ads_to_BigQuery-{}".format{i))
    gcs_export_uri_template = ...
    update_bigquery = GoogleCloudStorageToBigQueryOperator(...)

所有 5 个操作员将打印所有 5 个打印件。在我的情况下这是不正确的。 打印必须在GoogleCloudStorageToBigQueryOperator 内。

我该怎么做?

【问题讨论】:

  • 您也可以使用 BaseOperator 中的 log 方法记录 inside 运算符,例如 self.log.info("This was %s", "fun")。所有的打印语句都是在 dag 中完成的,还是你的 for 循环在自定义 Operator 中?
  • @judoole 我不确定我是否理解
  • 您是在 dag 文件中打印,还是在自创模块中打印,在 for 循环中创建 5 个 GoogleCloudStorageToBigQueryOperators?抱歉回复晚了。
  • @judoole 我根本不打印。我想将打印添加到我的代码中。我想在 GoogleCloudStorageToBigQueryOperator 中打印。
  • 我认为这是所有运营商的问题。您是在建议我将我的操作员包装在一个将打印然后执行操作员的函数中。我认为气流不支持这一点:\我提交了一张票issues.apache.org/jira/browse/AIRFLOW-3000

标签: python airflow


【解决方案1】:

一个可能足够的解决方案是创建一个包装类。示例:

class MyGoogleCloudStorageToBigQueryOperator(BaseOperator):
    template_fields = ('bucket', 'source_objects',
                   'schema_object', 'destination_project_dataset_table')

    @apply_defaults
    def __init__(self,
          bucket,
          destination_project_dataset_table,
          source_format,
          source_objects,
          schema_fields,
          params,
          create_disposition,
          write_disposition,
          skip_leading_rows,
          google_cloud_storage_conn_id,
          bigquery_conn_id,
          *args,
          **kwargs):
     super(InfogroupFilestreamOperator, self).__init__(*args, **kwargs)
     self.bucket= bucket
     self.destination_project_dataset_table=destination_project_dataset_table
     ...

  def execute(self, context):
      self.log.info("My name is %s", self.task_id)
      hook=GoogleCloudStorageToBigQueryOperator(
            task_id="doesnt_matter",
            bucket=self.bucket,             
            source_format=self.source_format,
            ...
            )

     return hook.execute(context)

MyGoogleCloudStorageToBigQueryOperator 然后可以在你的 for 循环中被实例化,而不是 GoogleCloudStorageToBigQueryOperator

【讨论】:

  • 为什么不class MyGoogleCloudStorageToBigQueryOperator(GoogleCloudStorageToBigQueryOperator): pass,覆盖execute,然后调用super().execute(context)
  • 是的,我想在这种情况下也可以。更好的解决方案。如果您想要一些额外的分数,请随时编辑,或者自己回答。我会投票。
【解决方案2】:

根据@judoole 的回答,您可以直接扩展GoogleCloudStorageToBigQueryOperator

class MyGoogleCloudStorageToBigQueryOperator(LoggingMixin, GoogleCloudStorageToBigQueryOperator):

    def execute(self, context):
        self.logger.info('Inside task {task_id}'.format(task_id=context['task_id']))
        super().execute(context)

更一般地说,您可以编写一个 mixin 类,它会自动为各种运算符执行这种类型的日志记录。

class LogTaskExecutionMixin(object):

    def execute(self, context):
        self.logger.info('Inside task {task_id}'.format(task_id=context['task_id']))
        super().execute(context)

class MyGoogleCloudStorageToBigQueryOperator(
    LogTaskExecutionMixin, LoggingMixin, GoogleCloudStorageToBigQueryOperator
):
    pass

这两种方法的想法是,您定义一个新的运算符,它在执行之前写入一条日志消息,但在其他方面与您扩展的运算符相同。

【讨论】:

    猜你喜欢
    • 2018-12-28
    • 2012-07-31
    • 1970-01-01
    • 1970-01-01
    • 2012-05-05
    • 1970-01-01
    • 1970-01-01
    • 2019-09-22
    • 1970-01-01
    相关资源
    最近更新 更多