#pragma once #include #include #include #include #include namespace DB { /** Распределённая таблица, находящаяся на нескольких серверах. * Использует данные заданной БД и таблицы на каждом сервере. * * Можно передать один адрес, а не несколько. * В этом случае, таблицу можно считать удалённой, а не распределённой. */ class StorageDistributed : public IStorage { friend class DistributedBlockOutputStream; friend class DirectoryMonitor; public: static StoragePtr create( const std::string & name_, /// Имя таблицы. NamesAndTypesListPtr columns_, /// Список столбцов. const String & remote_database_, /// БД на удалённых серверах. const String & remote_table_, /// Имя таблицы на удалённых серверах. const String & cluster_name, Context & context_, const ASTPtr & sharding_key_, const String & data_path_); static StoragePtr create( const std::string & name_, /// Имя таблицы. NamesAndTypesListPtr columns_, /// Список столбцов. const String & remote_database_, /// БД на удалённых серверах. const String & remote_table_, /// Имя таблицы на удалённых серверах. SharedPtr & owned_cluster_, Context & context_); std::string getName() const override { return "Distributed"; } std::string getTableName() const override { return name; } bool supportsSampling() const override { return true; } bool supportsFinal() const override { return true; } bool supportsPrewhere() const override { return true; } const NamesAndTypesList & getColumnsList() const override { return *columns; } NameAndTypePair getColumn(const String & column_name) const override; bool hasColumn(const String & column_name) const override; bool isRemote() const override { return true; } /// Сохранить временные таблицы, чтобы при следующем вызове метода read переслать их на удаленные серверы. void storeExternalTables(const Tables & tables_) override { external_tables = tables_; } BlockInputStreams read( const Names & column_names, ASTPtr query, const Settings & settings, QueryProcessingStage::Enum & processed_stage, size_t max_block_size = DEFAULT_BLOCK_SIZE, unsigned threads = 1) override; BlockOutputStreamPtr write(ASTPtr query) override; void drop() override {} void rename(const String & new_path_to_db, const String & new_database_name, const String & new_table_name) override { name = new_table_name; } /// в подтаблицах добавлять и удалять столбы нужно вручную /// структура подтаблиц не проверяется void alter(const AlterCommands & params, const String & database_name, const String & table_name, Context & context) override; void shutdown() override; const ExpressionActionsPtr & getShardingKeyExpr() const { return sharding_key_expr; } const String & getShardingKeyColumnName() const { return sharding_key_column_name; } const String & getPath() const { return path; } private: StorageDistributed( const std::string & name_, NamesAndTypesListPtr columns_, const String & remote_database_, const String & remote_table_, Cluster & cluster_, Context & context_, const ASTPtr & sharding_key_ = nullptr, const String & data_path_ = String{}); /// create directory monitor thread by subdirectory name void createDirectoryMonitor(const std::string & name); /// create directory monitors for each existing subdirectory void createDirectoryMonitors(); /// ensure directory monitor creation void requireDirectoryMonitor(const std::string & name); String name; NamesAndTypesListPtr columns; String remote_database; String remote_table; Context & context; /// Временные таблицы, которые необходимо отправить на сервер. Переменная очищается после каждого вызова метода read /// Для подготовки к отправке нужно использовтаь метод storeExternalTables Tables external_tables; /// Используется только, если таблица должна владеть объектом Cluster, которым больше никто не владеет - для реализации TableFunctionRemote. SharedPtr owned_cluster; /// Соединения с удалёнными серверами. Cluster & cluster; ExpressionActionsPtr sharding_key_expr; String sharding_key_column_name; bool write_enabled; String path; class DirectoryMonitor; std::unordered_map> directory_monitors; }; }