2018-05-23 20:19:33 +00:00
|
|
|
#include <Processors/IAccumulatingTransform.h>
|
|
|
|
|
|
|
|
|
|
|
|
namespace DB
|
|
|
|
{
|
2020-02-25 18:10:48 +00:00
|
|
|
namespace ErrorCodes
|
|
|
|
{
|
|
|
|
extern const int LOGICAL_ERROR;
|
|
|
|
}
|
2018-05-23 20:19:33 +00:00
|
|
|
|
|
|
|
IAccumulatingTransform::IAccumulatingTransform(Block input_header, Block output_header)
|
|
|
|
: IProcessor({std::move(input_header)}, {std::move(output_header)}),
|
|
|
|
input(inputs.front()), output(outputs.front())
|
|
|
|
{
|
|
|
|
}
|
|
|
|
|
2020-09-04 08:36:47 +00:00
|
|
|
InputPort * IAccumulatingTransform::addTotalsPort()
|
|
|
|
{
|
|
|
|
if (inputs.size() > 1)
|
|
|
|
throw Exception("Totals port was already added to IAccumulatingTransform", ErrorCodes::LOGICAL_ERROR);
|
|
|
|
|
|
|
|
return &inputs.emplace_back(getInputPort().getHeader(), this);
|
|
|
|
}
|
|
|
|
|
2018-05-23 20:19:33 +00:00
|
|
|
IAccumulatingTransform::Status IAccumulatingTransform::prepare()
|
|
|
|
{
|
2019-02-07 18:51:53 +00:00
|
|
|
/// Check can output.
|
|
|
|
if (output.isFinished())
|
|
|
|
{
|
2020-09-04 20:44:15 +00:00
|
|
|
for (auto & in : inputs)
|
|
|
|
in.close();
|
|
|
|
|
2019-02-07 18:51:53 +00:00
|
|
|
return Status::Finished;
|
|
|
|
}
|
2018-05-23 20:19:33 +00:00
|
|
|
|
2019-02-07 18:51:53 +00:00
|
|
|
if (!output.canPush())
|
2018-05-23 20:19:33 +00:00
|
|
|
{
|
2019-02-07 18:51:53 +00:00
|
|
|
input.setNotNeeded();
|
|
|
|
return Status::PortFull;
|
2018-05-23 20:19:33 +00:00
|
|
|
}
|
|
|
|
|
2019-02-07 18:51:53 +00:00
|
|
|
/// Output if has data.
|
2019-02-18 16:36:07 +00:00
|
|
|
if (current_output_chunk)
|
|
|
|
output.push(std::move(current_output_chunk));
|
2019-02-07 18:51:53 +00:00
|
|
|
|
|
|
|
if (finished_generate)
|
2018-05-23 20:19:33 +00:00
|
|
|
{
|
2019-02-07 18:51:53 +00:00
|
|
|
output.finish();
|
|
|
|
return Status::Finished;
|
2018-05-23 20:19:33 +00:00
|
|
|
}
|
|
|
|
|
2020-09-04 17:14:36 +00:00
|
|
|
/// Close input if flag was set manually.
|
|
|
|
if (finished_input)
|
|
|
|
input.close();
|
|
|
|
|
|
|
|
/// Read from totals port if has it.
|
2019-03-04 14:56:09 +00:00
|
|
|
if (input.isFinished())
|
|
|
|
{
|
2020-09-04 08:36:47 +00:00
|
|
|
if (inputs.size() > 1)
|
|
|
|
{
|
|
|
|
auto & totals_input = inputs.back();
|
|
|
|
if (!totals_input.isFinished())
|
|
|
|
{
|
|
|
|
totals_input.setNeeded();
|
|
|
|
if (!totals_input.hasData())
|
|
|
|
return Status::NeedData;
|
|
|
|
|
|
|
|
totals = totals_input.pull();
|
|
|
|
totals_input.close();
|
|
|
|
}
|
|
|
|
}
|
2019-03-04 14:56:09 +00:00
|
|
|
}
|
|
|
|
|
2020-09-04 17:14:36 +00:00
|
|
|
/// Generate output block.
|
|
|
|
if (input.isFinished())
|
2019-03-04 14:56:09 +00:00
|
|
|
{
|
2020-09-04 17:14:36 +00:00
|
|
|
finished_input = true;
|
2019-03-04 14:56:09 +00:00
|
|
|
return Status::Ready;
|
|
|
|
}
|
|
|
|
|
2019-02-07 18:51:53 +00:00
|
|
|
/// Check can input.
|
|
|
|
if (!has_input)
|
2018-05-23 20:19:33 +00:00
|
|
|
{
|
2019-02-07 18:51:53 +00:00
|
|
|
input.setNeeded();
|
|
|
|
if (!input.hasData())
|
|
|
|
return Status::NeedData;
|
|
|
|
|
2019-02-18 16:36:07 +00:00
|
|
|
current_input_chunk = input.pull();
|
2019-02-07 18:51:53 +00:00
|
|
|
has_input = true;
|
2018-05-23 20:19:33 +00:00
|
|
|
}
|
|
|
|
|
2019-02-07 18:51:53 +00:00
|
|
|
return Status::Ready;
|
2018-05-23 20:19:33 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
void IAccumulatingTransform::work()
|
|
|
|
{
|
2019-02-07 18:51:53 +00:00
|
|
|
if (!finished_input)
|
2018-05-23 20:19:33 +00:00
|
|
|
{
|
2019-02-18 16:36:07 +00:00
|
|
|
consume(std::move(current_input_chunk));
|
2019-02-07 18:51:53 +00:00
|
|
|
has_input = false;
|
2018-05-23 20:19:33 +00:00
|
|
|
}
|
|
|
|
else
|
|
|
|
{
|
2019-02-18 16:36:07 +00:00
|
|
|
current_output_chunk = generate();
|
|
|
|
if (!current_output_chunk)
|
2019-02-07 18:51:53 +00:00
|
|
|
finished_generate = true;
|
2018-05-23 20:19:33 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2019-02-26 18:40:08 +00:00
|
|
|
void IAccumulatingTransform::setReadyChunk(Chunk chunk)
|
|
|
|
{
|
|
|
|
if (current_output_chunk)
|
|
|
|
throw Exception("IAccumulatingTransform already has input. Cannot set another chunk. "
|
2019-05-13 10:03:47 +00:00
|
|
|
"Probably, setReadyChunk method was called twice per consume().", ErrorCodes::LOGICAL_ERROR);
|
2019-02-26 18:40:08 +00:00
|
|
|
|
|
|
|
current_output_chunk = std::move(chunk);
|
|
|
|
}
|
|
|
|
|
2018-05-23 20:19:33 +00:00
|
|
|
}
|
|
|
|
|