Packages
dqe
0.1.20
0.4.15
0.4.14
0.4.13
0.4.12
0.4.11
0.4.10
0.4.9
0.4.8
0.4.7
0.4.6
0.4.5
0.4.4
0.4.3
0.4.2
0.4.1
0.4.0
0.3.25
0.3.24
0.3.23
0.3.22
0.3.21
0.3.20
0.3.19
0.3.18
0.3.17
0.3.16
0.3.15
0.3.14
0.3.13
0.3.12
0.3.11
0.3.10
0.3.9
0.3.8
0.3.7
0.3.6
0.3.5
0.3.4
0.3.3
0.3.2
0.3.1
0.2.2
0.2.1
0.2.0
0.1.36
0.1.35
0.1.34
0.1.33
0.1.22
0.1.21
0.1.20
0.1.10
0.1.9
DalmatinerDB query engine
Current section
Files
Jump to
Current section
Files
src/dqe_mget.erl
-module(dqe_mget).
-behaviour(dflow).
-export([init/1, describe/1, start/2, emit/3, done/2]).
-record(state, {
acc = gb_trees:empty(),
count,
term_for_child = dict:new()
}).
init([SubQs]) ->
SubQs1 = [{make_ref(), SubQ} || SubQ <- SubQs],
{ok, #state{count = length(SubQs1)}, SubQs1}.
describe(_) ->
"sum".
start({_Start, _Count}, State) ->
{ok, State}.
emit(Child, {Data, Resolution},
State = #state{term_for_child = TFC, count = Count, acc = Tree}) ->
TFC1 = dict:update_counter(Child, 1, TFC),
Term = dict:fetch(Child, TFC1),
Tree1 = add_to_tree(Term, Data, Tree),
case shrink_tree(Tree1, Count, <<>>) of
{Tree2, <<>>} ->
{ok, State#state{acc = Tree2, term_for_child = TFC}};
{Tree2, Data1} ->
{emit, {Data1, Resolution},
State#state{acc = Tree2, term_for_child = TFC}}
end.
done({last, _Child}, State) ->
{done, State};
done(_, State) ->
{ok, State}.
add_to_tree(Term, Data, Tree) ->
case gb_trees:lookup(Term, Tree) of
none ->
gb_trees:insert(Term, {Data, 1}, Tree);
{value, {Sum, Count}} ->
Sum1 = mmath_comb:sum_r([Sum, Data]),
gb_trees:update(Term, {Sum1, Count+1}, Tree)
end.
shrink_tree(Tree, Count, Acc) ->
case gb_trees:is_empty(Tree) of
true ->
{Tree, Acc};
_ ->
case gb_trees:smallest(Tree) of
{Term0, {Data, FirstCount}} when FirstCount =:= Count ->
Tree1 = gb_trees:delete(Term0, Tree),
shrink_tree(Tree1, Count, <<Acc/binary, Data/binary>>);
_ ->
{Tree, Acc}
end
end.