-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy path1_cls_luigi_pipeline.py
More file actions
58 lines (45 loc) · 1.78 KB
/
Copy path1_cls_luigi_pipeline.py
File metadata and controls
58 lines (45 loc) · 1.78 KB
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
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
import luigi
from cls.fcl import FiniteCombinatoryLogic
from cls.subtypes import Subtypes
from cls_luigi.inhabitation_task import LuigiCombinator, ClsParameter, RepoMeta
from utils import print_tree
class TaskA(luigi.Task, LuigiCombinator):
abstract = False
def output(self):
return luigi.LocalTarget("output/taskA_output.txt")
def run(self):
with self.output().open('w') as f:
f.write("Task A completed")
class TaskB(luigi.Task, LuigiCombinator):
abstract = False
task_a = ClsParameter(tpe=TaskA.return_type())
def requires(self):
return self.task_a()
def output(self):
return luigi.LocalTarget("output/taskB_output.txt")
def run(self):
with self.input().open() as input_file, self.output().open('w') as output_file:
data = input_file.read()
output_file.write("Task B completed with input: " + data)
if __name__ == '__main__':
target = TaskB.return_type()
repository = RepoMeta.repository
fcl = FiniteCombinatoryLogic(repository, Subtypes(RepoMeta.subtypes))
inhabitation_result = fcl.inhabit(target)
max_tasks_when_infinite = 10
actual = inhabitation_result.size()
max_results = max_tasks_when_infinite
if not actual is None or actual == 0:
max_results = actual
validator = RepoMeta.get_unique_abstract_task_validator()
results = [t() for t in inhabitation_result.evaluated[0:max_results]
if validator.validate(t())]
for r in results:
print(print_tree(r))
if results:
print("Number of results", max_results)
print("Number of results after filtering", len(results))
print("Run Pipelines")
no_schedule_error = luigi.build(results, local_scheduler=True)
else:
print("No results!")