Python - Airflow再会

1. 前言

近期计划做一个任务调度系统,于是,重拾airflow,借机深入学习下。
主要调研和测试具体使用方法、能否满足我们的项目需求,以及可能存在哪些坑。

不了解airflow的朋友,可以参考我的上篇文章:
Python - Airflow任务调度系统初识

简单回顾一下两组关键名词:

  • Dag -> DagRun(Dag Instance)
  • Operator -> Task -> Task Instance

Core Concepts — Airflow Documentation (apache.org)

2. 使用方法

2.1. 编写Dag文件,并测试

image.png
  • 梳理实际用户需求,是否存在选择分支,是否存在人工审核步骤,是否存在任务重跑,等
  • 通过不同的Operator,实现一系列Task,以及任务的入参/出参、任务间通信、任务间依赖等(附件采用PythonOperator,实现了一个ETL的过程,采用xcom通信)
  • 在Test环境,测试Dag文件,是否符合预期

比较常用的两个Operator:

  • PythonOperator:可以通过args、kwargs设置Task运行时参数,通过xcom完成Task间通信
  • BranchOperator:可以在运行时,决定走哪个任务分支,比如时间、服务器资源、上个任务的输出、当前Dag的状态等


不采用json、yaml描述dag,便于airflow实现功能更加丰富的任务流,以及更自由的配置。
但是,如果让非软件开发人员编写Dag文件,无疑压力巨大。因此,需要人工审核。
如果是个2C软件,则可能需要jinja2来实时render一个dag.py文件。

Best Practices — Airflow Documentation (apache.org)

2.2. 上传Dag文件到scheduler和worker节点

image.png

Dag文件目录,需要在scheduler和worker机器上各存储一份,这个是由Airflow系统架构决定。
采用的方法,可以:

  • nas盘,挂载到同一个集群的不同机器上(最好在一个局域网内)
  • ansible手动push到服务器节点
  • 在服务器节点,启动脚本,实时或定时同步oss/gitlab文件到本地


如果Dag文件很多的话,文件的管理,可能需要namespace、mainfest等。
另外,是否需要先同步到worker节点,再同步到scheduler节点。

2.3. trigger dag

  • 自动trigger,例如cron、upstream-task
  • 手动trigger,通过命令行、airflow browser、airflow restful api

2.4. backfill

如果dag的下面两个参数为True,则dag一启动,则会查看过去的dagrun,如果存在未执行的dagrun,则回填。
所以,一般情况下,这两个参数需要手动设置为False。

depends_on_past=False
catchup=False

3. 日常维护

  • metadata数据过多,会导致系统响应缓慢,通过airflow browser就可以明显感觉到,需要手动清理db
  • 任务告警与相关责任人的管理
  • 任务流上线流程,需要严格测试,并增加权限控制
  • 任务流在生产环境运行时,需要引入pool的概念,来隔离不同任务流之间的资源竞争
  • Scheduler只能跑一份,所以存在单点故障风险。需要借助第三方组件airflow-scheduler-failover-controller实现Scheduler的高可用

4. 参考

5. 附,ETL脚本

import json
import pendulum
from airflow import DAG
from airflow.operators.python import PythonOperator


with DAG(
    "cp_pipeline_1",
    default_args={"retries": 2},
    description="computing pipeline",
    schedule=None,
    start_date=pendulum.datetime(2021, 1, 1, tz="UTC"),
    depends_on_past=False,
    catchup=False,
    tags=["computing-platform"],
) as dag:

    def extract(*args, **kwargs):
        print('extract kwargs', kwargs)
        ti = kwargs["ti"]
        data_string = '{"1001": 301.27, "1002": 433.21, "1003": 502.22}'
        ti.xcom_push("order_data", data_string)
        return {'extract', 'over'}  # return值也会进入到xcom,downstream的task也可以通过xcom获取

    def transform(*args, **kwargs):
        print('transform kwargs', kwargs)
        ti = kwargs["ti"]
        extract_data_string = ti.xcom_pull(task_ids="extract", key="order_data")
        order_data = json.loads(extract_data_string)

        total_order_value = 0
        for value in order_data.values():
            total_order_value += value

        total_value = {"total_order_value": total_order_value}
        total_value_json_string = json.dumps(total_value)
        ti.xcom_push("total_order_value", total_value_json_string)

    def load(*args, **kwargs):
        print('load kwargs', kwargs)
        ti = kwargs["ti"]
        total_value_string = ti.xcom_pull(task_ids="transform", key="total_order_value")
        total_order_value = json.loads(total_value_string)
        print('load result', total_order_value)

    extract_task = PythonOperator(
        task_id="extract",
        python_callable=extract,
        op_args=[],
        op_kwargs=dict(a1=1,b1='bb',c1=False),
    )

    transform_task = PythonOperator(
        task_id="transform",
        python_callable=transform,
        op_args=[],
        op_kwargs=dict(a2=1,b2='bb',c2=False),
    )

    load_task = PythonOperator(
        task_id="load",
        python_callable=load,
        op_args=[],
        op_kwargs=dict(a3=1,b3='bb',c3=False),
    )

    extract_task >> transform_task >> load_task
最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 205,033评论 6 478
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 87,725评论 2 381
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 151,473评论 0 338
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 54,846评论 1 277
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 63,848评论 5 368
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 48,691评论 1 282
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 38,053评论 3 399
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 36,700评论 0 258
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 42,856评论 1 300
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 35,676评论 2 323
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 37,787评论 1 333
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 33,430评论 4 321
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 39,034评论 3 307
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 29,990评论 0 19
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 31,218评论 1 260
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 45,174评论 2 352
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 42,526评论 2 343

推荐阅读更多精彩内容