【发布时间】: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 的用途吗?是表名吗?
-
我必须为艺术家、歌曲等不同类型的实体运行相同的代码。这里我刚刚提到了艺术家类型