summaryrefslogtreecommitdiffstats
path: root/yql/essentials/tests/sql/suites/agg_phases/udaf.yql
blob: e15b7a70af6aa985dc11b75f17edc39d106c4f5a (plain) (blame)
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;