2016-01-28 01:00:27 +00:00
|
|
|
#pragma once
|
|
|
|
|
|
|
|
#include <DB/Interpreters/InterserverIOHandler.h>
|
2016-03-01 17:47:53 +00:00
|
|
|
#include <DB/Storages/MergeTree/MergeTreeData.h>
|
2016-01-28 01:00:27 +00:00
|
|
|
#include <DB/IO/WriteBuffer.h>
|
2016-01-28 01:00:42 +00:00
|
|
|
#include <common/logger_useful.h>
|
2016-03-01 17:47:53 +00:00
|
|
|
#include <functional>
|
2016-01-28 01:00:27 +00:00
|
|
|
|
|
|
|
namespace DB
|
|
|
|
{
|
|
|
|
|
|
|
|
class StorageReplicatedMergeTree;
|
|
|
|
|
2016-03-01 17:47:53 +00:00
|
|
|
namespace ShardedPartitionUploader
|
2016-01-28 01:00:27 +00:00
|
|
|
{
|
|
|
|
|
|
|
|
/** Сервис для получения кусков из партиции таблицы *MergeTree.
|
|
|
|
*/
|
|
|
|
class Service final : public InterserverIOEndpoint
|
|
|
|
{
|
|
|
|
public:
|
2016-03-03 04:30:36 +00:00
|
|
|
Service(StoragePtr & storage_);
|
2016-01-28 01:00:27 +00:00
|
|
|
Service(const Service &) = delete;
|
|
|
|
Service & operator=(const Service &) = delete;
|
|
|
|
std::string getId(const std::string & node_id) const override;
|
2016-03-01 17:47:53 +00:00
|
|
|
void processQuery(const Poco::Net::HTMLForm & params, ReadBuffer & body, WriteBuffer & out) override;
|
2016-01-28 01:00:27 +00:00
|
|
|
|
|
|
|
private:
|
2016-03-03 04:30:36 +00:00
|
|
|
StoragePtr owned_storage;
|
2016-03-01 17:47:53 +00:00
|
|
|
MergeTreeData & data;
|
2016-03-25 11:48:45 +00:00
|
|
|
Logger * log = &Logger::get("ShardedPartitionUploader::Service");
|
2016-01-28 01:00:27 +00:00
|
|
|
};
|
|
|
|
|
|
|
|
/** Клиент для отправления кусков из партиции таблицы *MergeTree.
|
|
|
|
*/
|
|
|
|
class Client final
|
|
|
|
{
|
|
|
|
public:
|
2016-03-01 17:47:53 +00:00
|
|
|
using CancellationHook = std::function<void()>;
|
|
|
|
|
|
|
|
public:
|
|
|
|
Client(StorageReplicatedMergeTree & storage_);
|
|
|
|
|
2016-01-28 01:00:27 +00:00
|
|
|
Client(const Client &) = delete;
|
|
|
|
Client & operator=(const Client &) = delete;
|
2016-03-01 17:47:53 +00:00
|
|
|
|
|
|
|
void setCancellationHook(CancellationHook cancellation_hook_);
|
|
|
|
|
|
|
|
bool send(const std::string & part_name, size_t shard_no,
|
|
|
|
const InterserverIOEndpointLocation & to_location);
|
|
|
|
|
2016-01-28 01:00:27 +00:00
|
|
|
void cancel() { is_cancelled = true; }
|
|
|
|
|
|
|
|
private:
|
2016-03-01 17:47:53 +00:00
|
|
|
MergeTreeData::DataPartPtr findShardedPart(const std::string & name, size_t shard_no);
|
|
|
|
void abortIfRequested();
|
|
|
|
|
|
|
|
private:
|
|
|
|
StorageReplicatedMergeTree & storage;
|
|
|
|
MergeTreeData & data;
|
|
|
|
CancellationHook cancellation_hook;
|
2016-01-28 01:00:27 +00:00
|
|
|
std::atomic<bool> is_cancelled{false};
|
2016-03-25 11:48:45 +00:00
|
|
|
Logger * log = &Logger::get("ShardedPartitionUploader::Client");
|
2016-01-28 01:00:27 +00:00
|
|
|
};
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
}
|