2017-04-01 09:19:00 +00:00
|
|
|
#include <Storages/StorageLog.h>
|
2017-12-30 00:36:06 +00:00
|
|
|
#include <Storages/StorageFactory.h>
|
2011-11-05 23:31:19 +00:00
|
|
|
|
2017-04-01 09:19:00 +00:00
|
|
|
#include <Common/Exception.h>
|
2018-01-15 19:07:47 +00:00
|
|
|
#include <Common/StringUtils/StringUtils.h>
|
2019-12-23 16:57:16 +00:00
|
|
|
#include <Common/typeid_cast.h>
|
2010-03-18 19:32:14 +00:00
|
|
|
|
2020-02-14 14:28:33 +00:00
|
|
|
#include <Interpreters/evaluateConstantExpression.h>
|
|
|
|
|
2023-10-23 12:13:36 +00:00
|
|
|
#include <Parsers/ASTCheckQuery.h>
|
|
|
|
|
2021-11-01 00:39:38 +00:00
|
|
|
#include <IO/LimitReadBuffer.h>
|
2020-02-14 14:28:33 +00:00
|
|
|
#include <IO/ReadBufferFromFileBase.h>
|
2021-11-01 00:39:38 +00:00
|
|
|
#include <IO/ReadHelpers.h>
|
2020-02-20 16:39:32 +00:00
|
|
|
#include <IO/WriteBufferFromFileBase.h>
|
2021-11-01 00:39:38 +00:00
|
|
|
#include <IO/WriteHelpers.h>
|
2021-10-26 09:48:31 +00:00
|
|
|
#include <IO/copyData.h>
|
2020-10-26 17:24:15 +00:00
|
|
|
#include <Compression/CompressedReadBuffer.h>
|
2018-12-28 18:15:26 +00:00
|
|
|
#include <Compression/CompressedWriteBuffer.h>
|
2012-01-09 19:20:48 +00:00
|
|
|
|
2017-12-25 18:58:39 +00:00
|
|
|
#include <DataTypes/NestedUtils.h>
|
2012-08-29 20:07:24 +00:00
|
|
|
|
2017-05-24 21:06:29 +00:00
|
|
|
#include <Interpreters/Context.h>
|
2020-02-18 14:41:30 +00:00
|
|
|
#include "StorageLogSettings.h"
|
2021-10-10 20:38:21 +00:00
|
|
|
#include <Processors/Sources/NullSource.h>
|
2022-05-20 19:49:31 +00:00
|
|
|
#include <Processors/ISource.h>
|
2021-10-16 14:03:50 +00:00
|
|
|
#include <QueryPipeline/Pipe.h>
|
2021-07-23 19:33:59 +00:00
|
|
|
#include <Processors/Sinks/SinkToStorage.h>
|
2017-01-21 04:24:28 +00:00
|
|
|
|
2022-05-29 19:53:56 +00:00
|
|
|
#include <Backups/BackupEntriesCollector.h>
|
2022-01-31 06:35:07 +00:00
|
|
|
#include <Backups/BackupEntryFromAppendOnlyFile.h>
|
2023-04-25 17:44:03 +00:00
|
|
|
#include <Backups/BackupEntryFromMemory.h>
|
2021-10-26 09:48:31 +00:00
|
|
|
#include <Backups/BackupEntryFromSmallFile.h>
|
2023-04-23 10:25:46 +00:00
|
|
|
#include <Backups/BackupEntryWrappedWith.h>
|
2021-10-26 09:48:31 +00:00
|
|
|
#include <Backups/IBackup.h>
|
2022-05-31 09:33:23 +00:00
|
|
|
#include <Backups/RestorerFromBackup.h>
|
2021-10-26 09:48:31 +00:00
|
|
|
#include <Disks/TemporaryFileOnDisk.h>
|
2023-09-20 09:31:12 +00:00
|
|
|
#include <Storages/BlockNumberColumn.h>
|
2021-10-26 09:48:31 +00:00
|
|
|
|
2020-09-18 19:25:56 +00:00
|
|
|
#include <cassert>
|
2021-07-12 10:58:53 +00:00
|
|
|
#include <chrono>
|
2020-09-18 19:25:56 +00:00
|
|
|
|
2010-03-18 19:32:14 +00:00
|
|
|
|
2017-12-30 00:36:06 +00:00
|
|
|
#define DBMS_STORAGE_LOG_DATA_FILE_EXTENSION ".bin"
|
|
|
|
#define DBMS_STORAGE_LOG_MARKS_FILE_NAME "__marks.mrk"
|
2012-01-09 19:20:48 +00:00
|
|
|
|
|
|
|
|
2010-03-18 19:32:14 +00:00
|
|
|
namespace DB
|
|
|
|
{
|
|
|
|
|
2023-09-20 09:31:12 +00:00
|
|
|
CompressionCodecPtr getCompressionCodecDelta(UInt8 delta_bytes_size);
|
|
|
|
|
2016-01-11 21:46:36 +00:00
|
|
|
namespace ErrorCodes
|
|
|
|
{
|
2020-09-24 23:29:16 +00:00
|
|
|
extern const int TIMEOUT_EXCEEDED;
|
2016-01-11 21:46:36 +00:00
|
|
|
extern const int LOGICAL_ERROR;
|
|
|
|
extern const int DUPLICATE_COLUMN;
|
|
|
|
extern const int SIZES_OF_MARKS_FILES_ARE_INCONSISTENT;
|
2017-12-30 00:36:06 +00:00
|
|
|
extern const int NUMBER_OF_ARGUMENTS_DOESNT_MATCH;
|
2017-11-03 19:53:10 +00:00
|
|
|
extern const int INCORRECT_FILE_NAME;
|
2022-06-29 12:42:23 +00:00
|
|
|
extern const int CANNOT_RESTORE_TABLE;
|
2023-10-23 12:13:36 +00:00
|
|
|
extern const int NOT_IMPLEMENTED;
|
2016-01-11 21:46:36 +00:00
|
|
|
}
|
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
/// NOTE: The lock `StorageLog::rwlock` is NOT kept locked while reading,
|
|
|
|
/// because we read ranges of data that do not change.
|
2022-05-20 19:49:31 +00:00
|
|
|
class LogSource final : public ISource
|
2010-03-18 19:32:14 +00:00
|
|
|
{
|
2015-01-18 08:25:56 +00:00
|
|
|
public:
|
2020-01-31 15:10:10 +00:00
|
|
|
static Block getHeader(const NamesAndTypesList & columns)
|
|
|
|
{
|
|
|
|
Block res;
|
|
|
|
|
|
|
|
for (const auto & name_type : columns)
|
|
|
|
res.insert({ name_type.type->createColumn(), name_type.type, name_type.name });
|
|
|
|
|
2020-11-10 17:32:00 +00:00
|
|
|
return res;
|
2020-01-31 15:10:10 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
LogSource(
|
2021-11-01 00:39:38 +00:00
|
|
|
size_t block_size_,
|
|
|
|
const NamesAndTypesList & columns_,
|
|
|
|
const StorageLog & storage_,
|
|
|
|
size_t rows_limit_,
|
|
|
|
const std::vector<size_t> & offsets_,
|
|
|
|
const std::vector<size_t> & file_sizes_,
|
|
|
|
bool limited_by_file_sizes_,
|
|
|
|
ReadSettings read_settings_)
|
2022-05-20 19:49:31 +00:00
|
|
|
: ISource(getHeader(columns_))
|
2021-11-01 00:39:38 +00:00
|
|
|
, block_size(block_size_)
|
|
|
|
, columns(columns_)
|
|
|
|
, storage(storage_)
|
|
|
|
, rows_limit(rows_limit_)
|
|
|
|
, offsets(offsets_)
|
|
|
|
, file_sizes(file_sizes_)
|
|
|
|
, limited_by_file_sizes(limited_by_file_sizes_)
|
|
|
|
, read_settings(std::move(read_settings_))
|
2016-07-12 18:08:16 +00:00
|
|
|
{
|
|
|
|
}
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2016-12-19 23:55:13 +00:00
|
|
|
String getName() const override { return "Log"; }
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2015-01-18 08:25:56 +00:00
|
|
|
protected:
|
2020-01-31 15:10:10 +00:00
|
|
|
Chunk generate() override;
|
2016-08-24 00:39:38 +00:00
|
|
|
|
2015-01-18 08:25:56 +00:00
|
|
|
private:
|
2024-01-15 16:13:32 +00:00
|
|
|
NameAndTypePair getColumnOnDisk(const NameAndTypePair & column) const;
|
|
|
|
|
2021-11-01 00:39:38 +00:00
|
|
|
const size_t block_size;
|
|
|
|
const NamesAndTypesList columns;
|
2021-08-26 22:15:24 +00:00
|
|
|
const StorageLog & storage;
|
2021-11-01 00:39:38 +00:00
|
|
|
const size_t rows_limit; /// The maximum number of rows that can be read
|
2015-01-18 08:25:56 +00:00
|
|
|
size_t rows_read = 0;
|
2021-08-26 22:15:24 +00:00
|
|
|
bool is_finished = false;
|
2021-11-01 00:39:38 +00:00
|
|
|
const std::vector<size_t> offsets;
|
|
|
|
const std::vector<size_t> file_sizes;
|
|
|
|
const bool limited_by_file_sizes;
|
|
|
|
const ReadSettings read_settings;
|
2021-03-09 14:46:52 +00:00
|
|
|
|
2015-01-18 08:25:56 +00:00
|
|
|
struct Stream
|
|
|
|
{
|
2021-11-01 00:39:38 +00:00
|
|
|
Stream(const DiskPtr & disk, const String & data_path, size_t offset, size_t file_size, bool limited_by_file_size, ReadSettings read_settings_)
|
2015-01-18 08:25:56 +00:00
|
|
|
{
|
2021-11-01 00:39:38 +00:00
|
|
|
plain = disk->readFile(data_path, read_settings_.adjustBufferSize(file_size));
|
|
|
|
|
2015-01-18 08:25:56 +00:00
|
|
|
if (offset)
|
2020-10-29 14:14:23 +00:00
|
|
|
plain->seek(offset, SEEK_SET);
|
2021-11-01 00:39:38 +00:00
|
|
|
|
|
|
|
if (limited_by_file_size)
|
|
|
|
{
|
2023-02-22 16:54:35 +00:00
|
|
|
limited.emplace(*plain, file_size - offset, /* trow_exception */ false, /* exact_limit */ std::optional<size_t>());
|
2021-11-01 00:39:38 +00:00
|
|
|
compressed.emplace(*limited);
|
|
|
|
}
|
|
|
|
else
|
|
|
|
compressed.emplace(*plain);
|
2015-01-18 08:25:56 +00:00
|
|
|
}
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2020-10-26 17:24:15 +00:00
|
|
|
std::unique_ptr<ReadBufferFromFileBase> plain;
|
2021-11-01 00:39:38 +00:00
|
|
|
std::optional<LimitReadBuffer> limited;
|
|
|
|
std::optional<CompressedReadBuffer> compressed;
|
2015-01-18 08:25:56 +00:00
|
|
|
};
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2019-12-12 08:57:25 +00:00
|
|
|
using FileStreams = std::map<String, Stream>;
|
2015-01-18 08:25:56 +00:00
|
|
|
FileStreams streams;
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-03-09 14:46:52 +00:00
|
|
|
using DeserializeState = ISerialization::DeserializeBinaryBulkStatePtr;
|
2018-06-07 18:14:37 +00:00
|
|
|
using DeserializeStates = std::map<String, DeserializeState>;
|
|
|
|
DeserializeStates deserialize_states;
|
|
|
|
|
2021-03-09 14:46:52 +00:00
|
|
|
void readData(const NameAndTypePair & name_and_type, ColumnPtr & column, size_t max_rows_to_read, ISerialization::SubstreamsCache & cache);
|
2021-08-26 22:15:24 +00:00
|
|
|
bool isFinished();
|
2015-01-18 08:25:56 +00:00
|
|
|
};
|
|
|
|
|
2024-01-15 16:13:32 +00:00
|
|
|
NameAndTypePair LogSource::getColumnOnDisk(const NameAndTypePair & column) const
|
|
|
|
{
|
|
|
|
const auto & storage_columns = storage.columns_with_collected_nested;
|
|
|
|
|
|
|
|
/// A special case when we read subcolumn of shared offsets of Nested.
|
|
|
|
/// E.g. instead of requested column "n.arr1.size0" we must read column "n.size0" from disk.
|
|
|
|
auto name_in_storage = column.getNameInStorage();
|
|
|
|
if (column.getSubcolumnName() == "size0" && Nested::isSubcolumnOfNested(name_in_storage, storage_columns))
|
|
|
|
{
|
|
|
|
auto nested_name_in_storage = Nested::splitName(name_in_storage).first;
|
|
|
|
auto new_name = Nested::concatenateName(nested_name_in_storage, column.getSubcolumnName());
|
|
|
|
return storage_columns.getColumnOrSubcolumn(GetColumnsOptions::All, new_name);
|
|
|
|
}
|
|
|
|
|
|
|
|
return column;
|
|
|
|
}
|
2015-01-18 08:25:56 +00:00
|
|
|
|
2020-01-31 15:10:10 +00:00
|
|
|
Chunk LogSource::generate()
|
2010-03-18 19:32:14 +00:00
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
if (isFinished())
|
|
|
|
{
|
|
|
|
/// Close the files (before destroying the object).
|
|
|
|
/// When many sources are created, but simultaneously reading only a few of them,
|
|
|
|
/// buffers don't waste memory.
|
|
|
|
streams.clear();
|
2020-01-31 15:10:10 +00:00
|
|
|
return {};
|
2021-08-26 22:15:24 +00:00
|
|
|
}
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2017-03-13 18:01:46 +00:00
|
|
|
/// How many rows to read for the next block.
|
2012-06-22 16:54:51 +00:00
|
|
|
size_t max_rows_to_read = std::min(block_size, rows_limit - rows_read);
|
2021-03-09 14:46:52 +00:00
|
|
|
std::unordered_map<String, ISerialization::SubstreamsCache> caches;
|
2021-08-26 22:15:24 +00:00
|
|
|
Block res;
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2018-01-02 06:13:22 +00:00
|
|
|
for (const auto & name_type : columns)
|
2010-03-18 19:32:14 +00:00
|
|
|
{
|
2020-11-10 17:32:00 +00:00
|
|
|
ColumnPtr column;
|
2024-01-15 16:13:32 +00:00
|
|
|
auto name_type_on_disk = getColumnOnDisk(name_type);
|
|
|
|
|
2020-11-10 17:32:00 +00:00
|
|
|
try
|
2015-07-01 20:14:23 +00:00
|
|
|
{
|
2024-01-15 16:13:32 +00:00
|
|
|
column = name_type_on_disk.type->createColumn();
|
|
|
|
readData(name_type_on_disk, column, max_rows_to_read, caches[name_type_on_disk.getNameInStorage()]);
|
2015-07-01 20:14:23 +00:00
|
|
|
}
|
2020-11-10 17:32:00 +00:00
|
|
|
catch (Exception & e)
|
2015-07-01 20:14:23 +00:00
|
|
|
{
|
2024-01-15 16:13:32 +00:00
|
|
|
e.addMessage("while reading column " + name_type_on_disk.name + " at " + fullPath(storage.disk, storage.table_path));
|
2020-11-10 17:32:00 +00:00
|
|
|
throw;
|
2015-07-01 20:14:23 +00:00
|
|
|
}
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2020-01-31 15:10:10 +00:00
|
|
|
if (!column->empty())
|
2024-01-15 16:13:32 +00:00
|
|
|
res.insert(ColumnWithTypeAndName(column, name_type_on_disk.type, name_type_on_disk.name));
|
2010-03-18 19:32:14 +00:00
|
|
|
}
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2012-03-05 02:34:20 +00:00
|
|
|
if (res)
|
2012-06-22 17:00:59 +00:00
|
|
|
rows_read += res.rows();
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
if (!res)
|
|
|
|
is_finished = true;
|
|
|
|
|
|
|
|
if (isFinished())
|
2012-06-21 16:16:58 +00:00
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
/// Close the files (before destroying the object).
|
|
|
|
/// When many sources are created, but simultaneously reading only a few of them,
|
|
|
|
/// buffers don't waste memory.
|
2017-03-13 18:01:46 +00:00
|
|
|
streams.clear();
|
2012-06-21 16:16:58 +00:00
|
|
|
}
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2020-01-31 15:10:10 +00:00
|
|
|
UInt64 num_rows = res.rows();
|
|
|
|
return Chunk(res.getColumns(), num_rows);
|
2010-03-18 19:32:14 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
|
2020-11-10 17:32:00 +00:00
|
|
|
void LogSource::readData(const NameAndTypePair & name_and_type, ColumnPtr & column,
|
2021-03-09 14:46:52 +00:00
|
|
|
size_t max_rows_to_read, ISerialization::SubstreamsCache & cache)
|
2012-08-29 20:07:24 +00:00
|
|
|
{
|
2021-03-09 14:46:52 +00:00
|
|
|
ISerialization::DeserializeBinaryBulkSettings settings; /// TODO Use avg_value_size_hint.
|
2020-09-14 11:22:17 +00:00
|
|
|
const auto & [name, type] = name_and_type;
|
2021-11-01 02:13:07 +00:00
|
|
|
auto serialization = IDataType::getSerialization(name_and_type);
|
2018-06-07 18:14:37 +00:00
|
|
|
|
2020-10-21 23:02:20 +00:00
|
|
|
auto create_stream_getter = [&](bool stream_for_prefix)
|
2012-08-29 20:07:24 +00:00
|
|
|
{
|
2023-02-19 22:15:09 +00:00
|
|
|
return [&, stream_for_prefix] (const ISerialization::SubstreamPath & path) -> ReadBuffer *
|
2018-06-07 18:14:37 +00:00
|
|
|
{
|
2022-04-18 10:18:43 +00:00
|
|
|
if (cache.contains(ISerialization::getSubcolumnNameForStream(path)))
|
2020-11-10 17:32:00 +00:00
|
|
|
return nullptr;
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
String data_file_name = ISerialization::getFileNameForStream(name_and_type, path);
|
2017-11-28 02:13:46 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
const auto & data_file_it = storage.data_files_by_names.find(data_file_name);
|
|
|
|
if (data_file_it == storage.data_files_by_names.end())
|
2024-02-19 01:58:51 +00:00
|
|
|
throw Exception(ErrorCodes::LOGICAL_ERROR, "No information about file {} in StorageLog", data_file_name);
|
2021-08-26 22:15:24 +00:00
|
|
|
const auto & data_file = *data_file_it->second;
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
size_t offset = stream_for_prefix ? 0 : offsets[data_file.index];
|
2021-11-01 00:39:38 +00:00
|
|
|
size_t file_size = file_sizes[data_file.index];
|
2020-10-21 23:02:20 +00:00
|
|
|
|
2021-11-01 00:39:38 +00:00
|
|
|
auto it = streams.try_emplace(data_file_name, storage.disk, data_file.path, offset, file_size, limited_by_file_sizes, read_settings).first;
|
|
|
|
return &it->second.compressed.value();
|
2018-06-07 18:14:37 +00:00
|
|
|
};
|
2017-08-07 07:31:16 +00:00
|
|
|
};
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2022-04-18 10:18:43 +00:00
|
|
|
if (!deserialize_states.contains(name))
|
2018-06-07 18:14:37 +00:00
|
|
|
{
|
2020-10-21 23:02:20 +00:00
|
|
|
settings.getter = create_stream_getter(true);
|
2021-03-09 14:46:52 +00:00
|
|
|
serialization->deserializeBinaryBulkStatePrefix(settings, deserialize_states[name]);
|
2018-06-07 18:14:37 +00:00
|
|
|
}
|
|
|
|
|
2020-10-21 23:02:20 +00:00
|
|
|
settings.getter = create_stream_getter(false);
|
2021-03-09 14:46:52 +00:00
|
|
|
serialization->deserializeBinaryBulkWithMultipleStreams(column, max_rows_to_read, settings, deserialize_states[name], &cache);
|
2012-08-29 20:07:24 +00:00
|
|
|
}
|
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
bool LogSource::isFinished()
|
|
|
|
{
|
|
|
|
if (is_finished)
|
|
|
|
return true;
|
2012-08-29 20:07:24 +00:00
|
|
|
|
2021-11-01 00:39:38 +00:00
|
|
|
/// Check for row limit.
|
|
|
|
if (rows_read == rows_limit)
|
2021-08-26 22:15:24 +00:00
|
|
|
{
|
2021-11-01 00:39:38 +00:00
|
|
|
is_finished = true;
|
|
|
|
return true;
|
2021-08-26 22:15:24 +00:00
|
|
|
}
|
2021-11-01 00:39:38 +00:00
|
|
|
|
|
|
|
if (limited_by_file_sizes)
|
2021-08-26 22:15:24 +00:00
|
|
|
{
|
2021-11-01 00:39:38 +00:00
|
|
|
/// Check for EOF.
|
|
|
|
if (!streams.empty() && streams.begin()->second.compressed->eof())
|
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
is_finished = true;
|
2021-11-01 00:39:38 +00:00
|
|
|
return true;
|
|
|
|
}
|
2021-08-26 22:15:24 +00:00
|
|
|
}
|
|
|
|
|
2021-11-01 00:39:38 +00:00
|
|
|
return false;
|
2021-08-26 22:15:24 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
/// NOTE: The lock `StorageLog::rwlock` is kept locked in exclusive mode while writing.
|
2021-07-23 19:33:59 +00:00
|
|
|
class LogSink final : public SinkToStorage
|
2014-03-19 10:45:13 +00:00
|
|
|
{
|
2015-01-18 08:25:56 +00:00
|
|
|
public:
|
2021-08-26 22:15:24 +00:00
|
|
|
using WriteLock = std::unique_lock<std::shared_timed_mutex>;
|
|
|
|
|
2021-07-23 19:33:59 +00:00
|
|
|
explicit LogSink(
|
2021-08-26 22:15:24 +00:00
|
|
|
StorageLog & storage_, const StorageMetadataPtr & metadata_snapshot_, WriteLock && lock_)
|
2021-07-26 10:08:40 +00:00
|
|
|
: SinkToStorage(metadata_snapshot_->getSampleBlock())
|
2021-07-23 19:33:59 +00:00
|
|
|
, storage(storage_)
|
2020-06-16 15:51:29 +00:00
|
|
|
, metadata_snapshot(metadata_snapshot_)
|
2020-09-24 23:29:16 +00:00
|
|
|
, lock(std::move(lock_))
|
2015-01-18 08:25:56 +00:00
|
|
|
{
|
2020-09-24 23:29:16 +00:00
|
|
|
if (!lock)
|
2023-01-23 21:13:58 +00:00
|
|
|
throw Exception(ErrorCodes::TIMEOUT_EXCEEDED, "Lock timeout exceeded");
|
2021-01-05 01:49:15 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
/// Ensure that marks are loaded because we're going to update them.
|
|
|
|
storage.loadMarks(lock);
|
|
|
|
|
|
|
|
/// If there were no files, save zero file sizes to be able to rollback in case of error.
|
|
|
|
storage.saveFileSizes(lock);
|
2015-01-18 08:25:56 +00:00
|
|
|
}
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-07-23 19:33:59 +00:00
|
|
|
String getName() const override { return "LogSink"; }
|
|
|
|
|
|
|
|
~LogSink() override
|
2015-04-02 23:58:26 +00:00
|
|
|
{
|
|
|
|
try
|
|
|
|
{
|
2020-07-12 02:31:58 +00:00
|
|
|
if (!done)
|
|
|
|
{
|
|
|
|
/// Rollback partial writes.
|
2021-08-26 22:15:24 +00:00
|
|
|
|
|
|
|
/// No more writing.
|
2020-07-12 02:31:58 +00:00
|
|
|
streams.clear();
|
2021-08-26 22:15:24 +00:00
|
|
|
|
|
|
|
/// Truncate files to the older sizes.
|
2020-07-12 02:31:58 +00:00
|
|
|
storage.file_checker.repair();
|
2021-08-26 22:15:24 +00:00
|
|
|
|
|
|
|
/// Remove excessive marks.
|
|
|
|
storage.removeUnsavedMarks(lock);
|
2020-07-12 02:31:58 +00:00
|
|
|
}
|
2015-04-02 23:58:26 +00:00
|
|
|
}
|
|
|
|
catch (...)
|
|
|
|
{
|
|
|
|
tryLogCurrentException(__PRETTY_FUNCTION__);
|
|
|
|
}
|
|
|
|
}
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-07-23 19:33:59 +00:00
|
|
|
void consume(Chunk chunk) override;
|
|
|
|
void onFinish() override;
|
2015-04-02 23:58:26 +00:00
|
|
|
|
2015-01-18 08:25:56 +00:00
|
|
|
private:
|
|
|
|
StorageLog & storage;
|
2020-06-16 15:51:29 +00:00
|
|
|
StorageMetadataPtr metadata_snapshot;
|
2021-08-26 22:15:24 +00:00
|
|
|
WriteLock lock;
|
2015-04-02 23:58:26 +00:00
|
|
|
bool done = false;
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2015-01-18 08:25:56 +00:00
|
|
|
struct Stream
|
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
Stream(const DiskPtr & disk, const String & data_path, size_t initial_data_size, CompressionCodecPtr codec, size_t max_compress_block_size) :
|
2020-01-22 16:17:25 +00:00
|
|
|
plain(disk->writeFile(data_path, max_compress_block_size, WriteMode::Append)),
|
|
|
|
compressed(*plain, std::move(codec), max_compress_block_size),
|
2021-08-26 22:15:24 +00:00
|
|
|
plain_offset(initial_data_size)
|
2015-01-18 08:25:56 +00:00
|
|
|
{
|
|
|
|
}
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2020-01-22 16:17:25 +00:00
|
|
|
std::unique_ptr<WriteBuffer> plain;
|
2015-01-18 08:25:56 +00:00
|
|
|
CompressedWriteBuffer compressed;
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
/// How many bytes were in the file at the time the Stream was created.
|
|
|
|
size_t plain_offset;
|
|
|
|
|
|
|
|
/// Used to not write shared offsets of columns for nested structures multiple times.
|
|
|
|
bool written = false;
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2015-01-18 08:25:56 +00:00
|
|
|
void finalize()
|
|
|
|
{
|
|
|
|
compressed.next();
|
2023-05-28 10:29:38 +00:00
|
|
|
compressed.finalize();
|
|
|
|
|
2020-01-22 16:17:25 +00:00
|
|
|
plain->next();
|
2023-05-28 10:29:38 +00:00
|
|
|
plain->finalize();
|
2015-01-18 08:25:56 +00:00
|
|
|
}
|
|
|
|
};
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2019-12-12 08:57:25 +00:00
|
|
|
using FileStreams = std::map<String, Stream>;
|
2015-01-18 08:25:56 +00:00
|
|
|
FileStreams streams;
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-03-09 14:46:52 +00:00
|
|
|
using SerializeState = ISerialization::SerializeBinaryBulkStatePtr;
|
2018-06-07 18:14:37 +00:00
|
|
|
using SerializeStates = std::map<String, SerializeState>;
|
|
|
|
SerializeStates serialize_states;
|
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
ISerialization::OutputStreamGetter createStreamGetter(const NameAndTypePair & name_and_type);
|
2017-08-07 07:31:16 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
void writeData(const NameAndTypePair & name_and_type, const IColumn & column);
|
2015-01-18 08:25:56 +00:00
|
|
|
};
|
2014-03-19 10:45:13 +00:00
|
|
|
|
|
|
|
|
2021-07-23 19:33:59 +00:00
|
|
|
void LogSink::consume(Chunk chunk)
|
2010-03-18 19:32:14 +00:00
|
|
|
{
|
2021-09-03 17:29:36 +00:00
|
|
|
auto block = getHeader().cloneWithColumns(chunk.detachColumns());
|
2020-06-17 14:32:25 +00:00
|
|
|
metadata_snapshot->check(block, true);
|
2014-06-26 00:58:14 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
for (auto & stream : streams | boost::adaptors::map_values)
|
|
|
|
stream.written = false;
|
2016-07-13 10:35:00 +00:00
|
|
|
|
2010-03-18 19:32:14 +00:00
|
|
|
for (size_t i = 0; i < block.columns(); ++i)
|
|
|
|
{
|
2017-01-02 20:12:12 +00:00
|
|
|
const ColumnWithTypeAndName & column = block.safeGetByPosition(i);
|
2021-08-26 22:15:24 +00:00
|
|
|
writeData(NameAndTypePair(column.name, column.type), *column.column);
|
2012-08-29 20:07:24 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
|
2021-07-23 19:33:59 +00:00
|
|
|
void LogSink::onFinish()
|
2013-09-15 01:40:29 +00:00
|
|
|
{
|
2015-04-02 23:58:26 +00:00
|
|
|
if (done)
|
|
|
|
return;
|
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
for (auto & stream : streams | boost::adaptors::map_values)
|
|
|
|
stream.written = false;
|
|
|
|
|
2021-03-09 14:46:52 +00:00
|
|
|
ISerialization::SerializeBinaryBulkSettings settings;
|
2021-09-03 17:29:36 +00:00
|
|
|
for (const auto & column : getHeader())
|
2018-06-07 18:14:37 +00:00
|
|
|
{
|
|
|
|
auto it = serialize_states.find(column.name);
|
|
|
|
if (it != serialize_states.end())
|
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
settings.getter = createStreamGetter(NameAndTypePair(column.name, column.type));
|
2021-03-09 14:46:52 +00:00
|
|
|
auto serialization = column.type->getDefaultSerialization();
|
|
|
|
serialization->serializeBinaryBulkStateSuffix(settings, it->second);
|
2018-06-07 18:14:37 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2017-03-13 18:01:46 +00:00
|
|
|
/// Finish write.
|
2021-08-26 22:15:24 +00:00
|
|
|
for (auto & stream : streams | boost::adaptors::map_values)
|
|
|
|
stream.finalize();
|
|
|
|
streams.clear();
|
2014-08-04 06:36:24 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
storage.saveMarks(lock);
|
|
|
|
storage.saveFileSizes(lock);
|
2022-05-03 17:55:45 +00:00
|
|
|
storage.updateTotalRows(lock);
|
2014-08-04 06:36:24 +00:00
|
|
|
|
2020-07-12 02:31:58 +00:00
|
|
|
done = true;
|
2021-04-04 05:23:12 +00:00
|
|
|
|
|
|
|
/// unlock should be done from the same thread as lock, and dtor may be
|
|
|
|
/// called from different thread, so it should be done here (at least in
|
|
|
|
/// case of no exceptions occurred)
|
|
|
|
lock.unlock();
|
2013-09-15 01:40:29 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
ISerialization::OutputStreamGetter LogSink::createStreamGetter(const NameAndTypePair & name_and_type)
|
2018-06-07 18:14:37 +00:00
|
|
|
{
|
2021-03-09 14:46:52 +00:00
|
|
|
return [&] (const ISerialization::SubstreamPath & path) -> WriteBuffer *
|
2018-06-07 18:14:37 +00:00
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
String data_file_name = ISerialization::getFileNameForStream(name_and_type, path);
|
|
|
|
auto it = streams.find(data_file_name);
|
|
|
|
if (it == streams.end())
|
2024-02-19 01:58:51 +00:00
|
|
|
throw Exception(ErrorCodes::LOGICAL_ERROR, "Stream was not created when writing data in LogSink");
|
2021-08-26 22:15:24 +00:00
|
|
|
|
|
|
|
Stream & stream = it->second;
|
|
|
|
if (stream.written)
|
2018-06-07 18:14:37 +00:00
|
|
|
return nullptr;
|
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
return &stream.compressed;
|
2018-06-07 18:14:37 +00:00
|
|
|
};
|
|
|
|
}
|
|
|
|
|
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
void LogSink::writeData(const NameAndTypePair & name_and_type, const IColumn & column)
|
2012-08-29 20:07:24 +00:00
|
|
|
{
|
2021-03-09 14:46:52 +00:00
|
|
|
ISerialization::SerializeBinaryBulkSettings settings;
|
2020-09-14 11:22:17 +00:00
|
|
|
const auto & [name, type] = name_and_type;
|
2021-03-09 14:46:52 +00:00
|
|
|
auto serialization = type->getDefaultSerialization();
|
2018-06-07 18:14:37 +00:00
|
|
|
|
2021-03-09 14:46:52 +00:00
|
|
|
serialization->enumerateStreams([&] (const ISerialization::SubstreamPath & path)
|
2016-07-13 10:35:00 +00:00
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
String data_file_name = ISerialization::getFileNameForStream(name_and_type, path);
|
|
|
|
auto it = streams.find(data_file_name);
|
|
|
|
if (it == streams.end())
|
|
|
|
{
|
|
|
|
const auto & data_file_it = storage.data_files_by_names.find(data_file_name);
|
|
|
|
if (data_file_it == storage.data_files_by_names.end())
|
2024-02-19 01:58:51 +00:00
|
|
|
throw Exception(ErrorCodes::LOGICAL_ERROR, "No information about file {} in StorageLog", data_file_name);
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
const auto & data_file = *data_file_it->second;
|
|
|
|
const auto & columns = metadata_snapshot->getColumns();
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2023-09-20 09:31:12 +00:00
|
|
|
CompressionCodecPtr compression;
|
|
|
|
if (name_and_type.name == BlockNumberColumn::name)
|
|
|
|
compression = BlockNumberColumn::compression_codec;
|
|
|
|
else
|
|
|
|
compression = columns.getCodecOrDefault(name_and_type.name);
|
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
it = streams.try_emplace(data_file.name, storage.disk, data_file.path,
|
|
|
|
storage.file_checker.getFileSize(data_file.path),
|
2023-09-20 09:31:12 +00:00
|
|
|
compression, storage.max_compress_block_size).first;
|
2021-08-26 22:15:24 +00:00
|
|
|
}
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
auto & stream = it->second;
|
|
|
|
if (stream.written)
|
2018-06-07 18:14:37 +00:00
|
|
|
return;
|
2021-10-11 22:01:00 +00:00
|
|
|
});
|
2017-11-26 19:22:33 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
settings.getter = createStreamGetter(name_and_type);
|
2018-06-07 18:14:37 +00:00
|
|
|
|
2022-04-18 10:18:43 +00:00
|
|
|
if (!serialize_states.contains(name))
|
2022-09-08 15:16:39 +00:00
|
|
|
serialization->serializeBinaryBulkStatePrefix(column, settings, serialize_states[name]);
|
2018-06-07 18:14:37 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
if (storage.use_marks_file)
|
2017-11-26 19:22:33 +00:00
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
serialization->enumerateStreams([&] (const ISerialization::SubstreamPath & path)
|
|
|
|
{
|
|
|
|
String data_file_name = ISerialization::getFileNameForStream(name_and_type, path);
|
|
|
|
const auto & stream = streams.at(data_file_name);
|
|
|
|
if (stream.written)
|
|
|
|
return;
|
|
|
|
|
|
|
|
auto & data_file = *storage.data_files_by_names.at(data_file_name);
|
|
|
|
auto & marks = data_file.marks;
|
|
|
|
size_t prev_num_rows = marks.empty() ? 0 : marks.back().rows;
|
|
|
|
auto & mark = marks.emplace_back();
|
|
|
|
mark.rows = prev_num_rows + column.size();
|
|
|
|
mark.offset = stream.plain_offset + stream.plain->count();
|
|
|
|
});
|
|
|
|
}
|
2017-11-26 19:22:33 +00:00
|
|
|
|
2021-03-09 14:46:52 +00:00
|
|
|
serialization->serializeBinaryBulkWithMultipleStreams(column, 0, 0, settings, serialize_states[name]);
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-03-09 14:46:52 +00:00
|
|
|
serialization->enumerateStreams([&] (const ISerialization::SubstreamPath & path)
|
2012-08-29 20:07:24 +00:00
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
String data_file_name = ISerialization::getFileNameForStream(name_and_type, path);
|
|
|
|
auto & stream = streams.at(data_file_name);
|
|
|
|
if (stream.written)
|
2017-11-28 02:13:46 +00:00
|
|
|
return;
|
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
stream.written = true;
|
|
|
|
stream.compressed.next();
|
2021-10-11 22:01:00 +00:00
|
|
|
});
|
2010-03-18 19:32:14 +00:00
|
|
|
}
|
|
|
|
|
2013-02-26 13:06:01 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
void StorageLog::Mark::write(WriteBuffer & out) const
|
2016-07-12 18:08:16 +00:00
|
|
|
{
|
2023-09-15 15:19:52 +00:00
|
|
|
writeBinaryLittleEndian(rows, out);
|
|
|
|
writeBinaryLittleEndian(offset, out);
|
2021-08-26 22:15:24 +00:00
|
|
|
}
|
2014-06-26 00:58:14 +00:00
|
|
|
|
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
void StorageLog::Mark::read(ReadBuffer & in)
|
|
|
|
{
|
2023-09-15 15:19:52 +00:00
|
|
|
readBinaryLittleEndian(rows, in);
|
|
|
|
readBinaryLittleEndian(offset, in);
|
2021-08-26 22:15:24 +00:00
|
|
|
}
|
2014-06-26 00:58:14 +00:00
|
|
|
|
2021-10-10 06:13:00 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
namespace
|
|
|
|
{
|
|
|
|
/// NOTE: We extract the number of rows from the marks.
|
|
|
|
/// For normal columns, the number of rows in the block is specified in the marks.
|
|
|
|
/// For array columns and nested structures, there are more than one group of marks that correspond to different files
|
|
|
|
/// - for elements (file name.bin) - the total number of array elements in the block is specified,
|
|
|
|
/// - for array sizes (file name.size0.bin) - the number of rows (the whole arrays themselves) in the block is specified.
|
|
|
|
/// So for Array data type, first stream is array sizes; and number of array sizes is the number of arrays.
|
|
|
|
/// Thus we assume we can always get the real number of rows from the first column.
|
|
|
|
constexpr size_t INDEX_WITH_REAL_ROW_COUNT = 0;
|
2013-02-26 13:06:01 +00:00
|
|
|
}
|
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
|
2021-09-22 19:31:12 +00:00
|
|
|
StorageLog::~StorageLog() = default;
|
2021-09-20 09:05:34 +00:00
|
|
|
|
2014-09-30 03:08:47 +00:00
|
|
|
StorageLog::StorageLog(
|
2021-08-26 22:15:24 +00:00
|
|
|
const String & engine_name_,
|
2019-12-12 08:57:25 +00:00
|
|
|
DiskPtr disk_,
|
2019-12-26 14:03:32 +00:00
|
|
|
const String & relative_path_,
|
2019-12-04 16:06:55 +00:00
|
|
|
const StorageID & table_id_,
|
2018-03-06 20:18:34 +00:00
|
|
|
const ColumnsDescription & columns_,
|
2019-08-24 21:20:20 +00:00
|
|
|
const ConstraintsDescription & constraints_,
|
2021-04-23 12:18:23 +00:00
|
|
|
const String & comment,
|
2020-07-12 02:31:58 +00:00
|
|
|
bool attach,
|
2022-07-13 20:35:24 +00:00
|
|
|
ContextMutablePtr context_)
|
2019-12-04 16:06:55 +00:00
|
|
|
: IStorage(table_id_)
|
2022-07-13 20:35:24 +00:00
|
|
|
, WithMutableContext(context_)
|
2021-08-26 22:15:24 +00:00
|
|
|
, engine_name(engine_name_)
|
2020-01-13 11:41:42 +00:00
|
|
|
, disk(std::move(disk_))
|
|
|
|
, table_path(relative_path_)
|
2021-08-26 22:15:24 +00:00
|
|
|
, use_marks_file(engine_name == "Log")
|
|
|
|
, marks_file_path(table_path + DBMS_STORAGE_LOG_MARKS_FILE_NAME)
|
2020-01-13 11:41:42 +00:00
|
|
|
, file_checker(disk, table_path + "sizes.json")
|
2022-07-13 20:35:24 +00:00
|
|
|
, max_compress_block_size(context_->getSettingsRef().max_compress_block_size)
|
2010-03-18 19:32:14 +00:00
|
|
|
{
|
2020-06-19 15:39:41 +00:00
|
|
|
StorageInMemoryMetadata storage_metadata;
|
|
|
|
storage_metadata.setColumns(columns_);
|
|
|
|
storage_metadata.setConstraints(constraints_);
|
2021-04-23 12:18:23 +00:00
|
|
|
storage_metadata.setComment(comment);
|
2020-06-19 15:39:41 +00:00
|
|
|
setInMemoryMetadata(storage_metadata);
|
2019-08-24 21:20:20 +00:00
|
|
|
|
2019-10-25 19:07:47 +00:00
|
|
|
if (relative_path_.empty())
|
2023-01-23 21:13:58 +00:00
|
|
|
throw Exception(ErrorCodes::INCORRECT_FILE_NAME, "Storage {} requires data path", getName());
|
2014-06-26 00:58:14 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
/// Enumerate data files.
|
|
|
|
for (const auto & column : storage_metadata.getColumns().getAllPhysical())
|
|
|
|
addDataFiles(column);
|
|
|
|
|
|
|
|
/// Ensure the file checker is initialized.
|
|
|
|
if (file_checker.empty())
|
|
|
|
{
|
|
|
|
for (const auto & data_file : data_files)
|
|
|
|
file_checker.setEmpty(data_file.path);
|
|
|
|
if (use_marks_file)
|
|
|
|
file_checker.setEmpty(marks_file_path);
|
|
|
|
}
|
|
|
|
|
2020-07-12 02:31:58 +00:00
|
|
|
if (!attach)
|
|
|
|
{
|
|
|
|
/// create directories if they do not exist
|
|
|
|
disk->createDirectories(table_path);
|
|
|
|
}
|
|
|
|
else
|
|
|
|
{
|
|
|
|
try
|
|
|
|
{
|
|
|
|
file_checker.repair();
|
|
|
|
}
|
|
|
|
catch (...)
|
|
|
|
{
|
|
|
|
tryLogCurrentException(__PRETTY_FUNCTION__);
|
|
|
|
}
|
|
|
|
}
|
2022-05-03 17:55:45 +00:00
|
|
|
|
2024-01-15 16:13:32 +00:00
|
|
|
columns_with_collected_nested = ColumnsDescription{Nested::collect(columns_.getAll())};
|
2022-05-03 17:55:45 +00:00
|
|
|
total_bytes = file_checker.getTotalSize();
|
2012-08-29 20:07:24 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
void StorageLog::addDataFiles(const NameAndTypePair & column)
|
2012-08-29 20:07:24 +00:00
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
if (data_files_by_names.contains(column.name))
|
2023-01-23 21:13:58 +00:00
|
|
|
throw Exception(ErrorCodes::DUPLICATE_COLUMN, "Duplicate column with name {} in constructor of StorageLog.",
|
|
|
|
column.name);
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-03-09 14:46:52 +00:00
|
|
|
ISerialization::StreamCallback stream_callback = [&] (const ISerialization::SubstreamPath & substream_path)
|
2010-03-18 19:32:14 +00:00
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
String data_file_name = ISerialization::getFileNameForStream(column, substream_path);
|
|
|
|
if (!data_files_by_names.contains(data_file_name))
|
2013-07-16 14:55:01 +00:00
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
DataFile & data_file = data_files.emplace_back();
|
|
|
|
data_file.name = data_file_name;
|
|
|
|
data_file.path = table_path + data_file_name + DBMS_STORAGE_LOG_DATA_FILE_EXTENSION;
|
|
|
|
data_file.index = num_data_files++;
|
|
|
|
data_files_by_names.emplace(data_file_name, nullptr);
|
2013-07-16 14:55:01 +00:00
|
|
|
}
|
2017-08-07 07:31:16 +00:00
|
|
|
};
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
column.type->getDefaultSerialization()->enumerateStreams(stream_callback);
|
|
|
|
|
|
|
|
for (auto & data_file : data_files)
|
|
|
|
data_files_by_names[data_file.name] = &data_file;
|
2012-06-21 16:33:00 +00:00
|
|
|
}
|
2012-01-10 22:11:51 +00:00
|
|
|
|
2012-06-21 16:33:00 +00:00
|
|
|
|
2020-09-24 23:29:16 +00:00
|
|
|
void StorageLog::loadMarks(std::chrono::seconds lock_timeout)
|
2012-06-21 16:33:00 +00:00
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
if (!use_marks_file || marks_loaded)
|
|
|
|
return;
|
|
|
|
|
|
|
|
/// We load marks with an exclusive lock (i.e. the write lock) because we don't want
|
|
|
|
/// a data race between two threads trying to load marks simultaneously.
|
|
|
|
WriteLock lock{rwlock, lock_timeout};
|
2020-09-24 23:29:16 +00:00
|
|
|
if (!lock)
|
2023-01-23 21:13:58 +00:00
|
|
|
throw Exception(ErrorCodes::TIMEOUT_EXCEEDED, "Lock timeout exceeded");
|
2014-06-26 00:58:14 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
loadMarks(lock);
|
|
|
|
}
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2022-05-03 17:55:45 +00:00
|
|
|
void StorageLog::loadMarks(const WriteLock & lock /* already locked exclusively */)
|
2021-08-26 22:15:24 +00:00
|
|
|
{
|
|
|
|
if (!use_marks_file || marks_loaded)
|
|
|
|
return;
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
size_t num_marks = 0;
|
2019-12-25 08:24:13 +00:00
|
|
|
if (disk->exists(marks_file_path))
|
2013-02-26 13:06:01 +00:00
|
|
|
{
|
2019-12-25 08:24:13 +00:00
|
|
|
size_t file_size = disk->getFileSize(marks_file_path);
|
2021-08-26 22:15:24 +00:00
|
|
|
if (file_size % (num_data_files * sizeof(Mark)) != 0)
|
2023-01-23 21:13:58 +00:00
|
|
|
throw Exception(ErrorCodes::SIZES_OF_MARKS_FILES_ARE_INCONSISTENT, "Size of marks file is inconsistent");
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
num_marks = file_size / (num_data_files * sizeof(Mark));
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
for (auto & data_file : data_files)
|
|
|
|
data_file.marks.resize(num_marks);
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-08-24 21:45:58 +00:00
|
|
|
std::unique_ptr<ReadBuffer> marks_rb = disk->readFile(marks_file_path, ReadSettings().adjustBufferSize(32768));
|
2021-08-26 22:15:24 +00:00
|
|
|
for (size_t i = 0; i != num_marks; ++i)
|
2013-02-26 13:06:01 +00:00
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
for (auto & data_file : data_files)
|
2012-01-10 22:11:51 +00:00
|
|
|
{
|
|
|
|
Mark mark;
|
2021-08-26 22:15:24 +00:00
|
|
|
mark.read(*marks_rb);
|
|
|
|
data_file.marks[i] = mark;
|
2013-02-26 13:06:01 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
2017-08-07 07:31:16 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
marks_loaded = true;
|
|
|
|
num_marks_saved = num_marks;
|
2022-05-03 17:55:45 +00:00
|
|
|
|
|
|
|
/// We need marks to calculate the number of rows, and now we have the marks.
|
|
|
|
updateTotalRows(lock);
|
2021-08-26 22:15:24 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
void StorageLog::saveMarks(const WriteLock & /* already locked for writing */)
|
|
|
|
{
|
|
|
|
if (!use_marks_file)
|
|
|
|
return;
|
|
|
|
|
|
|
|
size_t num_marks = num_data_files ? data_files[0].marks.size() : 0;
|
|
|
|
if (num_marks_saved == num_marks)
|
|
|
|
return;
|
|
|
|
|
|
|
|
for (const auto & data_file : data_files)
|
|
|
|
{
|
|
|
|
if (data_file.marks.size() != num_marks)
|
2023-01-23 21:13:58 +00:00
|
|
|
throw Exception(ErrorCodes::LOGICAL_ERROR, "Wrong number of marks generated from block. Makes no sense.");
|
2021-08-26 22:15:24 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
size_t start = num_marks_saved;
|
|
|
|
auto marks_stream = disk->writeFile(marks_file_path, 4096, WriteMode::Append);
|
|
|
|
|
|
|
|
for (size_t i = start; i != num_marks; ++i)
|
|
|
|
{
|
|
|
|
for (const auto & data_file : data_files)
|
|
|
|
{
|
|
|
|
const auto & mark = data_file.marks[i];
|
|
|
|
mark.write(*marks_stream);
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
marks_stream->next();
|
|
|
|
marks_stream->finalize();
|
|
|
|
|
|
|
|
num_marks_saved = num_marks;
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
void StorageLog::removeUnsavedMarks(const WriteLock & /* already locked for writing */)
|
|
|
|
{
|
|
|
|
if (!use_marks_file)
|
|
|
|
return;
|
|
|
|
|
|
|
|
for (auto & data_file : data_files)
|
|
|
|
{
|
|
|
|
if (data_file.marks.size() > num_marks_saved)
|
|
|
|
data_file.marks.resize(num_marks_saved);
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
void StorageLog::saveFileSizes(const WriteLock & /* already locked for writing */)
|
|
|
|
{
|
|
|
|
for (const auto & data_file : data_files)
|
|
|
|
file_checker.update(data_file.path);
|
|
|
|
|
|
|
|
if (use_marks_file)
|
|
|
|
file_checker.update(marks_file_path);
|
|
|
|
|
|
|
|
file_checker.save();
|
2022-05-03 17:55:45 +00:00
|
|
|
total_bytes = file_checker.getTotalSize();
|
2010-03-18 19:32:14 +00:00
|
|
|
}
|
|
|
|
|
2013-02-07 13:03:19 +00:00
|
|
|
|
2020-04-07 14:05:51 +00:00
|
|
|
void StorageLog::rename(const String & new_path_to_table_data, const StorageID & new_table_id)
|
2012-06-18 06:19:13 +00:00
|
|
|
{
|
2020-09-18 19:25:56 +00:00
|
|
|
assert(table_path != new_path_to_table_data);
|
2020-09-17 19:50:43 +00:00
|
|
|
{
|
2022-08-05 19:41:02 +00:00
|
|
|
disk->createDirectories(new_path_to_table_data);
|
2020-09-17 19:50:43 +00:00
|
|
|
disk->moveDirectory(table_path, new_path_to_table_data);
|
2014-06-26 00:58:14 +00:00
|
|
|
|
2020-09-17 19:50:43 +00:00
|
|
|
table_path = new_path_to_table_data;
|
|
|
|
file_checker.setPath(table_path + "sizes.json");
|
2012-06-18 06:19:13 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
for (auto & data_file : data_files)
|
|
|
|
data_file.path = table_path + fileName(data_file.path);
|
2014-06-26 00:58:14 +00:00
|
|
|
|
2020-09-17 19:50:43 +00:00
|
|
|
marks_file_path = table_path + DBMS_STORAGE_LOG_MARKS_FILE_NAME;
|
|
|
|
}
|
2020-04-07 14:05:51 +00:00
|
|
|
renameInMemory(new_table_id);
|
2012-06-18 06:19:13 +00:00
|
|
|
}
|
|
|
|
|
2022-02-13 04:51:22 +00:00
|
|
|
static std::chrono::seconds getLockTimeout(ContextPtr context)
|
2018-04-21 00:35:20 +00:00
|
|
|
{
|
2022-02-13 04:51:22 +00:00
|
|
|
const Settings & settings = context->getSettingsRef();
|
|
|
|
Int64 lock_timeout = settings.lock_acquire_timeout.totalSeconds();
|
|
|
|
if (settings.max_execution_time.totalSeconds() != 0 && settings.max_execution_time.totalSeconds() < lock_timeout)
|
|
|
|
lock_timeout = settings.max_execution_time.totalSeconds();
|
|
|
|
return std::chrono::seconds{lock_timeout};
|
|
|
|
}
|
|
|
|
|
2022-07-13 20:35:24 +00:00
|
|
|
void StorageLog::truncate(const ASTPtr &, const StorageMetadataPtr &, ContextPtr local_context, TableExclusiveLockHolder &)
|
2022-02-13 04:51:22 +00:00
|
|
|
{
|
2022-07-13 20:35:24 +00:00
|
|
|
WriteLock lock{rwlock, getLockTimeout(local_context)};
|
2022-02-13 04:51:22 +00:00
|
|
|
if (!lock)
|
2023-01-23 21:13:58 +00:00
|
|
|
throw Exception(ErrorCodes::TIMEOUT_EXCEEDED, "Lock timeout exceeded");
|
2022-02-13 04:51:22 +00:00
|
|
|
|
2019-12-12 08:57:25 +00:00
|
|
|
disk->clearDirectory(table_path);
|
2018-04-21 00:35:20 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
for (auto & data_file : data_files)
|
|
|
|
{
|
|
|
|
data_file.marks.clear();
|
|
|
|
file_checker.setEmpty(data_file.path);
|
|
|
|
}
|
2018-04-21 00:35:20 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
if (use_marks_file)
|
|
|
|
file_checker.setEmpty(marks_file_path);
|
2012-06-18 06:19:13 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
marks_loaded = true;
|
|
|
|
num_marks_saved = 0;
|
2023-06-09 11:32:45 +00:00
|
|
|
total_rows = 0;
|
|
|
|
total_bytes = 0;
|
2023-06-29 10:02:41 +00:00
|
|
|
getContext()->clearMMappedFileCache();
|
2013-12-12 22:55:47 +00:00
|
|
|
}
|
|
|
|
|
2020-09-24 23:29:16 +00:00
|
|
|
|
2020-08-03 13:54:14 +00:00
|
|
|
Pipe StorageLog::read(
|
2011-08-09 15:57:33 +00:00
|
|
|
const Names & column_names,
|
2021-07-09 03:15:41 +00:00
|
|
|
const StorageSnapshotPtr & storage_snapshot,
|
2020-09-20 17:52:17 +00:00
|
|
|
SelectQueryInfo & /*query_info*/,
|
2022-07-13 20:35:24 +00:00
|
|
|
ContextPtr local_context,
|
2018-09-08 11:29:23 +00:00
|
|
|
QueryProcessingStage::Enum /*processed_stage*/,
|
2019-02-18 23:38:44 +00:00
|
|
|
size_t max_block_size,
|
2022-10-07 10:46:45 +00:00
|
|
|
size_t num_streams)
|
2010-03-18 19:32:14 +00:00
|
|
|
{
|
2021-07-09 03:15:41 +00:00
|
|
|
storage_snapshot->check(column_names);
|
2020-09-24 23:29:16 +00:00
|
|
|
|
2022-07-13 20:35:24 +00:00
|
|
|
auto lock_timeout = getLockTimeout(local_context);
|
2020-09-24 23:29:16 +00:00
|
|
|
loadMarks(lock_timeout);
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
ReadLock lock{rwlock, lock_timeout};
|
2020-09-24 23:29:16 +00:00
|
|
|
if (!lock)
|
2023-01-23 21:13:58 +00:00
|
|
|
throw Exception(ErrorCodes::TIMEOUT_EXCEEDED, "Lock timeout exceeded");
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
if (!num_data_files || !file_checker.getFileSize(data_files[INDEX_WITH_REAL_ROW_COUNT].path))
|
2021-11-09 12:36:25 +00:00
|
|
|
return Pipe(std::make_shared<NullSource>(storage_snapshot->getSampleBlockForColumns(column_names)));
|
2021-10-10 20:38:21 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
const Marks & marks_with_real_row_count = data_files[INDEX_WITH_REAL_ROW_COUNT].marks;
|
|
|
|
size_t num_marks = marks_with_real_row_count.size();
|
|
|
|
|
|
|
|
size_t max_streams = use_marks_file ? num_marks : 1;
|
|
|
|
if (num_streams > max_streams)
|
|
|
|
num_streams = max_streams;
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
std::vector<size_t> offsets;
|
|
|
|
offsets.resize(num_data_files, 0);
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-11-01 00:39:38 +00:00
|
|
|
std::vector<size_t> file_sizes;
|
|
|
|
file_sizes.resize(num_data_files, 0);
|
|
|
|
for (const auto & data_file : data_files)
|
|
|
|
file_sizes[data_file.index] = file_checker.getFileSize(data_file.path);
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-11-01 00:39:38 +00:00
|
|
|
/// For TinyLog (use_marks_file == false) there is no row limit and we just read
|
|
|
|
/// the data files up to their sizes.
|
|
|
|
bool limited_by_file_sizes = !use_marks_file;
|
|
|
|
size_t row_limit = std::numeric_limits<size_t>::max();
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2022-07-13 20:35:24 +00:00
|
|
|
ReadSettings read_settings = local_context->getReadSettings();
|
2021-08-26 22:15:24 +00:00
|
|
|
Pipes pipes;
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2024-01-15 16:13:32 +00:00
|
|
|
/// Converting to subcolumns of Nested is needed for
|
|
|
|
/// correct reading of parts of Nested with shared offsets.
|
|
|
|
auto options = GetColumnsOptions(GetColumnsOptions::All).withSubcolumns();
|
|
|
|
auto all_columns = storage_snapshot->getColumnsByNames(options, column_names);
|
|
|
|
all_columns = Nested::convertToSubcolumns(all_columns);
|
|
|
|
|
2017-08-07 07:31:16 +00:00
|
|
|
for (size_t stream = 0; stream < num_streams; ++stream)
|
2017-02-08 21:45:19 +00:00
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
if (use_marks_file)
|
|
|
|
{
|
|
|
|
size_t mark_begin = stream * num_marks / num_streams;
|
|
|
|
size_t mark_end = (stream + 1) * num_marks / num_streams;
|
2021-11-01 00:39:38 +00:00
|
|
|
size_t start_row = mark_begin ? marks_with_real_row_count[mark_begin - 1].rows : 0;
|
|
|
|
size_t end_row = mark_end ? marks_with_real_row_count[mark_end - 1].rows : 0;
|
|
|
|
row_limit = end_row - start_row;
|
2021-08-26 22:15:24 +00:00
|
|
|
for (const auto & data_file : data_files)
|
|
|
|
offsets[data_file.index] = data_file.marks[mark_begin].offset;
|
|
|
|
}
|
2017-11-27 21:21:09 +00:00
|
|
|
|
2020-01-31 15:10:10 +00:00
|
|
|
pipes.emplace_back(std::make_shared<LogSource>(
|
2017-08-07 07:31:16 +00:00
|
|
|
max_block_size,
|
2018-03-06 20:18:34 +00:00
|
|
|
all_columns,
|
2017-08-07 07:31:16 +00:00
|
|
|
*this,
|
2021-11-01 00:39:38 +00:00
|
|
|
row_limit,
|
2021-08-26 22:15:24 +00:00
|
|
|
offsets,
|
2021-11-01 00:39:38 +00:00
|
|
|
file_sizes,
|
|
|
|
limited_by_file_sizes,
|
2021-08-24 21:45:58 +00:00
|
|
|
read_settings));
|
2013-06-15 08:38:30 +00:00
|
|
|
}
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2020-09-24 23:29:16 +00:00
|
|
|
/// No need to hold lock while reading because we read fixed range of data that does not change while appending more data.
|
2020-08-06 12:24:05 +00:00
|
|
|
return Pipe::unitePipes(std::move(pipes));
|
2010-03-18 19:32:14 +00:00
|
|
|
}
|
|
|
|
|
2023-06-07 18:33:08 +00:00
|
|
|
SinkToStoragePtr StorageLog::write(const ASTPtr & /*query*/, const StorageMetadataPtr & metadata_snapshot, ContextPtr local_context, bool /*async_insert*/)
|
2010-03-18 19:32:14 +00:00
|
|
|
{
|
2022-07-13 20:35:24 +00:00
|
|
|
WriteLock lock{rwlock, getLockTimeout(local_context)};
|
2020-09-24 23:29:16 +00:00
|
|
|
if (!lock)
|
2023-01-23 21:13:58 +00:00
|
|
|
throw Exception(ErrorCodes::TIMEOUT_EXCEEDED, "Lock timeout exceeded");
|
2020-09-24 23:29:16 +00:00
|
|
|
|
2021-07-23 19:33:59 +00:00
|
|
|
return std::make_shared<LogSink>(*this, metadata_snapshot, std::move(lock));
|
2010-03-18 19:32:14 +00:00
|
|
|
}
|
|
|
|
|
2023-10-24 12:50:24 +00:00
|
|
|
IStorage::DataValidationTasksPtr StorageLog::getCheckTaskList(
|
|
|
|
const std::variant<std::monostate, ASTPtr, String> & check_task_filter, ContextPtr local_context)
|
2014-08-04 06:36:24 +00:00
|
|
|
{
|
2023-10-24 12:50:24 +00:00
|
|
|
if (!std::holds_alternative<std::monostate>(check_task_filter))
|
2023-10-23 12:13:36 +00:00
|
|
|
throw Exception(ErrorCodes::NOT_IMPLEMENTED, "CHECK PART/PARTITION are not supported for {}", getName());
|
|
|
|
|
2022-07-13 20:35:24 +00:00
|
|
|
ReadLock lock{rwlock, getLockTimeout(local_context)};
|
2020-09-24 23:29:16 +00:00
|
|
|
if (!lock)
|
2023-01-23 21:13:58 +00:00
|
|
|
throw Exception(ErrorCodes::TIMEOUT_EXCEEDED, "Lock timeout exceeded");
|
2023-10-23 12:13:36 +00:00
|
|
|
|
2023-08-14 09:58:08 +00:00
|
|
|
return std::make_unique<DataValidationTasks>(file_checker.getDataValidationTasks(), std::move(lock));
|
2023-08-01 10:57:59 +00:00
|
|
|
}
|
2017-12-30 00:36:06 +00:00
|
|
|
|
2023-10-23 10:12:30 +00:00
|
|
|
std::optional<CheckResult> StorageLog::checkDataNext(DataValidationTasksPtr & check_task_list)
|
2023-08-14 09:58:08 +00:00
|
|
|
{
|
2023-10-23 10:12:30 +00:00
|
|
|
return file_checker.checkNextEntry(assert_cast<DataValidationTasks *>(check_task_list.get())->file_checker_tasks);
|
2023-08-14 09:58:08 +00:00
|
|
|
}
|
2023-08-10 11:44:16 +00:00
|
|
|
|
2021-07-12 10:06:24 +00:00
|
|
|
IStorage::ColumnSizeByName StorageLog::getColumnSizes() const
|
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
ReadLock lock{rwlock, std::chrono::seconds(DBMS_DEFAULT_LOCK_ACQUIRE_TIMEOUT_SEC)};
|
2021-07-12 19:04:53 +00:00
|
|
|
if (!lock)
|
2023-01-23 21:13:58 +00:00
|
|
|
throw Exception(ErrorCodes::TIMEOUT_EXCEEDED, "Lock timeout exceeded");
|
2021-07-12 19:47:13 +00:00
|
|
|
|
2021-07-12 10:06:24 +00:00
|
|
|
ColumnSizeByName column_sizes;
|
2021-07-12 13:40:22 +00:00
|
|
|
|
2021-07-12 10:58:53 +00:00
|
|
|
for (const auto & column : getInMemoryMetadata().getColumns().getAllPhysical())
|
2021-07-12 10:06:24 +00:00
|
|
|
{
|
2021-07-12 13:40:22 +00:00
|
|
|
ISerialization::StreamCallback stream_callback = [&, this] (const ISerialization::SubstreamPath & substream_path)
|
2021-07-12 10:06:24 +00:00
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
String data_file_name = ISerialization::getFileNameForStream(column, substream_path);
|
|
|
|
auto it = data_files_by_names.find(data_file_name);
|
|
|
|
if (it != data_files_by_names.end())
|
|
|
|
{
|
|
|
|
const auto & data_file = *it->second;
|
|
|
|
column_sizes[column.name].data_compressed += file_checker.getFileSize(data_file.path);
|
|
|
|
}
|
2021-07-12 10:58:53 +00:00
|
|
|
};
|
|
|
|
|
|
|
|
auto serialization = column.type->getDefaultSerialization();
|
2021-10-11 22:01:00 +00:00
|
|
|
serialization->enumerateStreams(stream_callback);
|
2021-07-12 10:06:24 +00:00
|
|
|
}
|
2021-07-12 10:58:53 +00:00
|
|
|
|
2021-07-12 10:06:24 +00:00
|
|
|
return column_sizes;
|
|
|
|
}
|
|
|
|
|
2022-05-03 17:55:45 +00:00
|
|
|
void StorageLog::updateTotalRows(const WriteLock &)
|
|
|
|
{
|
|
|
|
if (!use_marks_file || !marks_loaded)
|
|
|
|
return;
|
|
|
|
|
|
|
|
if (num_data_files)
|
|
|
|
total_rows = data_files[INDEX_WITH_REAL_ROW_COUNT].marks.empty() ? 0 : data_files[INDEX_WITH_REAL_ROW_COUNT].marks.back().rows;
|
|
|
|
else
|
|
|
|
total_rows = 0;
|
|
|
|
}
|
|
|
|
|
|
|
|
std::optional<UInt64> StorageLog::totalRows(const Settings &) const
|
|
|
|
{
|
|
|
|
if (use_marks_file && marks_loaded)
|
|
|
|
return total_rows;
|
|
|
|
|
|
|
|
if (!total_bytes)
|
|
|
|
return 0;
|
|
|
|
|
|
|
|
return {};
|
|
|
|
}
|
|
|
|
|
|
|
|
std::optional<UInt64> StorageLog::totalBytes(const Settings &) const
|
|
|
|
{
|
|
|
|
return total_bytes;
|
|
|
|
}
|
2021-08-26 22:15:24 +00:00
|
|
|
|
2022-06-29 12:42:23 +00:00
|
|
|
void StorageLog::backupData(BackupEntriesCollector & backup_entries_collector, const String & data_path_in_backup, const std::optional<ASTs> & /* partitions */)
|
2021-10-26 09:48:31 +00:00
|
|
|
{
|
2023-05-03 11:51:36 +00:00
|
|
|
auto lock_timeout = getLockTimeout(backup_entries_collector.getContext());
|
2023-03-30 17:06:49 +00:00
|
|
|
|
2021-10-26 09:48:31 +00:00
|
|
|
loadMarks(lock_timeout);
|
|
|
|
|
|
|
|
ReadLock lock{rwlock, lock_timeout};
|
|
|
|
if (!lock)
|
2023-01-23 21:13:58 +00:00
|
|
|
throw Exception(ErrorCodes::TIMEOUT_EXCEEDED, "Lock timeout exceeded");
|
2021-10-26 09:48:31 +00:00
|
|
|
|
|
|
|
if (!num_data_files || !file_checker.getFileSize(data_files[INDEX_WITH_REAL_ROW_COUNT].path))
|
2022-05-29 19:53:56 +00:00
|
|
|
return;
|
2021-10-26 09:48:31 +00:00
|
|
|
|
2022-05-31 09:33:23 +00:00
|
|
|
fs::path data_path_in_backup_fs = data_path_in_backup;
|
2022-07-03 14:32:11 +00:00
|
|
|
auto temp_dir_owner = std::make_shared<TemporaryFileOnDisk>(disk, "tmp/");
|
2023-06-27 15:10:48 +00:00
|
|
|
fs::path temp_dir = temp_dir_owner->getRelativePath();
|
2021-10-26 09:48:31 +00:00
|
|
|
disk->createDirectories(temp_dir);
|
|
|
|
|
2023-07-23 09:23:36 +00:00
|
|
|
const auto & read_settings = backup_entries_collector.getReadSettings();
|
2023-05-03 23:27:16 +00:00
|
|
|
bool copy_encrypted = !backup_entries_collector.getBackupSettings().decrypt_files_from_encrypted_disks;
|
|
|
|
|
2021-10-26 09:48:31 +00:00
|
|
|
/// *.bin
|
|
|
|
for (const auto & data_file : data_files)
|
|
|
|
{
|
|
|
|
/// We make a copy of the data file because it can be changed later in write() or in truncate().
|
|
|
|
String data_file_name = fileName(data_file.path);
|
2022-05-29 19:53:56 +00:00
|
|
|
String hardlink_file_path = temp_dir / data_file_name;
|
2022-01-31 06:35:07 +00:00
|
|
|
disk->createHardLink(data_file.path, hardlink_file_path);
|
2023-04-23 10:25:46 +00:00
|
|
|
BackupEntryPtr backup_entry = std::make_unique<BackupEntryFromAppendOnlyFile>(
|
2023-05-03 23:27:16 +00:00
|
|
|
disk, hardlink_file_path, copy_encrypted, file_checker.getFileSize(data_file.path));
|
2023-04-23 10:25:46 +00:00
|
|
|
backup_entry = wrapBackupEntryWith(std::move(backup_entry), temp_dir_owner);
|
|
|
|
backup_entries_collector.addBackupEntry(data_path_in_backup_fs / data_file_name, std::move(backup_entry));
|
2021-10-26 09:48:31 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
/// __marks.mrk
|
|
|
|
if (use_marks_file)
|
|
|
|
{
|
|
|
|
/// We make a copy of the data file because it can be changed later in write() or in truncate().
|
|
|
|
String marks_file_name = fileName(marks_file_path);
|
2022-05-29 19:53:56 +00:00
|
|
|
String hardlink_file_path = temp_dir / marks_file_name;
|
2022-01-31 06:35:07 +00:00
|
|
|
disk->createHardLink(marks_file_path, hardlink_file_path);
|
2023-04-23 10:25:46 +00:00
|
|
|
BackupEntryPtr backup_entry = std::make_unique<BackupEntryFromAppendOnlyFile>(
|
2023-05-03 23:27:16 +00:00
|
|
|
disk, hardlink_file_path, copy_encrypted, file_checker.getFileSize(marks_file_path));
|
2023-04-23 10:25:46 +00:00
|
|
|
backup_entry = wrapBackupEntryWith(std::move(backup_entry), temp_dir_owner);
|
|
|
|
backup_entries_collector.addBackupEntry(data_path_in_backup_fs / marks_file_name, std::move(backup_entry));
|
2021-10-26 09:48:31 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
/// sizes.json
|
|
|
|
String files_info_path = file_checker.getPath();
|
2022-05-31 09:33:23 +00:00
|
|
|
backup_entries_collector.addBackupEntry(
|
2023-07-23 09:23:36 +00:00
|
|
|
data_path_in_backup_fs / fileName(files_info_path), std::make_unique<BackupEntryFromSmallFile>(disk, files_info_path, read_settings, copy_encrypted));
|
2021-10-26 09:48:31 +00:00
|
|
|
|
|
|
|
/// columns.txt
|
2022-05-31 09:33:23 +00:00
|
|
|
backup_entries_collector.addBackupEntry(
|
|
|
|
data_path_in_backup_fs / "columns.txt",
|
2022-05-29 19:53:56 +00:00
|
|
|
std::make_unique<BackupEntryFromMemory>(getInMemoryMetadata().getColumns().getAllPhysical().toString()));
|
2021-10-26 09:48:31 +00:00
|
|
|
|
|
|
|
/// count.txt
|
|
|
|
if (use_marks_file)
|
|
|
|
{
|
|
|
|
size_t num_rows = data_files[INDEX_WITH_REAL_ROW_COUNT].marks.empty() ? 0 : data_files[INDEX_WITH_REAL_ROW_COUNT].marks.back().rows;
|
2022-05-31 09:33:23 +00:00
|
|
|
backup_entries_collector.addBackupEntry(
|
|
|
|
data_path_in_backup_fs / "count.txt", std::make_unique<BackupEntryFromMemory>(toString(num_rows)));
|
2021-10-26 09:48:31 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2022-06-29 12:42:23 +00:00
|
|
|
void StorageLog::restoreDataFromBackup(RestorerFromBackup & restorer, const String & data_path_in_backup, const std::optional<ASTs> & /* partitions */)
|
2021-10-26 09:48:31 +00:00
|
|
|
{
|
2022-06-29 12:42:23 +00:00
|
|
|
auto backup = restorer.getBackup();
|
|
|
|
if (!backup->hasFiles(data_path_in_backup))
|
|
|
|
return;
|
2022-01-17 18:55:40 +00:00
|
|
|
|
2022-05-31 09:33:23 +00:00
|
|
|
if (!num_data_files)
|
|
|
|
return;
|
2021-10-26 09:48:31 +00:00
|
|
|
|
2022-06-29 12:42:23 +00:00
|
|
|
if (!restorer.isNonEmptyTableAllowed() && total_bytes)
|
2022-05-31 09:33:23 +00:00
|
|
|
RestorerFromBackup::throwTableIsNotEmpty(getStorageID());
|
2021-10-26 09:48:31 +00:00
|
|
|
|
2022-05-31 09:33:23 +00:00
|
|
|
auto lock_timeout = getLockTimeout(restorer.getContext());
|
|
|
|
restorer.addDataRestoreTask(
|
|
|
|
[storage = std::static_pointer_cast<StorageLog>(shared_from_this()), backup, data_path_in_backup, lock_timeout]
|
|
|
|
{ storage->restoreDataImpl(backup, data_path_in_backup, lock_timeout); });
|
|
|
|
}
|
2022-01-17 18:55:40 +00:00
|
|
|
|
2022-05-31 09:33:23 +00:00
|
|
|
void StorageLog::restoreDataImpl(const BackupPtr & backup, const String & data_path_in_backup, std::chrono::seconds lock_timeout)
|
|
|
|
{
|
|
|
|
WriteLock lock{rwlock, lock_timeout};
|
|
|
|
if (!lock)
|
2023-01-23 21:13:58 +00:00
|
|
|
throw Exception(ErrorCodes::TIMEOUT_EXCEEDED, "Lock timeout exceeded");
|
2021-10-26 09:48:31 +00:00
|
|
|
|
2022-05-31 09:33:23 +00:00
|
|
|
/// Load the marks if not loaded yet. We have to do that now because we're going to update these marks.
|
|
|
|
loadMarks(lock);
|
2021-10-26 09:48:31 +00:00
|
|
|
|
2022-05-31 09:33:23 +00:00
|
|
|
/// If there were no files, save zero file sizes to be able to rollback in case of error.
|
|
|
|
saveFileSizes(lock);
|
2021-10-26 09:48:31 +00:00
|
|
|
|
2022-05-31 09:33:23 +00:00
|
|
|
try
|
|
|
|
{
|
|
|
|
fs::path data_path_in_backup_fs = data_path_in_backup;
|
|
|
|
|
|
|
|
/// Append data files.
|
|
|
|
for (const auto & data_file : data_files)
|
2021-10-26 09:48:31 +00:00
|
|
|
{
|
2022-05-31 09:33:23 +00:00
|
|
|
String file_path_in_backup = data_path_in_backup_fs / fileName(data_file.path);
|
2022-06-29 12:42:23 +00:00
|
|
|
if (!backup->fileExists(file_path_in_backup))
|
2022-07-02 16:26:08 +00:00
|
|
|
throw Exception(ErrorCodes::CANNOT_RESTORE_TABLE, "File {} in backup is required to restore table", file_path_in_backup);
|
|
|
|
|
2023-03-13 22:43:15 +00:00
|
|
|
backup->copyFileToDisk(file_path_in_backup, disk, data_file.path, WriteMode::Append);
|
2022-05-31 09:33:23 +00:00
|
|
|
}
|
2021-10-26 09:48:31 +00:00
|
|
|
|
2022-05-31 09:33:23 +00:00
|
|
|
if (use_marks_file)
|
|
|
|
{
|
|
|
|
/// Append marks.
|
|
|
|
size_t num_extra_marks = 0;
|
|
|
|
String file_path_in_backup = data_path_in_backup_fs / fileName(marks_file_path);
|
2022-06-29 12:42:23 +00:00
|
|
|
if (!backup->fileExists(file_path_in_backup))
|
2022-07-02 16:26:08 +00:00
|
|
|
throw Exception(ErrorCodes::CANNOT_RESTORE_TABLE, "File {} in backup is required to restore table", file_path_in_backup);
|
|
|
|
|
2022-05-31 09:33:23 +00:00
|
|
|
size_t file_size = backup->getFileSize(file_path_in_backup);
|
|
|
|
if (file_size % (num_data_files * sizeof(Mark)) != 0)
|
2023-01-23 21:13:58 +00:00
|
|
|
throw Exception(ErrorCodes::SIZES_OF_MARKS_FILES_ARE_INCONSISTENT, "Size of marks file is inconsistent");
|
2022-05-31 09:33:23 +00:00
|
|
|
|
|
|
|
num_extra_marks = file_size / (num_data_files * sizeof(Mark));
|
|
|
|
|
|
|
|
size_t num_marks = data_files[0].marks.size();
|
|
|
|
for (auto & data_file : data_files)
|
|
|
|
data_file.marks.reserve(num_marks + num_extra_marks);
|
|
|
|
|
|
|
|
std::vector<size_t> old_data_sizes;
|
|
|
|
std::vector<size_t> old_num_rows;
|
|
|
|
old_data_sizes.resize(num_data_files);
|
|
|
|
old_num_rows.resize(num_data_files);
|
|
|
|
for (size_t i = 0; i != num_data_files; ++i)
|
2021-10-26 09:48:31 +00:00
|
|
|
{
|
2022-05-31 09:33:23 +00:00
|
|
|
old_data_sizes[i] = file_checker.getFileSize(data_files[i].path);
|
|
|
|
old_num_rows[i] = num_marks ? data_files[i].marks[num_marks - 1].rows : 0;
|
|
|
|
}
|
2021-10-26 09:48:31 +00:00
|
|
|
|
2023-03-13 22:43:15 +00:00
|
|
|
auto marks_rb = backup->readFile(file_path_in_backup);
|
2021-10-26 09:48:31 +00:00
|
|
|
|
2022-05-31 09:33:23 +00:00
|
|
|
for (size_t i = 0; i != num_extra_marks; ++i)
|
|
|
|
{
|
|
|
|
for (size_t j = 0; j != num_data_files; ++j)
|
2021-10-26 09:48:31 +00:00
|
|
|
{
|
2022-05-31 09:33:23 +00:00
|
|
|
Mark mark;
|
|
|
|
mark.read(*marks_rb);
|
|
|
|
mark.rows += old_num_rows[j]; /// Adjust the number of rows.
|
|
|
|
mark.offset += old_data_sizes[j]; /// Adjust the offset.
|
|
|
|
data_files[j].marks.push_back(mark);
|
2021-10-26 09:48:31 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2022-05-31 09:33:23 +00:00
|
|
|
/// Finish writing.
|
|
|
|
saveMarks(lock);
|
|
|
|
saveFileSizes(lock);
|
|
|
|
updateTotalRows(lock);
|
|
|
|
}
|
|
|
|
catch (...)
|
|
|
|
{
|
|
|
|
/// Rollback partial writes.
|
|
|
|
file_checker.repair();
|
|
|
|
removeUnsavedMarks(lock);
|
|
|
|
throw;
|
2022-01-17 18:55:40 +00:00
|
|
|
}
|
2021-10-26 09:48:31 +00:00
|
|
|
}
|
|
|
|
|
2017-12-30 00:36:06 +00:00
|
|
|
void registerStorageLog(StorageFactory & factory)
|
|
|
|
{
|
2020-02-18 14:41:30 +00:00
|
|
|
StorageFactory::StorageFeatures features{
|
|
|
|
.supports_settings = true
|
|
|
|
};
|
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
auto create_fn = [](const StorageFactory::Arguments & args)
|
2017-12-30 00:36:06 +00:00
|
|
|
{
|
|
|
|
if (!args.engine_args.empty())
|
2023-01-23 21:13:58 +00:00
|
|
|
throw Exception(ErrorCodes::NUMBER_OF_ARGUMENTS_DOESNT_MATCH, "Engine {} doesn't support any arguments ({} given)",
|
|
|
|
args.engine_name, args.engine_args.size());
|
2017-12-30 00:36:06 +00:00
|
|
|
|
2023-02-04 14:28:31 +00:00
|
|
|
String disk_name = getDiskName(*args.storage_def, args.getContext());
|
2021-04-10 23:33:54 +00:00
|
|
|
DiskPtr disk = args.getContext()->getDisk(disk_name);
|
2020-02-14 14:28:33 +00:00
|
|
|
|
2022-04-19 20:47:29 +00:00
|
|
|
return std::make_shared<StorageLog>(
|
2021-08-26 22:15:24 +00:00
|
|
|
args.engine_name,
|
2021-04-23 12:18:23 +00:00
|
|
|
disk,
|
|
|
|
args.relative_data_path,
|
|
|
|
args.table_id,
|
|
|
|
args.columns,
|
|
|
|
args.constraints,
|
|
|
|
args.comment,
|
|
|
|
args.attach,
|
2022-07-13 20:35:24 +00:00
|
|
|
args.getContext());
|
2021-08-26 22:15:24 +00:00
|
|
|
};
|
|
|
|
|
|
|
|
factory.registerStorage("Log", create_fn, features);
|
|
|
|
factory.registerStorage("TinyLog", create_fn, features);
|
2017-12-30 00:36:06 +00:00
|
|
|
}
|
|
|
|
|
2010-03-18 19:32:14 +00:00
|
|
|
}
|