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
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
|
/*
Double is an item
Python3 is a state
String is a result
Int64 is a serialized
*/
$script = @@#py
import json
class State:
def __init__(self, state):
self.state = state
def create(item):
return State(item)
def add(state, item):
return State(state.state + item)
def merge(state_a, state_b):
return State(state_a.state + state_b.state)
def get_result(state):
return str(state.state)
def serialize(state):
return int(state.state)
def deserialize(serialized):
return State(float(serialized))
@@;
$create = Python3::create(Callable<(Double)->Resource<Python3>>, $script);
$add = Python3::add(Callable<(Resource<Python3>,Double)->Resource<Python3>>, $script);
$merge = Python3::merge(Callable<(Resource<Python3>,Resource<Python3>)->Resource<Python3>>, $script);
$get_result = Python3::get_result(Callable<(Resource<Python3>)->String>, $script);
$serialize = Python3::serialize(Callable<(Resource<Python3>)->Int64>, $script);
$deserialize = Python3::deserialize(Callable<(Int64)->Resource<Python3>>, $script);
$default = '';
$factory = AggregationFactory(
"UDAF",
$create,
$add,
$merge,
$get_result,
$serialize,
$deserialize,
$default
);
$input = SELECT * FROM AS_TABLE([
<|key: 1, item: 1.0|>,
<|key: 1, item: 2.0|>,
<|key: 1, item: -2.0|>,
<|key: 2, item: 1.0|>,
<|key: 2, item: 1.0|>,
<|key: 3, item: 1.0|>,
<|key: 3, item: 2.0|>,
]);
$states = SELECT * FROM AS_TABLE([
<|key: 1, state: $serialize($create( 1))|>,
<|key: 1, state: $serialize($create( 2))|>,
<|key: 1, state: $serialize($create(-2))|>,
<|key: 2, state: $serialize($create( 1))|>,
<|key: 2, state: $serialize($create( 1))|>,
<|key: 3, state: $serialize($create( 1))|>,
<|key: 3, state: $serialize($create( 2))|>,
]);
$p = SELECT key, AGGREGATE_BY(item, $factory) AS state
FROM $input
GROUP BY key WITH combine;
$p = PROCESS $p;
SELECT * FROM $p;
$p = SELECT key, AGGREGATE_BY(item, $factory) AS result
FROM $input
GROUP BY key with finalize;
$p = PROCESS $p;
SELECT * FROM $p;
$p = SELECT key, AGGREGATE_BY(state, $factory) AS state
FROM $states
GROUP BY key WITH combinestate;
$p = PROCESS $p;
SELECT * FROM $p;
$p = SELECT key, AGGREGATE_BY(state, $factory) AS state
FROM $states
GROUP BY key WITH mergestate;
$p = PROCESS $p;
SELECT * FROM $p;
$p = SELECT key, AGGREGATE_BY(state, $factory) AS result
FROM $states
GROUP BY key WITH mergefinalize;
$p = PROCESS $p;
SELECT * FROM $p;
$p = SELECT key, AGGREGATE_BY(state, $factory) AS result
FROM (select key, just(state) as state from $states)
GROUP BY key with mergemanyfinalize;
$p = PROCESS $p;
SELECT * FROM $p;
|