blob: 1dd8eda157507b1af52512e77403f1afc5f0a907 (
plain) (
blame)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
|
#pragma once
#include <base/types.h>
#include <Core/MySQL/MySQLReplication.h>
#include <IO/ReadBufferFromPocoSocket.h>
#include <IO/ReadHelpers.h>
#include <IO/WriteBufferFromPocoSocket.h>
#include <IO/WriteHelpers.h>
#include <Poco/Net/NetException.h>
#include <Poco/Net/StreamSocket.h>
#include <Common/Exception.h>
#include <Common/NetException.h>
#include <Core/MySQL/IMySQLWritePacket.h>
namespace DB
{
using namespace MySQLProtocol;
using namespace MySQLReplication;
class MySQLClient
{
public:
MySQLClient(const String & host_, UInt16 port_, const String & user_, const String & password_);
MySQLClient(MySQLClient && other) noexcept;
void connect();
void disconnect();
void ping();
void setBinlogChecksum(const String & binlog_checksum);
/// Start replication stream by GTID.
/// replicate_db: replication database schema, events from other databases will be ignored.
/// gtid: executed gtid sets format like 'hhhhhhhh-hhhh-hhhh-hhhh-hhhhhhhhhhhh:x-y'.
void startBinlogDumpGTID(UInt32 slave_id, String replicate_db, std::unordered_set<String> replicate_tables, String gtid, const String & binlog_checksum);
BinlogEventPtr readOneBinlogEvent(UInt64 milliseconds = 0);
Position getPosition() const { return replication.getPosition(); }
private:
String host;
UInt16 port;
String user;
String password;
bool connected = false;
uint8_t sequence_id = 0;
uint32_t client_capabilities = 0;
const UInt8 charset_utf8 = 33;
const String mysql_native_password = "mysql_native_password";
MySQLFlavor replication;
std::shared_ptr<ReadBuffer> in;
std::shared_ptr<WriteBuffer> out;
std::unique_ptr<Poco::Net::StreamSocket> socket;
std::optional<Poco::Net::SocketAddress> address;
MySQLProtocol::PacketEndpointPtr packet_endpoint;
void handshake();
void registerSlaveOnMaster(UInt32 slave_id);
void writeCommand(char command, String query);
};
class WriteCommand : public IMySQLWritePacket
{
public:
char command;
String query;
WriteCommand(char command_, String query_) : command(command_), query(query_) { }
size_t getPayloadSize() const override { return 1 + query.size(); }
void writePayloadImpl(WriteBuffer & buffer) const override
{
buffer.write(static_cast<char>(command));
if (!query.empty())
{
buffer.write(query.data(), query.size());
}
}
};
}
|