ClickHouse/dbms/include/DB/Client/Connection.h

279 lines
9.9 KiB
C
Raw Normal View History

2012-05-21 20:38:34 +00:00
#pragma once
2015-09-29 19:19:54 +00:00
#include <common/logger_useful.h>
2012-10-16 19:20:58 +00:00
2012-05-16 18:03:00 +00:00
#include <Poco/Net/StreamSocket.h>
#include <DB/Common/Throttler.h>
2012-05-16 18:03:00 +00:00
#include <DB/Core/Block.h>
#include <DB/Core/Defines.h>
2012-05-16 18:03:00 +00:00
#include <DB/Core/Progress.h>
#include <DB/Core/Protocol.h>
#include <DB/Core/QueryProcessingStage.h>
#include <DB/DataStreams/IBlockInputStream.h>
#include <DB/DataStreams/IBlockOutputStream.h>
2015-01-18 08:25:56 +00:00
#include <DB/DataStreams/BlockStreamProfileInfo.h>
2012-05-16 18:03:00 +00:00
#include <DB/Interpreters/Settings.h>
#include <atomic>
2012-05-16 18:03:00 +00:00
namespace DB
{
using Poco::SharedPtr;
class ParallelReplicas;
/// Поток блоков читающих из таблицы и ее имя
typedef std::pair<BlockInputStreamPtr, std::string> ExternalTableData;
/// Вектор пар, описывающих таблицы
typedef std::vector<ExternalTableData> ExternalTablesData;
2012-05-16 18:03:00 +00:00
2012-05-16 18:03:00 +00:00
/** Соединение с сервером БД для использования в клиенте.
* Как использовать - см. Core/Protocol.h
* (Реализацию на стороне сервера - см. Server/TCPHandler.h)
2012-05-30 06:46:57 +00:00
*
* В качестве default_database может быть указана пустая строка
* - в этом случае сервер использует свою БД по-умолчанию.
2012-05-16 18:03:00 +00:00
*/
2012-10-22 19:03:01 +00:00
class Connection : private boost::noncopyable
2012-05-16 18:03:00 +00:00
{
friend class ParallelReplicas;
2012-05-16 18:03:00 +00:00
public:
2012-05-30 06:46:57 +00:00
Connection(const String & host_, UInt16 port_, const String & default_database_,
2013-08-10 09:04:45 +00:00
const String & user_, const String & password_,
2012-05-23 19:51:30 +00:00
const String & client_name_ = "client",
2012-07-26 19:42:20 +00:00
Protocol::Compression::Enum compression_ = Protocol::Compression::Enable,
Poco::Timespan connect_timeout_ = Poco::Timespan(DBMS_DEFAULT_CONNECT_TIMEOUT_SEC, 0),
Poco::Timespan receive_timeout_ = Poco::Timespan(DBMS_DEFAULT_RECEIVE_TIMEOUT_SEC, 0),
Poco::Timespan send_timeout_ = Poco::Timespan(DBMS_DEFAULT_SEND_TIMEOUT_SEC, 0),
Poco::Timespan ping_timeout_ = Poco::Timespan(DBMS_DEFAULT_PING_TIMEOUT_SEC, 0))
2013-08-10 09:04:45 +00:00
:
host(host_), port(port_), default_database(default_database_),
2015-05-28 21:41:28 +00:00
user(user_), password(password_), resolved_address(host, port),
client_name(client_name_),
compression(compression_),
connect_timeout(connect_timeout_), receive_timeout(receive_timeout_), send_timeout(send_timeout_),
ping_timeout(ping_timeout_),
log_wrapper(*this)
{
/// Соединеняемся не сразу, а при первой необходимости.
if (user.empty())
user = "default";
setDescription();
2015-05-28 21:41:28 +00:00
}
Connection(const String & host_, UInt16 port_, const Poco::Net::SocketAddress & resolved_address_,
const String & default_database_,
const String & user_, const String & password_,
const String & client_name_ = "client",
Protocol::Compression::Enum compression_ = Protocol::Compression::Enable,
Poco::Timespan connect_timeout_ = Poco::Timespan(DBMS_DEFAULT_CONNECT_TIMEOUT_SEC, 0),
Poco::Timespan receive_timeout_ = Poco::Timespan(DBMS_DEFAULT_RECEIVE_TIMEOUT_SEC, 0),
Poco::Timespan send_timeout_ = Poco::Timespan(DBMS_DEFAULT_SEND_TIMEOUT_SEC, 0),
Poco::Timespan ping_timeout_ = Poco::Timespan(DBMS_DEFAULT_PING_TIMEOUT_SEC, 0))
:
host(host_), port(port_),
default_database(default_database_),
2013-08-10 09:04:45 +00:00
user(user_), password(password_),
2015-05-28 21:41:28 +00:00
resolved_address(resolved_address_),
2015-02-01 07:21:19 +00:00
client_name(client_name_),
compression(compression_),
2012-10-16 19:20:58 +00:00
connect_timeout(connect_timeout_), receive_timeout(receive_timeout_), send_timeout(send_timeout_),
ping_timeout(ping_timeout_),
2015-05-28 21:41:28 +00:00
log_wrapper(*this)
2012-05-16 18:03:00 +00:00
{
2012-05-21 20:38:34 +00:00
/// Соединеняемся не сразу, а при первой необходимости.
2013-08-10 09:04:45 +00:00
if (user.empty())
user = "default";
setDescription();
2012-05-16 18:03:00 +00:00
}
virtual ~Connection() {};
/// Установить ограничитель сетевого трафика. Один ограничитель может использоваться одновременно для нескольких разных соединений.
void setThrottler(const ThrottlerPtr & throttler_)
{
throttler = throttler_;
}
2012-05-16 18:03:00 +00:00
/// Пакет, который может быть получен от сервера.
struct Packet
{
UInt64 type;
Block block;
SharedPtr<Exception> exception;
Progress progress;
2013-05-22 14:57:43 +00:00
BlockStreamProfileInfo profile_info;
2012-05-16 18:03:00 +00:00
Packet() : type(Protocol::Server::Hello) {}
};
/// Изменить базу данных по умолчанию. Изменения начинают использоваться только при следующем переподключении.
void setDefaultDatabase(const String & database);
2012-05-16 18:20:45 +00:00
void getServerVersion(String & name, UInt64 & version_major, UInt64 & version_minor, UInt64 & revision);
2012-05-16 18:03:00 +00:00
/// Для сообщений в логе и в эксепшенах.
const String & getDescription() const
{
return description;
}
2015-05-28 21:41:28 +00:00
const String & getHost() const
{
return host;
}
UInt16 getPort() const
{
return port;
}
/// Если последний флаг true, то затем необходимо вызвать sendExternalTablesData
void sendQuery(const String & query, const String & query_id_ = "", UInt64 stage = QueryProcessingStage::Complete,
2014-04-08 07:31:51 +00:00
const Settings * settings = nullptr, bool with_pending_data = false);
2012-05-16 18:03:00 +00:00
void sendCancel();
/// Отправить блок данных, на сервере сохранить во временную таблицу name
void sendData(const Block & block, const String & name = "");
/// Отправить все содержимое внешних таблиц
void sendExternalTablesData(ExternalTablesData & data);
2012-05-16 18:03:00 +00:00
/// Отправить блок данных, который уже был заранее сериализован (и, если надо, сжат), который следует прочитать из input-а.
/// можно передать размер сериализованного/сжатого блока.
void sendPreparedData(ReadBuffer & input, size_t size, const String & name = "");
/// Проверить, есть ли данные, которые можно прочитать.
2012-05-16 18:03:00 +00:00
bool poll(size_t timeout_microseconds = 0);
/// Проверить, есть ли данные в буфере для чтения.
bool hasReadBufferPendingData() const;
2012-05-16 18:03:00 +00:00
/// Получить пакет от сервера.
Packet receivePacket();
/// Если ещё не соединено, или соединение разорвано - соединиться. Если не получилось - кинуть исключение.
void forceConnected();
/** Разорвать соединение.
* Это может быть нужно, например, чтобы соединение не осталось висеть в
* рассинхронизированном состоянии (когда кто-то чего-то продолжает ждать) после эксепшена.
*/
void disconnect();
2015-10-12 14:53:16 +00:00
/** Заполнить информацию, которая необходима при получении блока для некоторых задач
* (пока только для запроса DESCRIBE TABLE с Distributed-таблицами).
*/
void fillBlockExtraInfo(BlockExtraInfo & info) const;
size_t outBytesCount() const { return !out.isNull() ? out->count() : 0; }
size_t inBytesCount() const { return !in.isNull() ? in->count() : 0; }
2012-05-16 18:03:00 +00:00
private:
String host;
UInt16 port;
2012-05-30 06:46:57 +00:00
String default_database;
2013-08-10 09:04:45 +00:00
String user;
String password;
2012-05-16 18:20:45 +00:00
/** Адрес может быть заранее отрезолвен и передан в конструктор. Тогда поля host и port имеют смысл только для логгирования.
* Иначе адрес резолвится в конструкторе. То есть, DNS балансировка не поддерживается.
*/
2015-05-28 21:41:28 +00:00
Poco::Net::SocketAddress resolved_address;
/// Для сообщений в логе и в эксепшенах.
String description;
void setDescription();
2012-05-23 19:51:30 +00:00
String client_name;
2015-02-01 07:21:19 +00:00
bool connected = false;
2012-05-21 20:38:34 +00:00
2012-05-16 18:20:45 +00:00
String server_name;
2015-02-01 07:21:19 +00:00
UInt64 server_version_major = 0;
UInt64 server_version_minor = 0;
UInt64 server_revision = 0;
2012-05-16 18:03:00 +00:00
Poco::Net::StreamSocket socket;
SharedPtr<ReadBuffer> in;
SharedPtr<WriteBuffer> out;
2012-05-16 18:03:00 +00:00
String query_id;
2012-05-16 18:03:00 +00:00
UInt64 compression; /// Сжимать ли данные при взаимодействии с сервером.
/// каким алгоритмом сжимать данные при INSERT и данные внешних таблиц
CompressionMethod network_compression_method = CompressionMethod::LZ4;
2012-05-16 18:03:00 +00:00
/** Если не nullptr, то используется, чтобы ограничить сетевой трафик.
* Учитывается только трафик при передаче блоков. Другие пакеты не учитываются.
*/
ThrottlerPtr throttler;
2012-07-26 19:42:20 +00:00
Poco::Timespan connect_timeout;
Poco::Timespan receive_timeout;
Poco::Timespan send_timeout;
Poco::Timespan ping_timeout;
2012-07-26 19:42:20 +00:00
2012-05-16 18:03:00 +00:00
/// Откуда читать результат выполнения запроса.
SharedPtr<ReadBuffer> maybe_compressed_in;
BlockInputStreamPtr block_in;
/// Куда писать данные INSERT-а.
SharedPtr<WriteBuffer> maybe_compressed_out;
BlockOutputStreamPtr block_out;
/// логгер, создаваемый лениво, чтобы не обращаться к DNS в конструкторе
class LoggerWrapper
{
public:
2015-05-28 21:41:28 +00:00
LoggerWrapper(Connection & parent_)
: log(nullptr), parent(parent_)
{
}
Logger * get()
{
if (!log)
log = &Logger::get("Connection (" + parent.getDescription() + ")");
return log;
}
private:
std::atomic<Logger *> log;
2015-05-28 21:41:28 +00:00
Connection & parent;
};
LoggerWrapper log_wrapper;
2012-10-16 19:20:58 +00:00
2012-05-16 18:03:00 +00:00
void connect();
void sendHello();
void receiveHello();
bool ping();
Block receiveData();
SharedPtr<Exception> receiveException();
Progress receiveProgress();
2013-05-22 14:57:43 +00:00
BlockStreamProfileInfo receiveProfileInfo();
void initBlockInput();
2012-05-16 18:03:00 +00:00
};
2012-05-21 20:38:34 +00:00
typedef SharedPtr<Connection> ConnectionPtr;
typedef std::vector<ConnectionPtr> Connections;
2012-05-16 18:03:00 +00:00
}