# AirFlow-Tutorial

- https://my.oschina.net/u/2306127/blog/1843515


In [1]:
"""
Code that goes along with the Airflow tutorial located at:
https://github.com/airbnb/airflow/blob/master/airflow/example_dags/tutorial.py
"""
from airflow import DAG
from airflow.operators.bash_operator import BashOperator
from datetime import datetime, timedelta


default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2015, 6, 1),
    'email': ['airflow@example.com'],
    'email_on_failure': False,
    'email_on_retry': False,
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
    # 'queue': 'bash_queue',
    # 'pool': 'backfill',
    # 'priority_weight': 10,
    # 'end_date': datetime(2016, 1, 1),
}

dag = DAG('tutorial', default_args=default_args)

# t1, t2 and t3 are examples of tasks created by instantiating operators
t1 = BashOperator(
    task_id='print_date',
    bash_command='date',
    dag=dag)

t2 = BashOperator(
    task_id='sleep',
    bash_command='sleep 5',
    retries=3,
    dag=dag)

templated_command = """
    {% for i in range(5) %}
        echo "{{ ds }}"
        echo "{{ macros.ds_add(ds, 7)}}"
        echo "{{ params.my_param }}"
    {% endfor %}
"""

t3 = BashOperator(
    task_id='templated',
    bash_command=templated_command,
    params={'my_param': 'Parameter I passed in'},
    dag=dag)

t2.set_upstream(t1)
t3.set_upstream(t1)

[2018-07-21 10:39:46,541] {__init__.py:57} INFO - Using executor SequentialExecutor


In [3]:
?t2

[0;31mType:[0m        BashOperator
[0;31mString form:[0m <Task(BashOperator): sleep>
[0;31mFile:[0m        /srv/conda/lib/python3.6/site-packages/airflow/operators/bash_operator.py
[0;31mDocstring:[0m  
Execute a Bash script, command or set of commands.

:param bash_command: The command, set of commands or reference to a
    bash script (must be '.sh') to be executed.
:type bash_command: string
:param xcom_push: If xcom_push is True, the last line written to stdout
    will also be pushed to an XCom when the bash command completes.
:type xcom_push: bool
:param env: If env is not None, it must be a mapping that defines the
    environment variables for the new process; these are used instead
    of inheriting the current process environment, which is the default
    behavior. (templated)
:type env: dict
:type output_encoding: output encoding of bash command


In [4]:
?t3

[0;31mType:[0m        BashOperator
[0;31mString form:[0m <Task(BashOperator): templated>
[0;31mFile:[0m        /srv/conda/lib/python3.6/site-packages/airflow/operators/bash_operator.py
[0;31mDocstring:[0m  
Execute a Bash script, command or set of commands.

:param bash_command: The command, set of commands or reference to a
    bash script (must be '.sh') to be executed.
:type bash_command: string
:param xcom_push: If xcom_push is True, the last line written to stdout
    will also be pushed to an XCom when the bash command completes.
:type xcom_push: bool
:param env: If env is not None, it must be a mapping that defines the
    environment variables for the new process; these are used instead
    of inheriting the current process environment, which is the default
    behavior. (templated)
:type env: dict
:type output_encoding: output encoding of bash command


In [5]:
?t1

[0;31mType:[0m        BashOperator
[0;31mString form:[0m <Task(BashOperator): print_date>
[0;31mFile:[0m        /srv/conda/lib/python3.6/site-packages/airflow/operators/bash_operator.py
[0;31mDocstring:[0m  
Execute a Bash script, command or set of commands.

:param bash_command: The command, set of commands or reference to a
    bash script (must be '.sh') to be executed.
:type bash_command: string
:param xcom_push: If xcom_push is True, the last line written to stdout
    will also be pushed to an XCom when the bash command completes.
:type xcom_push: bool
:param env: If env is not None, it must be a mapping that defines the
    environment variables for the new process; these are used instead
    of inheriting the current process environment, which is the default
    behavior. (templated)
:type env: dict
:type output_encoding: output encoding of bash command


In [10]:
dag.tree_view()

<Task(BashOperator): sleep>
    <Task(BashOperator): print_date>
<Task(BashOperator): templated>
    <Task(BashOperator): print_date>
