Packages
dqe
0.1.35
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_aggr2.erl
-module(dqe_aggr2).
-behaviour(dflow).
-include_lib("mmath/include/mmath.hrl").
-export([init/1, describe/1, start/2, emit/3, done/2]).
-record(state, {
aggr :: atom(),
time :: pos_integer(),
acc = <<>> :: binary(),
resolution :: pos_integer()
}).
init([Aggr, SubQ, Time]) ->
{ok, #state{aggr = Aggr, time = Time}, SubQ}.
start({_Start, _Count}, State) ->
{ok, State}.
describe(#state{aggr = Aggr, time = Time}) ->
[atom_to_list(Aggr), "(", integer_to_list(Time), "ms)"].
%% When we get the first data we can calculate both the applied
%% time and the upwards resolution.
emit(Child, {realized, {Data, Resolution}},
State = #state{resolution = undefined, time = Time}) ->
Time1 = dqe_time:apply_times(Time, Resolution),
emit(Child, {Data, Resolution},
State#state{resolution = Time1 * Resolution, time = Time1});
emit(_Child, {realized, {Data, _R}},
State = #state{aggr = Aggr, time = Time, acc = Acc}) ->
case execute(Aggr, <<Data/binary, Acc/binary>>, Time, <<>>) of
{Acc1, <<>>} ->
{ok, State#state{acc = Acc1}};
{Acc2, AccEmit} ->
{emit, {realized, {AccEmit, State#state.resolution}},
State#state{acc = Acc2}}
end.
done(_Child, State = #state{acc = <<>>}) ->
{done, State};
done(_Child, State = #state{aggr = Aggr, time = Time, acc = Acc}) ->
Data = mmath_aggr:Aggr(Acc, Time),
{done, {realized, {Data, State#state.resolution}}, State#state{acc = <<>>}}.
execute(Aggr, Acc, T1, AccEmit) when byte_size(Acc) >= T1 * 9 ->
MinSize = T1 * ?DATA_SIZE,
<<Data:MinSize/binary, Acc1/binary>> = Acc,
Result = mmath_aggr:Aggr(Data, T1),
execute(Aggr, Acc1, T1, <<AccEmit/binary, Result/binary>>);
execute(_, Acc, _, AccEmit) ->
{Acc, AccEmit}.