【问题标题】:How to add database partition based on airflow timestamp如何根据气流时间戳添加数据库分区
【发布时间】:2021-01-08 14:08:08
【问题描述】:

在这里,我尝试从开源获取数据,并根据气流时间戳将其作为分区添加到表中。但它会引发气流异常。

def partition_sql(entity_type):
            sql = """
                ALTER TABLE db.table
                ADD IF NOT EXISTS PARTITION (airflow_ts='{{ts}}')
                LOCATION 's3://db/table/update/airflow_ts={{ts}}';
            """
        return sql
    
with DAG(parameters)as dag:

    update = DockerOperator(
        task_id='update',
        cmd = 'python script.py 's3://db/table1/update/airflow_ts={{ts}}'
     )
  partition = AWSAthenaOperator(
        task_id='partition',
        query=partition_sql("artist"),
    )


 update >>partition

【问题讨论】:

  • 你需要什么格式的时间戳?
  • ISO 格式甚至只是日期都可以
  • 你能解释一下 entity_type 的用途吗?是表名吗?
  • 我必须为艺术家、歌曲等不同类型的实体运行相同的代码。这里我刚刚提到了艺术家类型

标签: airflow-scheduler airflow


【解决方案1】:

由于 cmdquery 字段均已模板化,因此应该可以:

items = ["artist"] #add more tables to be created dynamically
with DAG(
        dag_id="dag_name",
        default_args=default_args,
) as dag:
    for item in items:
        command = f'python script.py 's3://db/{item}/update/airflow_ts={{ ds }} '
        update = DockerOperator(
            task_id=f'update_table_{item}',
            cmd=command
        )
        sql = f'ALTER TABLE db.{item} ADD IF NOT EXISTS PARTITION (airflow_ts={{ ds }}) LOCATION s3://db/table/update/airflow_ts={{ ds }};'
        partition = AWSAthenaOperator(
            task_id=f'partition_table_{item}',
            query=sql
        )

您可以从{{ ds }} 更改为您喜欢的任何其他日期格式。您可以在macros 页面查看可用的格式或自定义一种格式。

请注意,您不必在此处将 SQL 保存在代码中。您可以按照here 的说明将其保存在.sql 文件中

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-04-04
    • 2019-06-06
    • 1970-01-01
    • 1970-01-01
    • 2015-12-03
    • 2018-04-28
    相关资源
    最近更新 更多