ClickHouse/dbms/include/DB/Columns/ColumnAggregateFunction.h

303 lines
8.5 KiB
C++
Raw Normal View History

2011-09-19 01:42:16 +00:00
#pragma once
#include <DB/Common/Arena.h>
2011-09-19 01:42:16 +00:00
#include <DB/AggregateFunctions/IAggregateFunction.h>
2014-06-05 23:52:28 +00:00
#include <DB/Columns/IColumn.h>
2011-09-19 01:42:16 +00:00
#include <DB/Core/Field.h>
#include <DB/IO/ReadBufferFromString.h>
#include <DB/IO/WriteBuffer.h>
#include <DB/IO/WriteHelpers.h>
2011-09-19 01:42:16 +00:00
namespace DB
{
namespace ErrorCodes
{
extern const int PARAMETER_OUT_OF_BOUND;
extern const int SIZES_OF_COLUMNS_DOESNT_MATCH;
}
2017-03-09 00:56:38 +00:00
/** State column of aggregate functions.
* Presented as an array of pointers to the states of aggregate functions (data).
* The states themselves are stored in one of the pools (arenas).
2014-06-08 22:17:44 +00:00
*
2017-03-09 00:56:38 +00:00
* It can be in two variants:
2014-06-08 22:17:44 +00:00
*
2017-03-09 00:56:38 +00:00
* 1. Own its values - that is, be responsible for destroying them.
* The column consists of the values "assigned to it" after the aggregation is performed (see Aggregator, convertToBlocks function),
* or from values created by itself (see `insert` method).
2017-03-09 00:56:38 +00:00
* In this case, `src` will be `nullptr`, and the column itself will be destroyed (call `IAggregateFunction::destroy`)
* states of aggregate functions in the destructor.
*
* 2. Do not own its values, but use values taken from another ColumnAggregateFunction column.
2017-03-09 00:56:38 +00:00
* For example, this is a column obtained by permutation/filtering or other transformations from another column.
* In this case, `src` will be `shared ptr` to the source column. Destruction of values will be handled by this source column.
2014-06-08 22:17:44 +00:00
*
2017-03-09 00:56:38 +00:00
* This solution is somewhat limited:
* - the variant in which the column contains a part of "it's own" and a part of "another's" values is not supported;
2017-03-09 00:56:38 +00:00
* - the option of having multiple source columns is not supported, which may be necessary for a more optimal merge of the two columns.
2014-06-08 22:17:44 +00:00
*
2017-03-09 00:56:38 +00:00
* These restrictions can be removed if you add an array of flags or even refcount,
* specifying which individual values should be destroyed and which ones should not.
2017-03-09 00:56:38 +00:00
* Clearly, this method would have a substantially non-zero price.
*/
class ColumnAggregateFunction final : public IColumn, public std::enable_shared_from_this<ColumnAggregateFunction>
{
public:
using Container_t = PaddedPODArray<AggregateDataPtr>;
2014-06-08 22:17:44 +00:00
private:
/// Memory pools. Aggregate states are allocated from them.
Arenas arenas;
2014-06-05 23:52:28 +00:00
/// Used for destroying states and for finalization of values.
AggregateFunctionPtr func;
/// Source column. Used (holds source from destruction),
/// if this column has been constructed from another and uses all or part of its values.
std::shared_ptr<const ColumnAggregateFunction> src;
2014-06-05 23:52:28 +00:00
/// Array of pointers to aggregation states, that are placed in arenas.
Container_t data;
2014-06-05 23:52:28 +00:00
public:
/// Create a new column that has another column as a source.
ColumnAggregateFunction(const ColumnAggregateFunction & other)
: arenas(other.arenas), func(other.func), src(other.shared_from_this())
2014-06-08 22:17:44 +00:00
{
}
2014-06-05 23:52:28 +00:00
2014-06-08 22:17:44 +00:00
ColumnAggregateFunction(const AggregateFunctionPtr & func_)
: func(func_)
2014-06-06 19:35:41 +00:00
{
}
2014-06-08 22:17:44 +00:00
ColumnAggregateFunction(const AggregateFunctionPtr & func_, const Arenas & arenas_)
: arenas(arenas_), func(func_)
{
}
2017-03-09 04:26:17 +00:00
~ColumnAggregateFunction()
{
2016-06-09 04:54:30 +00:00
if (!func->hasTrivialDestructor() && !src)
for (auto val : data)
func->destroy(val);
}
2017-03-09 04:26:17 +00:00
void set(const AggregateFunctionPtr & func_)
{
func = func_;
}
AggregateFunctionPtr getAggregateFunction() { return func; }
AggregateFunctionPtr getAggregateFunction() const { return func; }
2014-02-27 12:49:21 +00:00
/// Take shared ownership of Arena, that holds memory for states of aggregate functions.
void addArena(ArenaPtr arena_)
{
2014-06-08 22:17:44 +00:00
arenas.push_back(arena_);
}
/** Transform column with states of aggregate functions to column with final result values.
*/
ColumnPtr convertToValues() const;
2014-02-27 12:49:21 +00:00
std::string getName() const override { return "ColumnAggregateFunction"; }
size_t sizeOfField() const override { return sizeof(getData()[0]); }
2014-06-05 19:52:13 +00:00
size_t size() const override
2014-06-05 23:52:28 +00:00
{
2014-06-08 22:17:44 +00:00
return getData().size();
}
2011-09-19 01:42:16 +00:00
ColumnPtr cloneEmpty() const override
{
return std::make_shared<ColumnAggregateFunction>(func, Arenas(1, std::make_shared<Arena>()));
2014-06-05 23:52:28 +00:00
};
2011-09-19 01:42:16 +00:00
Field operator[](size_t n) const override
2011-09-19 01:42:16 +00:00
{
Field field = String();
{
WriteBufferFromString buffer(field.get<String &>());
func->serialize(getData()[n], buffer);
}
return field;
2011-09-19 01:42:16 +00:00
}
void get(size_t n, Field & res) const override
{
2014-06-05 23:52:28 +00:00
res = String();
{
WriteBufferFromString buffer(res.get<String &>());
func->serialize(getData()[n], buffer);
}
}
StringRef getDataAt(size_t n) const override
{
2014-06-08 22:17:44 +00:00
return StringRef(reinterpret_cast<const char *>(&getData()[n]), sizeof(getData()[n]));
}
void insertData(const char * pos, size_t length) override
{
2014-06-08 22:17:44 +00:00
getData().push_back(*reinterpret_cast<const AggregateDataPtr *>(pos));
2014-05-26 16:11:20 +00:00
}
void insertFrom(const IColumn & src, size_t n) override
2014-05-26 16:11:20 +00:00
{
/// Must create new state of aggregate function and take ownership of it,
/// because ownership of states of aggregate function cannot be shared for individual rows,
/// (only as a whole, see comment above).
insertDefault();
insertMergeFrom(src, n);
2014-06-08 22:17:44 +00:00
}
2014-06-05 23:52:28 +00:00
void insertFrom(ConstAggregateDataPtr place)
{
insertDefault();
insertMergeFrom(place);
}
/// Merge state at last row with specified state in another column.
void insertMergeFrom(ConstAggregateDataPtr place)
{
func->merge(getData().back(), place, &createOrGetArena());
}
2014-06-08 22:17:44 +00:00
void insertMergeFrom(const IColumn & src, size_t n)
{
insertMergeFrom(static_cast<const ColumnAggregateFunction &>(src).getData()[n]);
2014-06-05 23:52:28 +00:00
}
Arena & createOrGetArena()
{
if (unlikely(arenas.empty()))
arenas.emplace_back(std::make_shared<Arena>());
return *arenas.back().get();
}
void insert(const Field & x) override
2014-06-05 23:52:28 +00:00
{
IAggregateFunction * function = func.get();
2014-06-08 22:17:44 +00:00
Arena & arena = createOrGetArena();
2015-11-10 20:39:11 +00:00
getData().push_back(arena.alloc(function->sizeOfData()));
2014-06-08 22:17:44 +00:00
function->create(getData().back());
ReadBufferFromString read_buffer(x.get<const String &>());
2016-09-22 23:26:08 +00:00
function->deserialize(getData().back(), read_buffer, &arena);
2014-06-05 23:52:28 +00:00
}
void insertDefault() override
2014-06-05 23:52:28 +00:00
{
IAggregateFunction * function = func.get();
2016-04-14 05:03:33 +00:00
Arena & arena = createOrGetArena();
getData().push_back(arena.alloc(function->sizeOfData()));
function->create(getData().back());
2014-06-05 23:52:28 +00:00
}
StringRef serializeValueIntoArena(size_t n, Arena & arena, char const *& begin) const override
{
throw Exception("Method serializeValueIntoArena is not supported for " + getName(), ErrorCodes::NOT_IMPLEMENTED);
}
const char * deserializeAndInsertFromArena(const char * pos) override
{
throw Exception("Method deserializeAndInsertFromArena is not supported for " + getName(), ErrorCodes::NOT_IMPLEMENTED);
}
void updateHashWithValue(size_t n, SipHash & hash) const override;
size_t byteSize() const override;
2014-05-28 14:54:42 +00:00
size_t allocatedSize() const override;
2017-01-20 05:00:04 +00:00
void insertRangeFrom(const IColumn & from, size_t start, size_t length) override;
2013-05-03 05:23:14 +00:00
void popBack(size_t n) override
{
size_t size = data.size();
size_t new_size = size - n;
if (!src)
for (size_t i = new_size; i < size; ++i)
func->destroy(data[i]);
data.resize_assume_reserved(new_size);
}
2017-01-20 05:00:04 +00:00
ColumnPtr filter(const Filter & filter, ssize_t result_size_hint) const override;
2017-01-20 05:00:04 +00:00
ColumnPtr permute(const Permutation & perm, size_t limit) const override;
2013-05-03 05:23:14 +00:00
ColumnPtr replicate(const Offsets_t & offsets) const override
{
throw Exception("Method replicate is not supported for ColumnAggregateFunction.", ErrorCodes::NOT_IMPLEMENTED);
}
Columns scatter(ColumnIndex num_columns, const Selector & selector) const override
{
/// Columns with scattered values will point to this column as the owner of values.
Columns columns(num_columns);
for (auto & column : columns)
column = std::make_shared<ColumnAggregateFunction>(*this);
size_t num_rows = size();
{
size_t reserve_size = num_rows / num_columns * 1.1; /// 1.1 is just a guess. Better to use n-sigma rule.
if (reserve_size > 1)
for (auto & column : columns)
column->reserve(reserve_size);
}
for (size_t i = 0; i < num_rows; ++i)
static_cast<ColumnAggregateFunction &>(*columns[selector[i]]).data.push_back(data[i]);
return columns;
}
int compareAt(size_t n, size_t m, const IColumn & rhs_, int nan_direction_hint) const override
2011-09-19 01:42:16 +00:00
{
return 0;
}
2011-09-26 11:05:38 +00:00
void getPermutation(bool reverse, size_t limit, Permutation & res) const override
2011-09-26 11:05:38 +00:00
{
2014-06-08 22:17:44 +00:00
size_t s = getData().size();
res.resize(s);
2011-09-26 11:05:38 +00:00
for (size_t i = 0; i < s; ++i)
res[i] = i;
}
2014-06-05 23:52:28 +00:00
2017-03-09 04:26:17 +00:00
/** More efficient manipulation methods */
2014-06-05 23:52:28 +00:00
Container_t & getData()
{
return data;
2014-06-05 23:52:28 +00:00
}
const Container_t & getData() const
{
return data;
2014-06-05 23:52:28 +00:00
}
void getExtremes(Field & min, Field & max) const override
{
throw Exception("Method getExtremes is not supported for ColumnAggregateFunction.", ErrorCodes::NOT_IMPLEMENTED);
}
2011-09-19 01:42:16 +00:00
};
}