/
sequence.py
43 lines (29 loc) · 1.1 KB
/
sequence.py
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
""" A sequence of tasks incrementing a number
This workflow arranges 3 tasks in a row. Each task calls the same callable,
which increments a number and prints the current number and some information.
"""
from lightflow.models import Dag
from lightflow.tasks import PythonTask
# the callback function for all tasks
def inc_number(data, store, signal, context):
print('Task {task_name} being run in DAG {dag_name} '
'for workflow {workflow_name} ({workflow_id}) '
'on {worker_hostname}'.format(**context.to_dict()))
if 'value' not in data:
data['value'] = 0
data['value'] = data['value'] + 1
print('This is task #{}'.format(data['value']))
# create the main DAG
d = Dag('main_dag')
# create the 3 tasks that increment a number
task_1 = PythonTask(name='task_1',
callback=inc_number)
task_2 = PythonTask(name='task_2',
callback=inc_number)
task_3 = PythonTask(name='task_3',
callback=inc_number)
# set up the graph of the DAG as a linear sequence of tasks
d.define({
task_1: task_2,
task_2: task_3
})