2019-12-09 13:12:54 +00:00
|
|
|
#include <Functions/IFunctionImpl.h>
|
2018-09-08 22:04:39 +00:00
|
|
|
#include <Functions/FunctionFactory.h>
|
|
|
|
#include <Functions/FunctionHelpers.h>
|
|
|
|
#include <Columns/ColumnAggregateFunction.h>
|
|
|
|
#include <DataTypes/DataTypeAggregateFunction.h>
|
|
|
|
#include <Common/AlignedBuffer.h>
|
|
|
|
#include <Common/Arena.h>
|
|
|
|
#include <ext/scope_guard.h>
|
|
|
|
|
|
|
|
|
|
|
|
namespace DB
|
|
|
|
{
|
|
|
|
|
|
|
|
namespace ErrorCodes
|
|
|
|
{
|
|
|
|
extern const int ILLEGAL_COLUMN;
|
|
|
|
extern const int ILLEGAL_TYPE_OF_ARGUMENT;
|
2019-12-27 20:55:42 +00:00
|
|
|
extern const int NUMBER_OF_ARGUMENTS_DOESNT_MATCH;
|
2018-09-08 22:04:39 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
/** runningAccumulate(agg_state) - takes the states of the aggregate function and returns a column with values,
|
|
|
|
* are the result of the accumulation of these states for a set of block lines, from the first to the current line.
|
|
|
|
*
|
|
|
|
* Quite unusual function.
|
|
|
|
* Takes state of aggregate function (example runningAccumulate(uniqState(UserID))),
|
|
|
|
* and for each row of block, return result of aggregate function on merge of states of all previous rows and current row.
|
|
|
|
*
|
|
|
|
* So, result of function depends on partition of data to blocks and on order of data in block.
|
|
|
|
*/
|
|
|
|
class FunctionRunningAccumulate : public IFunction
|
|
|
|
{
|
|
|
|
public:
|
|
|
|
static constexpr auto name = "runningAccumulate";
|
|
|
|
static FunctionPtr create(const Context &)
|
|
|
|
{
|
|
|
|
return std::make_shared<FunctionRunningAccumulate>();
|
|
|
|
}
|
|
|
|
|
|
|
|
String getName() const override
|
|
|
|
{
|
|
|
|
return name;
|
|
|
|
}
|
|
|
|
|
2019-01-30 02:47:26 +00:00
|
|
|
bool isStateful() const override
|
|
|
|
{
|
|
|
|
return true;
|
|
|
|
}
|
|
|
|
|
2019-12-20 20:56:39 +00:00
|
|
|
bool isVariadic() const override { return true; }
|
|
|
|
|
|
|
|
size_t getNumberOfArguments() const override { return 0; }
|
2018-09-08 22:04:39 +00:00
|
|
|
|
|
|
|
bool isDeterministic() const override { return false; }
|
|
|
|
|
|
|
|
bool isDeterministicInScopeOfQuery() const override
|
|
|
|
{
|
|
|
|
return false;
|
|
|
|
}
|
|
|
|
|
|
|
|
DataTypePtr getReturnTypeImpl(const DataTypes & arguments) const override
|
|
|
|
{
|
2020-03-09 03:38:43 +00:00
|
|
|
if (arguments.empty() || arguments.size() > 2)
|
2019-12-28 20:54:50 +00:00
|
|
|
throw Exception("Incorrect number of arguments of function " + getName() + ". Must be 1 or 2.",
|
|
|
|
ErrorCodes::NUMBER_OF_ARGUMENTS_DOESNT_MATCH);
|
|
|
|
|
2018-09-08 22:04:39 +00:00
|
|
|
const DataTypeAggregateFunction * type = checkAndGetDataType<DataTypeAggregateFunction>(arguments[0].get());
|
|
|
|
if (!type)
|
|
|
|
throw Exception("Argument for function " + getName() + " must have type AggregateFunction - state of aggregate function.",
|
|
|
|
ErrorCodes::ILLEGAL_TYPE_OF_ARGUMENT);
|
|
|
|
|
|
|
|
return type->getReturnType();
|
|
|
|
}
|
|
|
|
|
2020-07-21 13:58:07 +00:00
|
|
|
void executeImpl(Block & block, const ColumnNumbers & arguments, size_t result, size_t /*input_rows_count*/) const override
|
2018-09-08 22:04:39 +00:00
|
|
|
{
|
|
|
|
const ColumnAggregateFunction * column_with_states
|
|
|
|
= typeid_cast<const ColumnAggregateFunction *>(&*block.getByPosition(arguments.at(0)).column);
|
2019-12-20 20:56:39 +00:00
|
|
|
|
2018-09-08 22:04:39 +00:00
|
|
|
if (!column_with_states)
|
|
|
|
throw Exception("Illegal column " + block.getByPosition(arguments.at(0)).column->getName()
|
|
|
|
+ " of first argument of function "
|
|
|
|
+ getName(),
|
|
|
|
ErrorCodes::ILLEGAL_COLUMN);
|
|
|
|
|
2019-12-20 20:56:39 +00:00
|
|
|
ColumnPtr column_with_groups;
|
|
|
|
|
2019-12-28 20:54:50 +00:00
|
|
|
if (arguments.size() == 2)
|
2019-12-20 20:56:39 +00:00
|
|
|
column_with_groups = block.getByPosition(arguments[1]).column;
|
|
|
|
|
2018-09-08 22:04:39 +00:00
|
|
|
AggregateFunctionPtr aggregate_function_ptr = column_with_states->getAggregateFunction();
|
|
|
|
const IAggregateFunction & agg_func = *aggregate_function_ptr;
|
|
|
|
|
|
|
|
AlignedBuffer place(agg_func.sizeOfData(), agg_func.alignOfData());
|
|
|
|
|
2019-12-20 20:56:39 +00:00
|
|
|
/// Will pass empty arena if agg_func does not allocate memory in arena
|
2018-09-08 22:04:39 +00:00
|
|
|
std::unique_ptr<Arena> arena = agg_func.allocatesMemoryInArena() ? std::make_unique<Arena>() : nullptr;
|
|
|
|
|
|
|
|
auto result_column_ptr = agg_func.getReturnType()->createColumn();
|
|
|
|
IColumn & result_column = *result_column_ptr;
|
|
|
|
result_column.reserve(column_with_states->size());
|
|
|
|
|
|
|
|
const auto & states = column_with_states->getData();
|
2019-12-20 20:56:39 +00:00
|
|
|
|
2019-12-28 20:54:50 +00:00
|
|
|
bool state_created = false;
|
2019-12-27 22:52:03 +00:00
|
|
|
SCOPE_EXIT({
|
2019-12-28 20:54:50 +00:00
|
|
|
if (state_created)
|
2019-12-27 22:52:03 +00:00
|
|
|
agg_func.destroy(place.data());
|
|
|
|
});
|
|
|
|
|
2019-12-28 20:54:50 +00:00
|
|
|
size_t row_number = 0;
|
2018-09-08 22:04:39 +00:00
|
|
|
for (const auto & state_to_add : states)
|
|
|
|
{
|
2019-12-28 20:54:50 +00:00
|
|
|
if (row_number == 0 || (column_with_groups && column_with_groups->compareAt(row_number, row_number - 1, *column_with_groups, 1) != 0))
|
2019-12-20 20:56:39 +00:00
|
|
|
{
|
2019-12-28 20:54:50 +00:00
|
|
|
if (state_created)
|
2019-12-27 22:52:03 +00:00
|
|
|
{
|
2019-12-28 20:54:50 +00:00
|
|
|
agg_func.destroy(place.data());
|
|
|
|
state_created = false;
|
2019-12-27 22:52:03 +00:00
|
|
|
}
|
|
|
|
|
2019-12-28 20:54:50 +00:00
|
|
|
agg_func.create(place.data());
|
|
|
|
state_created = true;
|
2019-12-20 20:56:39 +00:00
|
|
|
}
|
|
|
|
|
2018-09-08 22:04:39 +00:00
|
|
|
agg_func.merge(place.data(), state_to_add, arena.get());
|
2020-06-17 19:36:27 +00:00
|
|
|
agg_func.insertResultInto(place.data(), result_column, arena.get());
|
2019-12-20 20:56:39 +00:00
|
|
|
|
2019-12-28 20:54:50 +00:00
|
|
|
++row_number;
|
2018-09-08 22:04:39 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
block.getByPosition(result).column = std::move(result_column_ptr);
|
|
|
|
}
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
|
|
void registerFunctionRunningAccumulate(FunctionFactory & factory)
|
|
|
|
{
|
|
|
|
factory.registerFunction<FunctionRunningAccumulate>();
|
|
|
|
}
|
|
|
|
|
|
|
|
}
|