2015-01-14 10:06:30 +00:00
|
|
|
|
#pragma once
|
|
|
|
|
|
2015-02-10 20:48:17 +00:00
|
|
|
|
#include <DB/Common/Throttler.h>
|
2015-01-14 10:06:30 +00:00
|
|
|
|
#include <DB/Client/Connection.h>
|
2015-01-15 18:33:20 +00:00
|
|
|
|
#include <DB/Client/ConnectionPool.h>
|
2015-01-14 10:06:30 +00:00
|
|
|
|
|
2015-02-10 20:48:17 +00:00
|
|
|
|
|
2015-01-14 10:06:30 +00:00
|
|
|
|
namespace DB
|
|
|
|
|
{
|
2015-02-06 22:32:54 +00:00
|
|
|
|
|
2015-01-14 10:06:30 +00:00
|
|
|
|
|
2015-02-10 20:48:17 +00:00
|
|
|
|
/** Для получения данных сразу из нескольких реплик (соединений) в рамках одного потока.
|
|
|
|
|
* В качестве вырожденного случая, может также работать с одним соединением.
|
|
|
|
|
*
|
|
|
|
|
* Интерфейс почти совпадает с Connection.
|
|
|
|
|
*/
|
|
|
|
|
class ParallelReplicas final : private boost::noncopyable
|
|
|
|
|
{
|
|
|
|
|
public:
|
|
|
|
|
/// Принимает готовое соединение.
|
|
|
|
|
ParallelReplicas(Connection * connection_, Settings * settings_, ThrottlerPtr throttler_);
|
|
|
|
|
|
|
|
|
|
/// Принимает пул, из которого нужно будет достать одно или несколько соединений.
|
|
|
|
|
ParallelReplicas(IConnectionPool * pool_, Settings * settings_, ThrottlerPtr throttler_);
|
|
|
|
|
|
|
|
|
|
/// Отправить на реплики всё содержимое внешних таблиц.
|
|
|
|
|
void sendExternalTablesData(std::vector<ExternalTablesData> & data);
|
2015-01-14 10:06:30 +00:00
|
|
|
|
|
2015-02-10 20:48:17 +00:00
|
|
|
|
/// Отправить запрос на реплики.
|
|
|
|
|
void sendQuery(const String & query, const String & query_id = "",
|
|
|
|
|
UInt64 stage = QueryProcessingStage::Complete, bool with_pending_data = false);
|
|
|
|
|
|
|
|
|
|
/// Получить пакет от какой-нибудь реплики.
|
|
|
|
|
Connection::Packet receivePacket();
|
|
|
|
|
|
|
|
|
|
/// Отменить запросы к репликам
|
|
|
|
|
void sendCancel();
|
|
|
|
|
|
|
|
|
|
/** На каждой реплике читать и пропускать все пакеты до EndOfStream или Exception.
|
|
|
|
|
* Возвращает EndOfStream, если не было получено никакого исключения. В противном
|
|
|
|
|
* случае возвращает последний полученный пакет типа Exception.
|
|
|
|
|
*/
|
|
|
|
|
Connection::Packet drain();
|
2015-01-14 10:06:30 +00:00
|
|
|
|
|
2015-02-10 20:48:17 +00:00
|
|
|
|
/// Получить адреса реплик в виде строки.
|
|
|
|
|
std::string dumpAddresses() const;
|
2015-01-14 10:06:30 +00:00
|
|
|
|
|
2015-02-10 20:48:17 +00:00
|
|
|
|
/// Возвращает количесто реплик.
|
|
|
|
|
size_t size() const { return replica_map.size(); }
|
2015-01-22 21:54:16 +00:00
|
|
|
|
|
2015-02-10 20:48:17 +00:00
|
|
|
|
/// Проверить, есть ли действительные реплики.
|
|
|
|
|
bool hasActiveReplicas() const { return active_replica_count > 0; }
|
2015-01-15 12:09:26 +00:00
|
|
|
|
|
2015-02-10 20:48:17 +00:00
|
|
|
|
private:
|
|
|
|
|
/// Реплики хэшированные по id сокета
|
|
|
|
|
using ReplicaMap = std::unordered_map<int, Connection *>;
|
2015-01-15 12:09:26 +00:00
|
|
|
|
|
|
|
|
|
|
2015-02-10 20:48:17 +00:00
|
|
|
|
/// Зарегистрировать реплику.
|
|
|
|
|
void registerReplica(Connection * connection);
|
2015-01-15 15:05:03 +00:00
|
|
|
|
|
2015-02-10 20:48:17 +00:00
|
|
|
|
/// Получить реплику, на которой можно прочитать данные.
|
|
|
|
|
ReplicaMap::iterator getReplicaForReading();
|
2015-01-15 16:55:50 +00:00
|
|
|
|
|
2015-02-10 20:48:17 +00:00
|
|
|
|
/** Проверить, есть ли данные, которые можно прочитать на каких-нибудь репликах.
|
|
|
|
|
* Возвращает одну такую реплику, если она найдётся.
|
|
|
|
|
*/
|
|
|
|
|
ReplicaMap::iterator waitForReadEvent();
|
2015-01-14 10:06:30 +00:00
|
|
|
|
|
2015-02-10 20:48:17 +00:00
|
|
|
|
/// Пометить реплику как недействительную.
|
|
|
|
|
void invalidateReplica(ReplicaMap::iterator it);
|
2015-02-05 22:31:03 +00:00
|
|
|
|
|
|
|
|
|
|
2015-02-10 20:48:17 +00:00
|
|
|
|
Settings * settings;
|
|
|
|
|
ReplicaMap replica_map;
|
2015-02-04 13:07:10 +00:00
|
|
|
|
|
2015-02-10 20:48:17 +00:00
|
|
|
|
/// Если не nullptr, то используется, чтобы ограничить сетевой трафик.
|
|
|
|
|
ThrottlerPtr throttler;
|
2015-02-06 10:41:03 +00:00
|
|
|
|
|
2015-02-10 20:48:17 +00:00
|
|
|
|
std::vector<ConnectionPool::Entry> pool_entries;
|
|
|
|
|
ConnectionPool::Entry pool_entry;
|
2015-02-06 14:46:04 +00:00
|
|
|
|
|
2015-02-10 20:48:17 +00:00
|
|
|
|
/// Текущее количество действительных соединений к репликам.
|
|
|
|
|
size_t active_replica_count;
|
|
|
|
|
/// Запрос выполняется параллельно на нескольких репликах.
|
|
|
|
|
bool supports_parallel_execution;
|
|
|
|
|
/// Отправили запрос
|
|
|
|
|
bool sent_query = false;
|
|
|
|
|
/// Отменили запрос
|
|
|
|
|
bool cancelled = false;
|
|
|
|
|
};
|
2015-02-06 14:46:04 +00:00
|
|
|
|
|
2015-01-14 10:06:30 +00:00
|
|
|
|
}
|