2019-10-10 16:30:30 +00:00
|
|
|
#include <Storages/MergeTree/MergeTreeReaderWide.h>
|
2020-02-25 09:49:45 +00:00
|
|
|
|
2017-04-01 09:19:00 +00:00
|
|
|
#include <Columns/ColumnArray.h>
|
2020-02-25 09:49:45 +00:00
|
|
|
#include <DataTypes/DataTypeArray.h>
|
|
|
|
#include <DataTypes/NestedUtils.h>
|
2020-02-25 08:53:14 +00:00
|
|
|
#include <Interpreters/inplaceBlockConversions.h>
|
2020-02-25 09:49:45 +00:00
|
|
|
#include <Storages/MergeTree/IMergeTreeReader.h>
|
2020-02-10 20:27:06 +00:00
|
|
|
#include <Storages/MergeTree/MergeTreeDataPartWide.h>
|
2020-02-25 09:49:45 +00:00
|
|
|
#include <Common/escapeForFileName.h>
|
2017-07-13 20:58:19 +00:00
|
|
|
#include <Common/typeid_cast.h>
|
2016-07-19 10:57:57 +00:00
|
|
|
|
|
|
|
namespace DB
|
|
|
|
{
|
|
|
|
|
|
|
|
namespace
|
|
|
|
{
|
2017-04-01 07:20:54 +00:00
|
|
|
using OffsetColumns = std::map<std::string, ColumnPtr>;
|
|
|
|
constexpr auto DATA_FILE_EXTENSION = ".bin";
|
2016-07-19 10:57:57 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
namespace ErrorCodes
|
|
|
|
{
|
2017-04-01 07:20:54 +00:00
|
|
|
extern const int MEMORY_LIMIT_EXCEEDED;
|
2016-11-20 12:43:20 +00:00
|
|
|
}
|
2016-07-21 16:22:24 +00:00
|
|
|
|
2020-02-10 20:27:06 +00:00
|
|
|
MergeTreeReaderWide::MergeTreeReaderWide(
|
2020-03-23 02:12:31 +00:00
|
|
|
DataPartWidePtr data_part_,
|
|
|
|
NamesAndTypesList columns_,
|
2020-06-17 16:39:58 +00:00
|
|
|
const StorageMetadataPtr & metadata_snapshot_,
|
2020-02-10 20:27:06 +00:00
|
|
|
UncompressedCache * uncompressed_cache_,
|
|
|
|
MarkCache * mark_cache_,
|
2020-03-23 02:12:31 +00:00
|
|
|
MarkRanges mark_ranges_,
|
|
|
|
MergeTreeReaderSettings settings_,
|
|
|
|
IMergeTreeDataPart::ValueSizeMap avg_value_size_hints_,
|
2020-02-10 20:27:06 +00:00
|
|
|
const ReadBufferFromFileBase::ProfileCallback & profile_callback_,
|
|
|
|
clockid_t clock_type_)
|
2020-02-25 09:49:45 +00:00
|
|
|
: IMergeTreeReader(
|
2020-06-17 16:39:58 +00:00
|
|
|
std::move(data_part_),
|
|
|
|
std::move(columns_),
|
|
|
|
metadata_snapshot_,
|
|
|
|
uncompressed_cache_,
|
|
|
|
std::move(mark_cache_),
|
|
|
|
std::move(mark_ranges_),
|
|
|
|
std::move(settings_),
|
|
|
|
std::move(avg_value_size_hints_))
|
2016-07-19 10:57:57 +00:00
|
|
|
{
|
2017-04-01 07:20:54 +00:00
|
|
|
try
|
|
|
|
{
|
2017-12-25 21:57:29 +00:00
|
|
|
for (const NameAndTypePair & column : columns)
|
2020-02-25 08:53:14 +00:00
|
|
|
{
|
2020-04-08 16:20:52 +00:00
|
|
|
auto column_from_part = getColumnFromPart(column);
|
2020-09-14 11:22:17 +00:00
|
|
|
addStreams(column_from_part, profile_callback_, clock_type_);
|
2020-02-25 08:53:14 +00:00
|
|
|
}
|
2017-04-01 07:20:54 +00:00
|
|
|
}
|
|
|
|
catch (...)
|
|
|
|
{
|
|
|
|
storage.reportBrokenPart(data_part->name);
|
|
|
|
throw;
|
|
|
|
}
|
2016-07-19 10:57:57 +00:00
|
|
|
}
|
|
|
|
|
2019-11-13 01:57:45 +00:00
|
|
|
|
2019-12-19 13:10:57 +00:00
|
|
|
size_t MergeTreeReaderWide::readRows(size_t from_mark, bool continue_reading, size_t max_rows_to_read, Columns & res_columns)
|
2016-07-19 10:57:57 +00:00
|
|
|
{
|
2017-06-16 20:11:02 +00:00
|
|
|
size_t read_rows = 0;
|
2017-04-01 07:20:54 +00:00
|
|
|
try
|
|
|
|
{
|
2019-09-23 19:22:02 +00:00
|
|
|
size_t num_columns = columns.size();
|
2020-04-14 19:47:19 +00:00
|
|
|
checkNumberOfColumns(num_columns);
|
2019-09-23 19:22:02 +00:00
|
|
|
|
2017-04-01 07:20:54 +00:00
|
|
|
/// Pointers to offset columns that are common to the nested data structure columns.
|
|
|
|
/// If append is true, then the value will be equal to nullptr and will be used only to
|
|
|
|
/// check that the offsets column has been already read.
|
|
|
|
OffsetColumns offset_columns;
|
2021-03-09 14:46:52 +00:00
|
|
|
std::unordered_map<String, ISerialization::SubstreamsCache> caches;
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-07-26 00:34:36 +00:00
|
|
|
/// Request reading of data in advance,
|
|
|
|
/// so if reading can be asynchronous, it will also be performed in parallel for all columns.
|
2019-09-23 19:22:02 +00:00
|
|
|
auto name_and_type = columns.begin();
|
|
|
|
for (size_t pos = 0; pos < num_columns; ++pos, ++name_and_type)
|
2021-07-26 00:34:36 +00:00
|
|
|
{
|
|
|
|
auto column_from_part = getColumnFromPart(*name_and_type);
|
|
|
|
try
|
|
|
|
{
|
|
|
|
auto & cache = caches[column_from_part.getNameInStorage()];
|
|
|
|
prefetch(column_from_part, from_mark, continue_reading, cache);
|
|
|
|
}
|
|
|
|
catch (Exception & e)
|
|
|
|
{
|
|
|
|
/// Better diagnostics.
|
|
|
|
e.addMessage("(while reading column " + column_from_part.name + ")");
|
|
|
|
throw;
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
name_and_type = columns.begin();
|
|
|
|
for (size_t pos = 0; pos < num_columns; ++pos, ++name_and_type)
|
2017-04-01 07:20:54 +00:00
|
|
|
{
|
2020-09-14 11:22:17 +00:00
|
|
|
auto column_from_part = getColumnFromPart(*name_and_type);
|
|
|
|
const auto & [name, type] = column_from_part;
|
2019-09-23 19:22:02 +00:00
|
|
|
|
2017-04-01 07:20:54 +00:00
|
|
|
/// The column is already present in the block so we will append the values to the end.
|
2019-09-23 19:22:02 +00:00
|
|
|
bool append = res_columns[pos] != nullptr;
|
2017-12-15 21:11:24 +00:00
|
|
|
if (!append)
|
2020-02-25 08:53:14 +00:00
|
|
|
res_columns[pos] = type->createColumn();
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2020-11-10 17:32:00 +00:00
|
|
|
auto & column = res_columns[pos];
|
2017-04-01 07:20:54 +00:00
|
|
|
try
|
|
|
|
{
|
2017-12-15 21:11:24 +00:00
|
|
|
size_t column_size_before_reading = column->size();
|
2020-12-22 15:03:48 +00:00
|
|
|
auto & cache = caches[column_from_part.getNameInStorage()];
|
2017-11-21 02:23:41 +00:00
|
|
|
|
2020-11-10 17:32:00 +00:00
|
|
|
readData(column_from_part, column, from_mark, continue_reading, max_rows_to_read, cache);
|
2017-11-21 02:23:41 +00:00
|
|
|
|
|
|
|
/// For elements of Nested, column_size_before_reading may be greater than column size
|
|
|
|
/// if offsets are not empty and were already read, but elements are empty.
|
2019-09-23 19:22:02 +00:00
|
|
|
if (!column->empty())
|
2017-12-15 21:11:24 +00:00
|
|
|
read_rows = std::max(read_rows, column->size() - column_size_before_reading);
|
2017-04-01 07:20:54 +00:00
|
|
|
}
|
|
|
|
catch (Exception & e)
|
|
|
|
{
|
|
|
|
/// Better diagnostics.
|
2019-09-23 19:22:02 +00:00
|
|
|
e.addMessage("(while reading column " + name + ")");
|
2017-04-01 07:20:54 +00:00
|
|
|
throw;
|
|
|
|
}
|
|
|
|
|
2019-09-23 19:22:02 +00:00
|
|
|
if (column->empty())
|
|
|
|
res_columns[pos] = nullptr;
|
2017-04-01 07:20:54 +00:00
|
|
|
}
|
|
|
|
|
2019-09-23 19:22:02 +00:00
|
|
|
/// NOTE: positions for all streams must be kept in sync.
|
|
|
|
/// In particular, even if for some streams there are no rows to be read,
|
2017-04-01 07:20:54 +00:00
|
|
|
/// you must ensure that no seeks are skipped and at this point they all point to to_mark.
|
|
|
|
}
|
|
|
|
catch (Exception & e)
|
|
|
|
{
|
|
|
|
if (e.code() != ErrorCodes::MEMORY_LIMIT_EXCEEDED)
|
|
|
|
storage.reportBrokenPart(data_part->name);
|
|
|
|
|
|
|
|
/// Better diagnostics.
|
2020-02-27 16:47:40 +00:00
|
|
|
e.addMessage("(while reading from part " + data_part->getFullPath() + " "
|
2019-09-23 19:22:02 +00:00
|
|
|
"from mark " + toString(from_mark) + " "
|
|
|
|
"with max_rows_to_read = " + toString(max_rows_to_read) + ")");
|
2017-04-01 07:20:54 +00:00
|
|
|
throw;
|
|
|
|
}
|
|
|
|
catch (...)
|
|
|
|
{
|
|
|
|
storage.reportBrokenPart(data_part->name);
|
|
|
|
|
|
|
|
throw;
|
|
|
|
}
|
2017-06-16 20:11:02 +00:00
|
|
|
|
|
|
|
return read_rows;
|
2016-07-19 10:57:57 +00:00
|
|
|
}
|
|
|
|
|
2020-09-14 11:22:17 +00:00
|
|
|
void MergeTreeReaderWide::addStreams(const NameAndTypePair & name_and_type,
|
2017-08-07 07:31:16 +00:00
|
|
|
const ReadBufferFromFileBase::ProfileCallback & profile_callback, clockid_t clock_type)
|
2016-07-19 10:57:57 +00:00
|
|
|
{
|
2021-03-09 14:46:52 +00:00
|
|
|
ISerialization::StreamCallback callback = [&] (const ISerialization::SubstreamPath & substream_path)
|
2017-04-01 07:20:54 +00:00
|
|
|
{
|
2021-03-09 14:46:52 +00:00
|
|
|
String stream_name = ISerialization::getFileNameForStream(name_and_type, substream_path);
|
2017-11-20 02:15:15 +00:00
|
|
|
|
2017-08-07 07:31:16 +00:00
|
|
|
if (streams.count(stream_name))
|
|
|
|
return;
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2019-06-16 20:13:54 +00:00
|
|
|
bool data_file_exists = data_part->checksums.files.count(stream_name + DATA_FILE_EXTENSION);
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2017-08-07 07:31:16 +00:00
|
|
|
/** If data file is missing then we will not try to open it.
|
|
|
|
* It is necessary since it allows to add new column to structure of the table without creating new files for old parts.
|
|
|
|
*/
|
2017-11-21 02:23:41 +00:00
|
|
|
if (!data_file_exists)
|
2017-04-01 07:20:54 +00:00
|
|
|
return;
|
2019-12-18 16:41:11 +00:00
|
|
|
|
2019-02-05 14:50:25 +00:00
|
|
|
streams.emplace(stream_name, std::make_unique<MergeTreeReaderStream>(
|
2020-05-09 21:24:15 +00:00
|
|
|
data_part->volume->getDisk(), data_part->getFullRelativePath() + stream_name, DATA_FILE_EXTENSION,
|
2020-02-27 16:47:40 +00:00
|
|
|
data_part->getMarksCount(), all_mark_ranges, settings, mark_cache,
|
2019-02-27 20:02:48 +00:00
|
|
|
uncompressed_cache, data_part->getFileSizeOrZero(stream_name + DATA_FILE_EXTENSION),
|
2019-06-19 10:07:56 +00:00
|
|
|
&data_part->index_granularity_info,
|
2018-11-15 14:06:54 +00:00
|
|
|
profile_callback, clock_type));
|
2017-08-07 07:31:16 +00:00
|
|
|
};
|
|
|
|
|
2021-03-09 14:46:52 +00:00
|
|
|
auto serialization = data_part->getSerializationForColumn(name_and_type);
|
|
|
|
serialization->enumerateStreams(callback);
|
|
|
|
serializations.emplace(name_and_type.name, std::move(serialization));
|
2016-07-19 10:57:57 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
|
2021-07-26 00:34:36 +00:00
|
|
|
static ReadBuffer * getStream(
|
|
|
|
bool stream_for_prefix,
|
|
|
|
const ISerialization::SubstreamPath & substream_path,
|
|
|
|
MergeTreeReaderWide::FileStreams & streams,
|
|
|
|
const NameAndTypePair & name_and_type,
|
|
|
|
size_t from_mark, bool continue_reading,
|
2021-03-09 14:46:52 +00:00
|
|
|
ISerialization::SubstreamsCache & cache)
|
2016-07-19 10:57:57 +00:00
|
|
|
{
|
2021-07-26 00:34:36 +00:00
|
|
|
/// If substream have already been read.
|
|
|
|
if (cache.count(ISerialization::getSubcolumnNameForStream(substream_path)))
|
|
|
|
return nullptr;
|
|
|
|
|
|
|
|
String stream_name = ISerialization::getFileNameForStream(name_and_type, substream_path);
|
|
|
|
|
|
|
|
auto it = streams.find(stream_name);
|
|
|
|
if (it == streams.end())
|
|
|
|
return nullptr;
|
|
|
|
|
|
|
|
MergeTreeReaderStream & stream = *it->second;
|
|
|
|
|
|
|
|
if (stream_for_prefix)
|
|
|
|
stream.seekToStart();
|
|
|
|
else if (!continue_reading)
|
|
|
|
stream.seekToMark(from_mark);
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-07-26 00:34:36 +00:00
|
|
|
return stream.data_buffer;
|
|
|
|
}
|
2017-04-01 07:20:54 +00:00
|
|
|
|
|
|
|
|
2021-07-26 00:34:36 +00:00
|
|
|
void MergeTreeReaderWide::prefetch(
|
|
|
|
const NameAndTypePair & name_and_type,
|
|
|
|
size_t from_mark,
|
|
|
|
bool continue_reading,
|
|
|
|
ISerialization::SubstreamsCache & cache)
|
|
|
|
{
|
|
|
|
const auto & name = name_and_type.name;
|
|
|
|
auto & serialization = serializations[name];
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-07-26 00:34:36 +00:00
|
|
|
serialization->enumerateStreams([&](const ISerialization::SubstreamPath & substream_path)
|
|
|
|
{
|
|
|
|
if (ReadBuffer * buf = getStream(false, substream_path, streams, name_and_type, from_mark, continue_reading, cache))
|
|
|
|
buf->prefetch();
|
|
|
|
});
|
|
|
|
}
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2017-08-07 07:31:16 +00:00
|
|
|
|
2021-07-26 00:34:36 +00:00
|
|
|
void MergeTreeReaderWide::readData(
|
|
|
|
const NameAndTypePair & name_and_type, ColumnPtr & column,
|
|
|
|
size_t from_mark, bool continue_reading, size_t max_rows_to_read,
|
|
|
|
ISerialization::SubstreamsCache & cache)
|
|
|
|
{
|
2020-09-14 11:22:17 +00:00
|
|
|
double & avg_value_size_hint = avg_value_size_hints[name_and_type.name];
|
2021-03-09 14:46:52 +00:00
|
|
|
ISerialization::DeserializeBinaryBulkSettings deserialize_settings;
|
2019-12-12 18:55:19 +00:00
|
|
|
deserialize_settings.avg_value_size_hint = avg_value_size_hint;
|
2018-05-21 16:21:15 +00:00
|
|
|
|
2021-03-09 14:46:52 +00:00
|
|
|
const auto & name = name_and_type.name;
|
2021-07-26 00:34:36 +00:00
|
|
|
auto & serialization = serializations[name];
|
2021-03-09 14:46:52 +00:00
|
|
|
|
|
|
|
if (deserialize_binary_bulk_state_map.count(name) == 0)
|
2018-06-07 18:14:37 +00:00
|
|
|
{
|
2021-07-26 00:34:36 +00:00
|
|
|
deserialize_settings.getter = [&](const ISerialization::SubstreamPath & substream_path)
|
|
|
|
{
|
|
|
|
return getStream(true, substream_path, streams, name_and_type, from_mark, continue_reading, cache);
|
|
|
|
};
|
2021-03-09 14:46:52 +00:00
|
|
|
serialization->deserializeBinaryBulkStatePrefix(deserialize_settings, deserialize_binary_bulk_state_map[name]);
|
2018-06-07 18:14:37 +00:00
|
|
|
}
|
2018-05-21 16:21:15 +00:00
|
|
|
|
2021-07-26 00:34:36 +00:00
|
|
|
deserialize_settings.getter = [&](const ISerialization::SubstreamPath & substream_path)
|
|
|
|
{
|
|
|
|
return getStream(false, substream_path, streams, name_and_type, from_mark, continue_reading, cache);
|
|
|
|
};
|
2019-12-12 18:55:19 +00:00
|
|
|
deserialize_settings.continuous_reading = continue_reading;
|
2021-03-09 14:46:52 +00:00
|
|
|
auto & deserialize_state = deserialize_binary_bulk_state_map[name];
|
|
|
|
|
2021-07-26 00:34:36 +00:00
|
|
|
serialization->deserializeBinaryBulkWithMultipleStreams(column, max_rows_to_read, deserialize_settings, deserialize_state, &cache);
|
2020-11-10 17:32:00 +00:00
|
|
|
IDataType::updateAvgValueSizeHint(*column, avg_value_size_hint);
|
2016-07-19 10:57:57 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
}
|