ClickHouse/dbms/src/DataStreams/SquashingTransform.cpp

100 lines
2.6 KiB
C++
Raw Normal View History

#include <DataStreams/SquashingTransform.h>
namespace DB
{
SquashingTransform::SquashingTransform(size_t min_block_size_rows_, size_t min_block_size_bytes_, bool reserve_memory_)
: min_block_size_rows(min_block_size_rows_)
, min_block_size_bytes(min_block_size_bytes_)
, reserve_memory(reserve_memory_)
{
}
2018-09-08 19:23:48 +00:00
SquashingTransform::Result SquashingTransform::add(MutableColumns && columns)
{
2018-09-09 02:23:24 +00:00
/// End of input stream.
2018-09-08 19:23:48 +00:00
if (columns.empty())
return Result(std::move(accumulated_columns));
/// Just read block is already enough.
2018-09-08 19:23:48 +00:00
if (isEnoughSize(columns))
{
/// If no accumulated data, return just read block.
2018-09-08 19:23:48 +00:00
if (accumulated_columns.empty())
return Result(std::move(columns));
2018-09-09 02:23:24 +00:00
/// Return accumulated data (maybe it has small size) and place new block to accumulated data.
2018-09-08 19:23:48 +00:00
columns.swap(accumulated_columns);
return Result(std::move(columns));
}
/// Accumulated block is already enough.
2018-09-08 19:23:48 +00:00
if (!accumulated_columns.empty() && isEnoughSize(accumulated_columns))
{
/// Return accumulated data and place new block to accumulated data.
2018-09-08 19:23:48 +00:00
columns.swap(accumulated_columns);
return Result(std::move(columns));
}
2018-09-08 19:23:48 +00:00
append(std::move(columns));
2018-09-08 19:23:48 +00:00
if (isEnoughSize(accumulated_columns))
{
2018-09-08 19:23:48 +00:00
MutableColumns res;
res.swap(accumulated_columns);
return Result(std::move(res));
}
/// Squashed block is not ready.
return false;
}
2018-09-08 19:23:48 +00:00
void SquashingTransform::append(MutableColumns && columns)
{
2018-09-08 19:23:48 +00:00
if (accumulated_columns.empty())
{
2018-09-08 19:23:48 +00:00
accumulated_columns = std::move(columns);
return;
}
2018-09-08 19:23:48 +00:00
for (size_t i = 0, size = columns.size(); i < size; ++i)
{
auto & column = accumulated_columns[i];
if (reserve_memory)
column->reserve(min_block_size_bytes);
column->insertRangeFrom(*columns[i], 0, columns[i]->size());
}
2018-09-08 19:23:48 +00:00
}
bool SquashingTransform::isEnoughSize(const MutableColumns & columns)
{
size_t rows = 0;
size_t bytes = 0;
2018-09-08 19:23:48 +00:00
for (const auto & column : columns)
{
2018-09-08 19:23:48 +00:00
if (!rows)
rows = column->size();
else if (rows != column->size())
throw Exception("Sizes of columns doesn't match", ErrorCodes::SIZES_OF_COLUMNS_DOESNT_MATCH);
bytes += column->byteSize();
}
2018-09-08 19:23:48 +00:00
return isEnoughSize(rows, bytes);
}
bool SquashingTransform::isEnoughSize(size_t rows, size_t bytes) const
{
return (!min_block_size_rows && !min_block_size_bytes)
|| (min_block_size_rows && rows >= min_block_size_rows)
|| (min_block_size_bytes && bytes >= min_block_size_bytes);
}
}