Packages
dqe
0.3.11
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_fun_flow2.erl
-module(dqe_fun_flow2).
-behaviour(dflow).
-export([init/1, describe/1, start/2, emit/3, done/2]).
-record(state, {
dqe_fun :: atom(),
chunk :: pos_integer(),
acc = <<>> :: binary(),
fun_state :: dqe_fun:fun_state()
}).
init([Fun, FunState, SubQ]) ->
Chunk = Fun:chunk(FunState),
{ok, #state{dqe_fun = Fun,
chunk = Chunk,
fun_state = FunState}, SubQ}.
start(_, State) ->
{ok, State}.
describe(#state{dqe_fun = Fun, fun_state = State}) ->
Fun:describe(State).
emit(_Child, Data, State = #state{dqe_fun = Fun, fun_state = FunState,
chunk = 1})
when is_list(Data) ->
{Result, FunState1} = Fun:run([Data], FunState),
{emit, Result, State#state{fun_state = FunState1}};
emit(_Child, Data, State = #state{dqe_fun = Fun, fun_state = FunState,
chunk = ChunkSize, acc = Acc})
when byte_size(Acc) + byte_size(Data) >= ChunkSize ->
Size = ((byte_size(Acc) + byte_size(Data)) div ChunkSize) * ChunkSize,
<<ToCompute:Size/binary, Acc1/binary>> = <<Acc/binary, Data/binary>>,
{Result, FunState1} = Fun:run([ToCompute], FunState),
{emit, Result, State#state{fun_state = FunState1, acc = Acc1}};
emit(_Child, Data, State = #state{acc = Acc}) ->
{ok, State#state{acc = <<Acc/binary, Data/binary>>}}.
done(_Child, State = #state{acc = <<>>}) ->
{done, State};
done(_Child,
State = #state{dqe_fun = Fun, fun_state = FunState, acc = Acc}) ->
{Result, FunState1} = Fun:run([Acc], FunState),
{done, Result, State#state{acc = <<>>, fun_state = FunState1}}.