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());
        }
    }
};
}
 |