2016-10-25 06:49:24 +00:00
|
|
|
#include <sys/stat.h>
|
|
|
|
#include <sys/types.h>
|
2018-05-09 04:22:30 +00:00
|
|
|
#include <errno.h>
|
2016-10-25 06:49:24 +00:00
|
|
|
|
2015-08-16 07:01:41 +00:00
|
|
|
#include <map>
|
2017-11-20 04:15:43 +00:00
|
|
|
#include <optional>
|
2015-08-16 07:01:41 +00:00
|
|
|
|
2017-04-01 09:19:00 +00:00
|
|
|
#include <Common/escapeForFileName.h>
|
|
|
|
#include <Common/Exception.h>
|
2015-08-16 07:01:41 +00:00
|
|
|
|
2020-02-20 16:39:32 +00:00
|
|
|
#include <IO/WriteBufferFromFileBase.h>
|
2021-10-26 09:48:31 +00:00
|
|
|
#include <Compression/CompressedReadBuffer.h>
|
2018-12-28 18:15:26 +00:00
|
|
|
#include <Compression/CompressedReadBufferFromFile.h>
|
|
|
|
#include <Compression/CompressedWriteBuffer.h>
|
2017-04-01 09:19:00 +00:00
|
|
|
#include <IO/ReadHelpers.h>
|
|
|
|
#include <IO/WriteHelpers.h>
|
2021-10-26 09:48:31 +00:00
|
|
|
#include <IO/copyData.h>
|
2015-08-16 07:01:41 +00:00
|
|
|
|
2021-10-15 20:18:20 +00:00
|
|
|
#include <Formats/NativeReader.h>
|
|
|
|
#include <Formats/NativeWriter.h>
|
2015-08-16 07:01:41 +00:00
|
|
|
|
2018-02-18 03:23:48 +00:00
|
|
|
#include <DataTypes/DataTypeFactory.h>
|
|
|
|
|
2017-04-01 09:19:00 +00:00
|
|
|
#include <Columns/ColumnArray.h>
|
2015-08-16 07:01:41 +00:00
|
|
|
|
2017-05-24 21:06:29 +00:00
|
|
|
#include <Interpreters/Context.h>
|
2017-01-21 04:24:28 +00:00
|
|
|
|
2020-02-14 14:28:33 +00:00
|
|
|
#include <Interpreters/evaluateConstantExpression.h>
|
|
|
|
#include <Parsers/ASTLiteral.h>
|
2017-12-30 00:36:06 +00:00
|
|
|
#include <Storages/StorageFactory.h>
|
2020-02-14 14:28:33 +00:00
|
|
|
#include <Storages/StorageStripeLog.h>
|
2020-02-18 14:41:30 +00:00
|
|
|
#include "StorageLogSettings.h"
|
2020-02-14 10:57:09 +00:00
|
|
|
#include <Processors/Sources/SourceWithProgress.h>
|
|
|
|
#include <Processors/Sources/NullSource.h>
|
2021-07-23 14:25:35 +00:00
|
|
|
#include <Processors/Sinks/SinkToStorage.h>
|
2021-10-16 14:03:50 +00:00
|
|
|
#include <QueryPipeline/Pipe.h>
|
2015-08-16 07:01:41 +00:00
|
|
|
|
2022-01-31 06:35:07 +00:00
|
|
|
#include <Backups/BackupEntryFromAppendOnlyFile.h>
|
2021-10-26 09:48:31 +00:00
|
|
|
#include <Backups/BackupEntryFromSmallFile.h>
|
|
|
|
#include <Backups/IBackup.h>
|
2022-01-19 18:56:08 +00:00
|
|
|
#include <Backups/IRestoreTask.h>
|
2021-10-26 09:48:31 +00:00
|
|
|
#include <Disks/TemporaryFileOnDisk.h>
|
|
|
|
|
|
|
|
#include <base/insertAtEnd.h>
|
2015-08-16 07:01:41 +00:00
|
|
|
|
2020-09-18 19:25:56 +00:00
|
|
|
#include <cassert>
|
|
|
|
|
2015-08-16 07:01:41 +00:00
|
|
|
|
|
|
|
namespace DB
|
|
|
|
{
|
|
|
|
|
2016-01-11 21:46:36 +00:00
|
|
|
namespace ErrorCodes
|
|
|
|
{
|
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;
|
2020-09-24 23:29:16 +00:00
|
|
|
extern const int TIMEOUT_EXCEEDED;
|
2021-10-26 09:48:31 +00:00
|
|
|
extern const int NOT_IMPLEMENTED;
|
2016-01-11 21:46:36 +00:00
|
|
|
}
|
|
|
|
|
2015-08-16 07:01:41 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
/// NOTE: The lock `StorageStripeLog::rwlock` is NOT kept locked while reading,
|
|
|
|
/// because we read ranges of data that do not change.
|
2020-02-14 10:57:09 +00:00
|
|
|
class StripeLogSource final : public SourceWithProgress
|
2015-08-16 07:01:41 +00:00
|
|
|
{
|
|
|
|
public:
|
2020-02-14 10:57:09 +00:00
|
|
|
static Block getHeader(
|
2021-07-09 03:15:41 +00:00
|
|
|
const StorageSnapshotPtr & storage_snapshot,
|
2020-02-14 10:57:09 +00:00
|
|
|
const Names & column_names,
|
|
|
|
IndexForNativeFormat::Blocks::const_iterator index_begin,
|
|
|
|
IndexForNativeFormat::Blocks::const_iterator index_end)
|
2017-04-01 07:20:54 +00:00
|
|
|
{
|
2020-02-14 10:57:09 +00:00
|
|
|
if (index_begin == index_end)
|
2021-07-09 03:15:41 +00:00
|
|
|
return storage_snapshot->getSampleBlockForColumns(column_names);
|
2020-02-14 10:57:09 +00:00
|
|
|
|
|
|
|
/// TODO: check if possible to always return storage.getSampleBlock()
|
|
|
|
|
|
|
|
Block header;
|
|
|
|
|
|
|
|
for (const auto & column : index_begin->columns)
|
2018-02-21 04:38:26 +00:00
|
|
|
{
|
2020-02-14 10:57:09 +00:00
|
|
|
auto type = DataTypeFactory::instance().get(column.type);
|
|
|
|
header.insert(ColumnWithTypeAndName{ type, column.name });
|
2018-02-21 04:38:26 +00:00
|
|
|
}
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2020-02-14 10:57:09 +00:00
|
|
|
return header;
|
|
|
|
}
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2020-06-16 14:25:08 +00:00
|
|
|
StripeLogSource(
|
2021-08-26 22:15:24 +00:00
|
|
|
const StorageStripeLog & storage_,
|
2021-07-09 03:15:41 +00:00
|
|
|
const StorageSnapshotPtr & storage_snapshot_,
|
2020-06-16 14:25:08 +00:00
|
|
|
const Names & column_names,
|
2021-08-24 21:45:58 +00:00
|
|
|
ReadSettings read_settings_,
|
2021-08-26 22:15:24 +00:00
|
|
|
std::shared_ptr<const IndexForNativeFormat> indices_,
|
2020-02-14 10:57:09 +00:00
|
|
|
IndexForNativeFormat::Blocks::const_iterator index_begin_,
|
2021-11-01 00:39:38 +00:00
|
|
|
IndexForNativeFormat::Blocks::const_iterator index_end_,
|
|
|
|
size_t file_size_)
|
2021-11-09 12:36:25 +00:00
|
|
|
: SourceWithProgress(getHeader(storage_snapshot_, column_names, index_begin_, index_end_))
|
2020-06-16 14:25:08 +00:00
|
|
|
, storage(storage_)
|
2021-07-09 03:15:41 +00:00
|
|
|
, storage_snapshot(storage_snapshot_)
|
2021-08-24 21:45:58 +00:00
|
|
|
, read_settings(std::move(read_settings_))
|
2021-08-26 22:15:24 +00:00
|
|
|
, indices(indices_)
|
2020-06-16 14:25:08 +00:00
|
|
|
, index_begin(index_begin_)
|
|
|
|
, index_end(index_end_)
|
2021-11-01 00:39:38 +00:00
|
|
|
, file_size(file_size_)
|
2018-01-09 01:51:08 +00:00
|
|
|
{
|
2018-08-10 04:02:56 +00:00
|
|
|
}
|
2018-01-06 18:10:44 +00:00
|
|
|
|
2020-02-14 10:57:09 +00:00
|
|
|
String getName() const override { return "StripeLog"; }
|
|
|
|
|
2015-08-16 07:01:41 +00:00
|
|
|
protected:
|
2020-02-14 10:57:09 +00:00
|
|
|
Chunk generate() override
|
2017-04-01 07:20:54 +00:00
|
|
|
{
|
|
|
|
Block res;
|
2018-01-09 01:51:08 +00:00
|
|
|
start();
|
2015-12-16 02:32:49 +00:00
|
|
|
|
2017-04-01 07:20:54 +00:00
|
|
|
if (block_in)
|
|
|
|
{
|
|
|
|
res = block_in->read();
|
2015-10-01 03:30:50 +00:00
|
|
|
|
2017-04-01 07:20:54 +00:00
|
|
|
/// Freeing memory before destroying the object.
|
|
|
|
if (!res)
|
|
|
|
{
|
2017-11-20 04:41:56 +00:00
|
|
|
block_in.reset();
|
|
|
|
data_in.reset();
|
2021-08-26 22:15:24 +00:00
|
|
|
indices.reset();
|
2017-04-01 07:20:54 +00:00
|
|
|
}
|
|
|
|
}
|
2015-10-01 03:30:50 +00:00
|
|
|
|
2020-02-14 10:57:09 +00:00
|
|
|
return Chunk(res.getColumns(), res.rows());
|
2017-04-01 07:20:54 +00:00
|
|
|
}
|
2015-08-16 07:01:41 +00:00
|
|
|
|
|
|
|
private:
|
2021-08-26 22:15:24 +00:00
|
|
|
const StorageStripeLog & storage;
|
2021-07-09 03:15:41 +00:00
|
|
|
StorageSnapshotPtr storage_snapshot;
|
2021-08-24 21:45:58 +00:00
|
|
|
ReadSettings read_settings;
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
std::shared_ptr<const IndexForNativeFormat> indices;
|
2017-04-01 07:20:54 +00:00
|
|
|
IndexForNativeFormat::Blocks::const_iterator index_begin;
|
|
|
|
IndexForNativeFormat::Blocks::const_iterator index_end;
|
2021-11-01 00:39:38 +00:00
|
|
|
size_t file_size;
|
2021-08-26 22:15:24 +00:00
|
|
|
|
2018-02-21 04:38:26 +00:00
|
|
|
Block header;
|
2017-04-01 07:20:54 +00:00
|
|
|
|
|
|
|
/** optional - to create objects only on first reading
|
|
|
|
* and delete objects (release buffers) after the source is exhausted
|
|
|
|
* - to save RAM when using a large number of sources.
|
|
|
|
*/
|
|
|
|
bool started = false;
|
2017-11-20 04:15:43 +00:00
|
|
|
std::optional<CompressedReadBufferFromFile> data_in;
|
2021-10-08 17:21:19 +00:00
|
|
|
std::optional<NativeReader> block_in;
|
2018-01-09 01:51:08 +00:00
|
|
|
|
|
|
|
void start()
|
|
|
|
{
|
|
|
|
if (!started)
|
|
|
|
{
|
|
|
|
started = true;
|
|
|
|
|
2019-12-25 08:24:13 +00:00
|
|
|
String data_file_path = storage.table_path + "data.bin";
|
2021-11-01 00:39:38 +00:00
|
|
|
data_in.emplace(storage.disk->readFile(data_file_path, read_settings.adjustBufferSize(file_size)));
|
2018-02-18 02:46:39 +00:00
|
|
|
block_in.emplace(*data_in, 0, index_begin, index_end);
|
2018-01-09 01:51:08 +00:00
|
|
|
}
|
|
|
|
}
|
2015-08-16 07:01:41 +00:00
|
|
|
};
|
|
|
|
|
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
/// NOTE: The lock `StorageStripeLog::rwlock` is kept locked in exclusive mode while writing.
|
2021-07-23 14:25:35 +00:00
|
|
|
class StripeLogSink final : public SinkToStorage
|
2015-08-16 07:01:41 +00:00
|
|
|
{
|
|
|
|
public:
|
2021-08-26 22:15:24 +00:00
|
|
|
using WriteLock = std::unique_lock<std::shared_timed_mutex>;
|
|
|
|
|
2021-07-23 14:25:35 +00:00
|
|
|
explicit StripeLogSink(
|
2021-08-26 22:15:24 +00:00
|
|
|
StorageStripeLog & storage_, const StorageMetadataPtr & metadata_snapshot_, WriteLock && lock_)
|
2021-07-26 10:08:40 +00:00
|
|
|
: SinkToStorage(metadata_snapshot_->getSampleBlock())
|
2021-07-23 14:25:35 +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_))
|
2021-08-26 22:15:24 +00:00
|
|
|
, data_out_compressed(storage.disk->writeFile(storage.data_file_path, DBMS_DEFAULT_BUFFER_SIZE, WriteMode::Append))
|
2020-07-12 02:31:58 +00:00
|
|
|
, data_out(std::make_unique<CompressedWriteBuffer>(
|
2021-08-26 22:15:24 +00:00
|
|
|
*data_out_compressed, CompressionCodecFactory::instance().getDefaultCodec(), storage.max_compress_block_size))
|
2017-04-01 07:20:54 +00:00
|
|
|
{
|
2020-09-24 23:29:16 +00:00
|
|
|
if (!lock)
|
|
|
|
throw Exception("Lock timeout exceeded", ErrorCodes::TIMEOUT_EXCEEDED);
|
2021-01-05 01:49:15 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
/// Ensure that indices are loaded because we're going to update them.
|
|
|
|
storage.loadIndices(lock);
|
|
|
|
|
|
|
|
/// If there were no files, save zero file sizes to be able to rollback in case of error.
|
|
|
|
storage.saveFileSizes(lock);
|
|
|
|
|
|
|
|
size_t initial_data_size = storage.file_checker.getFileSize(storage.data_file_path);
|
|
|
|
block_out = std::make_unique<NativeWriter>(*data_out, 0, metadata_snapshot->getSampleBlock(), false, &storage.indices, initial_data_size);
|
2017-04-01 07:20:54 +00:00
|
|
|
}
|
|
|
|
|
2021-07-23 14:25:35 +00:00
|
|
|
String getName() const override { return "StripeLogSink"; }
|
|
|
|
|
|
|
|
~StripeLogSink() override
|
2017-04-01 07:20:54 +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
|
|
|
data_out.reset();
|
|
|
|
data_out_compressed.reset();
|
|
|
|
|
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 indices.
|
|
|
|
storage.removeUnsavedIndices(lock);
|
2020-07-12 02:31:58 +00:00
|
|
|
}
|
2017-04-01 07:20:54 +00:00
|
|
|
}
|
|
|
|
catch (...)
|
|
|
|
{
|
|
|
|
tryLogCurrentException(__PRETTY_FUNCTION__);
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2021-07-23 14:25:35 +00:00
|
|
|
void consume(Chunk chunk) override
|
2017-04-01 07:20:54 +00:00
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
block_out->write(getHeader().cloneWithColumns(chunk.detachColumns()));
|
2017-04-01 07:20:54 +00:00
|
|
|
}
|
|
|
|
|
2021-07-23 14:25:35 +00:00
|
|
|
void onFinish() override
|
2017-04-01 07:20:54 +00:00
|
|
|
{
|
|
|
|
if (done)
|
|
|
|
return;
|
|
|
|
|
2020-07-12 02:31:58 +00:00
|
|
|
data_out->next();
|
2019-12-12 08:57:25 +00:00
|
|
|
data_out_compressed->next();
|
2020-07-30 13:42:05 +00:00
|
|
|
data_out_compressed->finalize();
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
/// Save the new indices.
|
|
|
|
storage.saveIndices(lock);
|
|
|
|
|
|
|
|
/// Save the new file sizes.
|
|
|
|
storage.saveFileSizes(lock);
|
2017-04-01 07:20:54 +00:00
|
|
|
|
|
|
|
done = true;
|
2021-04-04 05:24:01 +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();
|
2017-04-01 07:20:54 +00:00
|
|
|
}
|
2015-08-16 07:01:41 +00:00
|
|
|
|
|
|
|
private:
|
2017-04-01 07:20:54 +00:00
|
|
|
StorageStripeLog & storage;
|
2020-06-16 15:51:29 +00:00
|
|
|
StorageMetadataPtr metadata_snapshot;
|
2021-08-26 22:15:24 +00:00
|
|
|
WriteLock lock;
|
2015-08-16 07:01:41 +00:00
|
|
|
|
2019-12-12 08:57:25 +00:00
|
|
|
std::unique_ptr<WriteBuffer> data_out_compressed;
|
2020-07-12 02:31:58 +00:00
|
|
|
std::unique_ptr<CompressedWriteBuffer> data_out;
|
2021-08-26 22:15:24 +00:00
|
|
|
std::unique_ptr<NativeWriter> block_out;
|
2015-08-16 07:01:41 +00:00
|
|
|
|
2017-04-01 07:20:54 +00:00
|
|
|
bool done = false;
|
2015-08-16 07:01:41 +00:00
|
|
|
};
|
|
|
|
|
|
|
|
|
|
|
|
StorageStripeLog::StorageStripeLog(
|
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,
|
2017-04-01 07:20:54 +00:00
|
|
|
bool attach,
|
|
|
|
size_t max_compress_block_size_)
|
2019-12-04 16:06:55 +00:00
|
|
|
: IStorage(table_id_)
|
2020-01-13 11:41:42 +00:00
|
|
|
, disk(std::move(disk_))
|
|
|
|
, table_path(relative_path_)
|
2021-08-26 22:15:24 +00:00
|
|
|
, data_file_path(table_path + "data.bin")
|
|
|
|
, index_file_path(table_path + "index.mrk")
|
2020-01-13 11:41:42 +00:00
|
|
|
, file_checker(disk, table_path + "sizes.json")
|
2021-08-26 22:15:24 +00:00
|
|
|
, max_compress_block_size(max_compress_block_size_)
|
2020-05-30 21:57:37 +00:00
|
|
|
, log(&Poco::Logger::get("StorageStripeLog"))
|
2015-08-16 07:01:41 +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())
|
2017-11-03 19:53:10 +00:00
|
|
|
throw Exception("Storage " + getName() + " requires data path", ErrorCodes::INCORRECT_FILE_NAME);
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
/// Ensure the file checker is initialized.
|
|
|
|
if (file_checker.empty())
|
|
|
|
{
|
|
|
|
file_checker.setEmpty(data_file_path);
|
|
|
|
file_checker.setEmpty(index_file_path);
|
|
|
|
}
|
|
|
|
|
2017-04-01 07:20:54 +00:00
|
|
|
if (!attach)
|
|
|
|
{
|
2019-12-12 08:57:25 +00:00
|
|
|
/// create directories if they do not exist
|
|
|
|
disk->createDirectories(table_path);
|
2020-07-12 02:31:58 +00:00
|
|
|
}
|
|
|
|
else
|
|
|
|
{
|
|
|
|
try
|
|
|
|
{
|
|
|
|
file_checker.repair();
|
|
|
|
}
|
|
|
|
catch (...)
|
|
|
|
{
|
|
|
|
tryLogCurrentException(__PRETTY_FUNCTION__);
|
|
|
|
}
|
2017-04-01 07:20:54 +00:00
|
|
|
}
|
2015-08-16 07:01:41 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
StorageStripeLog::~StorageStripeLog() = default;
|
|
|
|
|
|
|
|
|
2020-04-07 14:05:51 +00:00
|
|
|
void StorageStripeLog::rename(const String & new_path_to_table_data, const StorageID & new_table_id)
|
2015-08-16 07:01:41 +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
|
|
|
{
|
|
|
|
disk->moveDirectory(table_path, new_path_to_table_data);
|
2015-08-16 07:01:41 +00:00
|
|
|
|
2020-09-17 19:50:43 +00:00
|
|
|
table_path = new_path_to_table_data;
|
2021-08-26 22:15:24 +00:00
|
|
|
data_file_path = table_path + "data.bin";
|
|
|
|
index_file_path = table_path + "index.mrk";
|
2020-09-17 19:50:43 +00:00
|
|
|
file_checker.setPath(table_path + "sizes.json");
|
|
|
|
}
|
2020-04-07 14:05:51 +00:00
|
|
|
renameInMemory(new_table_id);
|
2015-08-16 07:01:41 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
|
2021-04-10 23:33:54 +00:00
|
|
|
static std::chrono::seconds getLockTimeout(ContextPtr context)
|
2020-09-24 23:29:16 +00:00
|
|
|
{
|
2021-04-10 23:33:54 +00:00
|
|
|
const Settings & settings = context->getSettingsRef();
|
2020-09-24 23:29:16 +00:00
|
|
|
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};
|
|
|
|
}
|
|
|
|
|
|
|
|
|
2020-08-03 13:54:14 +00:00
|
|
|
Pipe StorageStripeLog::read(
|
2017-04-01 07:20:54 +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*/,
|
2021-04-10 23:33:54 +00:00
|
|
|
ContextPtr context,
|
2018-09-08 11:29:23 +00:00
|
|
|
QueryProcessingStage::Enum /*processed_stage*/,
|
2017-12-01 21:13:25 +00:00
|
|
|
const size_t /*max_block_size*/,
|
2017-06-02 15:54:39 +00:00
|
|
|
unsigned num_streams)
|
2015-08-16 07:01:41 +00:00
|
|
|
{
|
2021-07-09 03:15:41 +00:00
|
|
|
storage_snapshot->check(column_names);
|
2015-08-16 08:18:34 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
auto lock_timeout = getLockTimeout(context);
|
|
|
|
loadIndices(lock_timeout);
|
2015-08-16 08:18:34 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
ReadLock lock{rwlock, lock_timeout};
|
|
|
|
if (!lock)
|
|
|
|
throw Exception("Lock timeout exceeded", ErrorCodes::TIMEOUT_EXCEEDED);
|
2020-02-14 10:57:09 +00:00
|
|
|
|
2021-11-01 00:39:38 +00:00
|
|
|
size_t data_file_size = file_checker.getFileSize(data_file_path);
|
|
|
|
if (!data_file_size)
|
2021-07-09 03:15:41 +00:00
|
|
|
return Pipe(std::make_shared<NullSource>(storage_snapshot->getSampleBlockForColumns(column_names)));
|
2015-12-08 20:04:11 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
auto indices_for_selected_columns
|
|
|
|
= std::make_shared<IndexForNativeFormat>(indices.extractIndexForColumns(NameSet{column_names.begin(), column_names.end()}));
|
2021-08-24 21:45:58 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
size_t size = indices_for_selected_columns->blocks.size();
|
2017-06-02 15:54:39 +00:00
|
|
|
if (num_streams > size)
|
|
|
|
num_streams = size;
|
2015-08-16 08:18:34 +00:00
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
ReadSettings read_settings = context->getReadSettings();
|
|
|
|
Pipes pipes;
|
|
|
|
|
2017-06-02 15:54:39 +00:00
|
|
|
for (size_t stream = 0; stream < num_streams; ++stream)
|
2017-04-01 07:20:54 +00:00
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
IndexForNativeFormat::Blocks::const_iterator begin = indices_for_selected_columns->blocks.begin();
|
|
|
|
IndexForNativeFormat::Blocks::const_iterator end = indices_for_selected_columns->blocks.begin();
|
2015-08-16 08:18:34 +00:00
|
|
|
|
2017-06-02 15:54:39 +00:00
|
|
|
std::advance(begin, stream * size / num_streams);
|
|
|
|
std::advance(end, (stream + 1) * size / num_streams);
|
2015-08-16 08:18:34 +00:00
|
|
|
|
2020-02-14 10:57:09 +00:00
|
|
|
pipes.emplace_back(std::make_shared<StripeLogSource>(
|
2021-11-09 12:36:25 +00:00
|
|
|
*this, storage_snapshot, column_names, read_settings, indices_for_selected_columns, begin, end, data_file_size));
|
2017-04-01 07:20:54 +00:00
|
|
|
}
|
2015-08-16 08:18:34 +00:00
|
|
|
|
2017-04-01 07:20:54 +00:00
|
|
|
/// We do not keep read lock directly at the time of reading, because we read ranges of data that do not change.
|
2015-08-16 08:18:34 +00:00
|
|
|
|
2020-08-06 12:24:05 +00:00
|
|
|
return Pipe::unitePipes(std::move(pipes));
|
2015-08-16 07:01:41 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
|
2021-07-23 14:25:35 +00:00
|
|
|
SinkToStoragePtr StorageStripeLog::write(const ASTPtr & /*query*/, const StorageMetadataPtr & metadata_snapshot, ContextPtr context)
|
2015-08-16 07:01:41 +00:00
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
WriteLock lock{rwlock, getLockTimeout(context)};
|
2020-09-24 23:29:16 +00:00
|
|
|
if (!lock)
|
|
|
|
throw Exception("Lock timeout exceeded", ErrorCodes::TIMEOUT_EXCEEDED);
|
|
|
|
|
2021-07-23 14:25:35 +00:00
|
|
|
return std::make_shared<StripeLogSink>(*this, metadata_snapshot, std::move(lock));
|
2015-08-16 07:01:41 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
|
2021-04-10 23:33:54 +00:00
|
|
|
CheckResults StorageStripeLog::checkData(const ASTPtr & /* query */, ContextPtr context)
|
2015-08-16 07:01:41 +00:00
|
|
|
{
|
2021-08-26 22:15:24 +00:00
|
|
|
ReadLock lock{rwlock, getLockTimeout(context)};
|
2020-09-24 23:29:16 +00:00
|
|
|
if (!lock)
|
|
|
|
throw Exception("Lock timeout exceeded", ErrorCodes::TIMEOUT_EXCEEDED);
|
|
|
|
|
2017-04-01 07:20:54 +00:00
|
|
|
return file_checker.check();
|
2015-08-16 07:01:41 +00:00
|
|
|
}
|
|
|
|
|
2021-08-26 22:15:24 +00:00
|
|
|
|
2021-04-10 23:33:54 +00:00
|
|
|
void StorageStripeLog::truncate(const ASTPtr &, const StorageMetadataPtr &, ContextPtr, TableExclusiveLockHolder &)
|
2018-04-21 00:35:20 +00:00
|
|
|
{
|
2019-12-12 08:57:25 +00:00
|
|
|
disk->clearDirectory(table_path);
|
2021-08-26 22:15:24 +00:00
|
|
|
|
|
|
|
indices.clear();
|
|
|
|
file_checker.setEmpty(data_file_path);
|
|
|
|
file_checker.setEmpty(index_file_path);
|
|
|
|
|
|
|
|
indices_loaded = true;
|
|
|
|
num_indices_saved = 0;
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
void StorageStripeLog::loadIndices(std::chrono::seconds lock_timeout)
|
|
|
|
{
|
|
|
|
if (indices_loaded)
|
|
|
|
return;
|
|
|
|
|
|
|
|
/// We load indices with an exclusive lock (i.e. the write lock) because we don't want
|
|
|
|
/// a data race between two threads trying to load indices simultaneously.
|
|
|
|
WriteLock lock{rwlock, lock_timeout};
|
|
|
|
if (!lock)
|
|
|
|
throw Exception("Lock timeout exceeded", ErrorCodes::TIMEOUT_EXCEEDED);
|
|
|
|
|
|
|
|
loadIndices(lock);
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
void StorageStripeLog::loadIndices(const WriteLock & /* already locked exclusively */)
|
|
|
|
{
|
|
|
|
if (indices_loaded)
|
|
|
|
return;
|
|
|
|
|
|
|
|
if (disk->exists(index_file_path))
|
|
|
|
{
|
|
|
|
CompressedReadBufferFromFile index_in(disk->readFile(index_file_path, ReadSettings{}.adjustBufferSize(4096)));
|
|
|
|
indices.read(index_in);
|
|
|
|
}
|
|
|
|
|
|
|
|
indices_loaded = true;
|
|
|
|
num_indices_saved = indices.blocks.size();
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
void StorageStripeLog::saveIndices(const WriteLock & /* already locked for writing */)
|
|
|
|
{
|
|
|
|
size_t num_indices = indices.blocks.size();
|
|
|
|
if (num_indices_saved == num_indices)
|
|
|
|
return;
|
|
|
|
|
|
|
|
size_t start = num_indices_saved;
|
|
|
|
auto index_out_compressed = disk->writeFile(index_file_path, DBMS_DEFAULT_BUFFER_SIZE, WriteMode::Append);
|
|
|
|
auto index_out = std::make_unique<CompressedWriteBuffer>(*index_out_compressed);
|
|
|
|
|
|
|
|
for (size_t i = start; i != num_indices; ++i)
|
|
|
|
indices.blocks[i].write(*index_out);
|
|
|
|
|
|
|
|
index_out->next();
|
|
|
|
index_out_compressed->next();
|
|
|
|
index_out_compressed->finalize();
|
|
|
|
|
|
|
|
num_indices_saved = num_indices;
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
void StorageStripeLog::removeUnsavedIndices(const WriteLock & /* already locked for writing */)
|
|
|
|
{
|
|
|
|
if (indices.blocks.size() > num_indices_saved)
|
|
|
|
indices.blocks.resize(num_indices_saved);
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
void StorageStripeLog::saveFileSizes(const WriteLock & /* already locked for writing */)
|
|
|
|
{
|
|
|
|
file_checker.update(data_file_path);
|
|
|
|
file_checker.update(index_file_path);
|
|
|
|
file_checker.save();
|
2018-04-21 00:35:20 +00:00
|
|
|
}
|
|
|
|
|
2017-12-30 00:36:06 +00:00
|
|
|
|
2022-02-22 13:31:50 +00:00
|
|
|
BackupEntries StorageStripeLog::backupData(ContextPtr context, const ASTs & partitions)
|
2021-10-26 09:48:31 +00:00
|
|
|
{
|
|
|
|
if (!partitions.empty())
|
|
|
|
throw Exception(ErrorCodes::NOT_IMPLEMENTED, "Table engine {} doesn't support partitions", getName());
|
|
|
|
|
|
|
|
auto lock_timeout = getLockTimeout(context);
|
|
|
|
loadIndices(lock_timeout);
|
|
|
|
|
|
|
|
ReadLock lock{rwlock, lock_timeout};
|
|
|
|
if (!lock)
|
|
|
|
throw Exception("Lock timeout exceeded", ErrorCodes::TIMEOUT_EXCEEDED);
|
|
|
|
|
|
|
|
if (!file_checker.getFileSize(data_file_path))
|
|
|
|
return {};
|
|
|
|
|
|
|
|
auto temp_dir_owner = std::make_shared<TemporaryFileOnDisk>(disk, "tmp/backup_");
|
|
|
|
auto temp_dir = temp_dir_owner->getPath();
|
|
|
|
disk->createDirectories(temp_dir);
|
|
|
|
|
|
|
|
BackupEntries backup_entries;
|
|
|
|
|
|
|
|
/// data.bin
|
|
|
|
{
|
|
|
|
/// 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-01-31 06:35:07 +00:00
|
|
|
String hardlink_file_path = temp_dir + "/" + data_file_name;
|
|
|
|
disk->createHardLink(data_file_path, hardlink_file_path);
|
2021-10-26 09:48:31 +00:00
|
|
|
backup_entries.emplace_back(
|
|
|
|
data_file_name,
|
2022-01-31 06:35:07 +00:00
|
|
|
std::make_unique<BackupEntryFromAppendOnlyFile>(
|
|
|
|
disk, hardlink_file_path, file_checker.getFileSize(data_file_path), std::nullopt, temp_dir_owner));
|
2021-10-26 09:48:31 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
/// index.mrk
|
|
|
|
{
|
|
|
|
/// We make a copy of the data file because it can be changed later in write() or in truncate().
|
|
|
|
String index_file_name = fileName(index_file_path);
|
2022-01-31 06:35:07 +00:00
|
|
|
String hardlink_file_path = temp_dir + "/" + index_file_name;
|
|
|
|
disk->createHardLink(index_file_path, hardlink_file_path);
|
2021-10-26 09:48:31 +00:00
|
|
|
backup_entries.emplace_back(
|
|
|
|
index_file_name,
|
2022-01-31 06:35:07 +00:00
|
|
|
std::make_unique<BackupEntryFromAppendOnlyFile>(
|
|
|
|
disk, hardlink_file_path, file_checker.getFileSize(index_file_path), std::nullopt, temp_dir_owner));
|
2021-10-26 09:48:31 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
/// sizes.json
|
|
|
|
String files_info_path = file_checker.getPath();
|
|
|
|
backup_entries.emplace_back(fileName(files_info_path), std::make_unique<BackupEntryFromSmallFile>(disk, files_info_path));
|
|
|
|
|
|
|
|
/// columns.txt
|
|
|
|
backup_entries.emplace_back(
|
|
|
|
"columns.txt", std::make_unique<BackupEntryFromMemory>(getInMemoryMetadata().getColumns().getAllPhysical().toString()));
|
|
|
|
|
|
|
|
/// count.txt
|
|
|
|
size_t num_rows = 0;
|
|
|
|
for (const auto & block : indices.blocks)
|
|
|
|
num_rows += block.num_rows;
|
|
|
|
backup_entries.emplace_back("count.txt", std::make_unique<BackupEntryFromMemory>(toString(num_rows)));
|
|
|
|
|
|
|
|
return backup_entries;
|
|
|
|
}
|
|
|
|
|
2022-01-19 18:56:08 +00:00
|
|
|
class StripeLogRestoreTask : public IRestoreTask
|
2021-10-26 09:48:31 +00:00
|
|
|
{
|
2022-01-17 18:55:40 +00:00
|
|
|
using WriteLock = StorageStripeLog::WriteLock;
|
2021-10-26 09:48:31 +00:00
|
|
|
|
2022-01-17 18:55:40 +00:00
|
|
|
public:
|
|
|
|
StripeLogRestoreTask(
|
|
|
|
const std::shared_ptr<StorageStripeLog> storage_,
|
|
|
|
const BackupPtr & backup_,
|
|
|
|
const String & data_path_in_backup_,
|
|
|
|
ContextMutablePtr context_)
|
|
|
|
: storage(storage_), backup(backup_), data_path_in_backup(data_path_in_backup_), context(context_)
|
2021-10-26 09:48:31 +00:00
|
|
|
{
|
2022-01-17 18:55:40 +00:00
|
|
|
}
|
|
|
|
|
2022-01-19 18:56:08 +00:00
|
|
|
RestoreTasks run() override
|
2022-01-17 18:55:40 +00:00
|
|
|
{
|
|
|
|
WriteLock lock{storage->rwlock, getLockTimeout(context)};
|
2021-10-26 09:48:31 +00:00
|
|
|
if (!lock)
|
|
|
|
throw Exception("Lock timeout exceeded", ErrorCodes::TIMEOUT_EXCEEDED);
|
|
|
|
|
2022-01-17 18:55:40 +00:00
|
|
|
auto & file_checker = storage->file_checker;
|
|
|
|
|
2021-10-26 09:48:31 +00:00
|
|
|
/// Load the indices if not loaded yet. We have to do that now because we're going to update these indices.
|
2022-01-17 18:55:40 +00:00
|
|
|
storage->loadIndices(lock);
|
2021-10-26 09:48:31 +00:00
|
|
|
|
|
|
|
/// If there were no files, save zero file sizes to be able to rollback in case of error.
|
2022-01-17 18:55:40 +00:00
|
|
|
storage->saveFileSizes(lock);
|
2021-10-26 09:48:31 +00:00
|
|
|
|
|
|
|
try
|
|
|
|
{
|
|
|
|
/// Append the data file.
|
2022-01-17 18:55:40 +00:00
|
|
|
auto old_data_size = file_checker.getFileSize(storage->data_file_path);
|
2021-10-26 09:48:31 +00:00
|
|
|
{
|
2022-01-17 18:55:40 +00:00
|
|
|
const auto & data_file_path = storage->data_file_path;
|
2021-10-26 09:48:31 +00:00
|
|
|
String file_path_in_backup = data_path_in_backup + fileName(data_file_path);
|
2021-11-06 14:32:48 +00:00
|
|
|
auto backup_entry = backup->readFile(file_path_in_backup);
|
2022-01-17 18:55:40 +00:00
|
|
|
const auto & disk = storage->disk;
|
2021-10-26 09:48:31 +00:00
|
|
|
auto in = backup_entry->getReadBuffer();
|
2022-01-17 18:55:40 +00:00
|
|
|
auto out = disk->writeFile(data_file_path, storage->max_compress_block_size, WriteMode::Append);
|
2021-10-26 09:48:31 +00:00
|
|
|
copyData(*in, *out);
|
|
|
|
}
|
|
|
|
|
|
|
|
/// Append the index.
|
|
|
|
{
|
2022-01-17 18:55:40 +00:00
|
|
|
const auto & index_file_path = storage->index_file_path;
|
2021-11-06 14:32:48 +00:00
|
|
|
String index_path_in_backup = data_path_in_backup + fileName(index_file_path);
|
2021-10-26 09:48:31 +00:00
|
|
|
IndexForNativeFormat extra_indices;
|
2021-11-06 14:32:48 +00:00
|
|
|
auto backup_entry = backup->readFile(index_path_in_backup);
|
2021-10-26 09:48:31 +00:00
|
|
|
auto index_in = backup_entry->getReadBuffer();
|
|
|
|
CompressedReadBuffer index_compressed_in{*index_in};
|
|
|
|
extra_indices.read(index_compressed_in);
|
|
|
|
|
|
|
|
/// Adjust the offsets.
|
|
|
|
for (auto & block : extra_indices.blocks)
|
|
|
|
{
|
|
|
|
for (auto & column : block.columns)
|
|
|
|
column.location.offset_in_compressed_file += old_data_size;
|
|
|
|
}
|
|
|
|
|
2022-01-17 18:55:40 +00:00
|
|
|
insertAtEnd(storage->indices.blocks, std::move(extra_indices.blocks));
|
2021-10-26 09:48:31 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
/// Finish writing.
|
2022-01-17 18:55:40 +00:00
|
|
|
storage->saveIndices(lock);
|
|
|
|
storage->saveFileSizes(lock);
|
|
|
|
return {};
|
2021-10-26 09:48:31 +00:00
|
|
|
}
|
|
|
|
catch (...)
|
|
|
|
{
|
|
|
|
/// Rollback partial writes.
|
|
|
|
file_checker.repair();
|
2022-01-17 18:55:40 +00:00
|
|
|
storage->removeUnsavedIndices(lock);
|
2021-10-26 09:48:31 +00:00
|
|
|
throw;
|
|
|
|
}
|
2022-01-17 18:55:40 +00:00
|
|
|
}
|
2021-10-26 09:48:31 +00:00
|
|
|
|
2022-01-17 18:55:40 +00:00
|
|
|
private:
|
|
|
|
std::shared_ptr<StorageStripeLog> storage;
|
|
|
|
BackupPtr backup;
|
|
|
|
String data_path_in_backup;
|
|
|
|
ContextMutablePtr context;
|
|
|
|
};
|
|
|
|
|
|
|
|
|
2022-04-19 18:15:27 +00:00
|
|
|
RestoreTaskPtr StorageStripeLog::restoreData(ContextMutablePtr context, const ASTs & partitions, const BackupPtr & backup, const String & data_path_in_backup, const StorageRestoreSettings &, const std::shared_ptr<IRestoreCoordination> &)
|
2022-01-17 18:55:40 +00:00
|
|
|
{
|
|
|
|
if (!partitions.empty())
|
|
|
|
throw Exception(ErrorCodes::NOT_IMPLEMENTED, "Table engine {} doesn't support partitions", getName());
|
|
|
|
|
|
|
|
return std::make_unique<StripeLogRestoreTask>(
|
|
|
|
typeid_cast<std::shared_ptr<StorageStripeLog>>(shared_from_this()), backup, data_path_in_backup, context);
|
2018-04-21 00:35:20 +00:00
|
|
|
}
|
|
|
|
|
2017-12-30 00:36:06 +00:00
|
|
|
|
|
|
|
void registerStorageStripeLog(StorageFactory & factory)
|
|
|
|
{
|
2020-02-18 14:41:30 +00:00
|
|
|
StorageFactory::StorageFeatures features{
|
|
|
|
.supports_settings = true
|
|
|
|
};
|
|
|
|
|
2017-12-30 00:36:06 +00:00
|
|
|
factory.registerStorage("StripeLog", [](const StorageFactory::Arguments & args)
|
|
|
|
{
|
2020-02-18 14:41:30 +00:00
|
|
|
if (!args.engine_args.empty())
|
2017-12-30 00:36:06 +00:00
|
|
|
throw Exception(
|
2020-02-18 14:41:30 +00:00
|
|
|
"Engine " + args.engine_name + " doesn't support any arguments (" + toString(args.engine_args.size()) + " given)",
|
2017-12-30 00:36:06 +00:00
|
|
|
ErrorCodes::NUMBER_OF_ARGUMENTS_DOESNT_MATCH);
|
|
|
|
|
2020-02-18 14:41:30 +00:00
|
|
|
String disk_name = getDiskName(*args.storage_def);
|
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<StorageStripeLog>(
|
2021-04-23 12:18:23 +00:00
|
|
|
disk,
|
|
|
|
args.relative_data_path,
|
|
|
|
args.table_id,
|
|
|
|
args.columns,
|
|
|
|
args.constraints,
|
|
|
|
args.comment,
|
|
|
|
args.attach,
|
|
|
|
args.getContext()->getSettings().max_compress_block_size);
|
2020-02-18 14:41:30 +00:00
|
|
|
}, features);
|
2017-12-30 00:36:06 +00:00
|
|
|
}
|
|
|
|
|
2015-08-16 07:01:41 +00:00
|
|
|
}
|