This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new b87c1778607 branch-4.1: [fix](streamingjob) Prevent MySQL CDC data
loss on keepalive reconnect for non-GTID #66998 (#67269)
b87c1778607 is described below
commit b87c1778607077a9ca1ee6958b4999efc39c3e6b
Author: github-actions[bot]
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Fri Aug 28 21:11:05 2026 +0800
branch-4.1: [fix](streamingjob) Prevent MySQL CDC data loss on keepalive
reconnect for non-GTID #66998 (#67269)
Cherry-picked from #66998
Co-authored-by: wudi <[email protected]>
---
.licenserc.yaml | 1 +
fs_brokers/cdc_client/pom.xml | 3 +-
.../shyiko/mysql/binlog/BinaryLogClient.java | 1477 ++++++++++++++++++++
.../mysql/MySqlStreamingChangeEventSource.java | 17 +-
.../apache/doris/cdcclient/common/Constants.java | 2 +-
.../source/reader/mysql/MySqlSourceReader.java | 5 +-
.../BinaryLogClientTransactionReplayTest.java | 135 ++
.../MySqlBinaryLogClientKeepAliveITCase.java | 208 +++
8 files changed, 1836 insertions(+), 12 deletions(-)
diff --git a/.licenserc.yaml b/.licenserc.yaml
index 9f031b1ff0c..bcf4fca020a 100644
--- a/.licenserc.yaml
+++ b/.licenserc.yaml
@@ -114,4 +114,5 @@ header:
- "tools/FlameGraph/*"
- "thirdparty/LICENSE.txt"
- "fs_brokers/cdc_client/src/main/java/io/debezium/**"
+ - "fs_brokers/cdc_client/src/main/java/com/github/shyiko/**"
comment: on-failure
diff --git a/fs_brokers/cdc_client/pom.xml b/fs_brokers/cdc_client/pom.xml
index 0a0c7355318..9d0ad7c8087 100644
--- a/fs_brokers/cdc_client/pom.xml
+++ b/fs_brokers/cdc_client/pom.xml
@@ -276,6 +276,7 @@ under the License.
<artifactId>maven-failsafe-plugin</artifactId>
<version>${maven-failsafe-plugin.version}</version>
<configuration>
+
<classesDirectory>${project.build.outputDirectory}</classesDirectory>
<includes>
<include>**/*ITCase.java</include>
</includes>
@@ -355,4 +356,4 @@ under the License.
</plugin>
</plugins>
</build>
-</project>
\ No newline at end of file
+</project>
diff --git
a/fs_brokers/cdc_client/src/main/java/com/github/shyiko/mysql/binlog/BinaryLogClient.java
b/fs_brokers/cdc_client/src/main/java/com/github/shyiko/mysql/binlog/BinaryLogClient.java
new file mode 100644
index 00000000000..49d4b71de41
--- /dev/null
+++
b/fs_brokers/cdc_client/src/main/java/com/github/shyiko/mysql/binlog/BinaryLogClient.java
@@ -0,0 +1,1477 @@
+/*
+ * Copyright 2013 Stanley Shyiko
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package com.github.shyiko.mysql.binlog;
+
+import com.github.shyiko.mysql.binlog.event.AnnotateRowsEventData;
+import com.github.shyiko.mysql.binlog.event.Event;
+import com.github.shyiko.mysql.binlog.event.EventHeader;
+import com.github.shyiko.mysql.binlog.event.EventHeaderV4;
+import com.github.shyiko.mysql.binlog.event.EventType;
+import com.github.shyiko.mysql.binlog.event.GtidEventData;
+import com.github.shyiko.mysql.binlog.event.MariadbGtidEventData;
+import com.github.shyiko.mysql.binlog.event.MariadbGtidListEventData;
+import com.github.shyiko.mysql.binlog.event.QueryEventData;
+import com.github.shyiko.mysql.binlog.event.RotateEventData;
+import
com.github.shyiko.mysql.binlog.event.deserialization.AnnotateRowsEventDataDeserializer;
+import com.github.shyiko.mysql.binlog.event.deserialization.ChecksumType;
+import
com.github.shyiko.mysql.binlog.event.deserialization.EventDataDeserializationException;
+import
com.github.shyiko.mysql.binlog.event.deserialization.EventDataDeserializer;
+import com.github.shyiko.mysql.binlog.event.deserialization.EventDeserializer;
+import
com.github.shyiko.mysql.binlog.event.deserialization.EventDeserializer.EventDataWrapper;
+import
com.github.shyiko.mysql.binlog.event.deserialization.GtidEventDataDeserializer;
+import
com.github.shyiko.mysql.binlog.event.deserialization.MariadbGtidEventDataDeserializer;
+import
com.github.shyiko.mysql.binlog.event.deserialization.MariadbGtidListEventDataDeserializer;
+import
com.github.shyiko.mysql.binlog.event.deserialization.QueryEventDataDeserializer;
+import
com.github.shyiko.mysql.binlog.event.deserialization.RotateEventDataDeserializer;
+import com.github.shyiko.mysql.binlog.io.ByteArrayInputStream;
+import com.github.shyiko.mysql.binlog.jmx.BinaryLogClientMXBean;
+import com.github.shyiko.mysql.binlog.network.AuthenticationException;
+import com.github.shyiko.mysql.binlog.network.Authenticator;
+import com.github.shyiko.mysql.binlog.network.ClientCapabilities;
+import com.github.shyiko.mysql.binlog.network.DefaultSSLSocketFactory;
+import com.github.shyiko.mysql.binlog.network.SSLMode;
+import com.github.shyiko.mysql.binlog.network.SSLSocketFactory;
+import com.github.shyiko.mysql.binlog.network.ServerException;
+import com.github.shyiko.mysql.binlog.network.SocketFactory;
+import com.github.shyiko.mysql.binlog.network.TLSHostnameVerifier;
+import com.github.shyiko.mysql.binlog.network.protocol.ErrorPacket;
+import com.github.shyiko.mysql.binlog.network.protocol.GreetingPacket;
+import com.github.shyiko.mysql.binlog.network.protocol.Packet;
+import com.github.shyiko.mysql.binlog.network.protocol.PacketChannel;
+import com.github.shyiko.mysql.binlog.network.protocol.ResultSetRowPacket;
+import com.github.shyiko.mysql.binlog.network.protocol.command.Command;
+import
com.github.shyiko.mysql.binlog.network.protocol.command.DumpBinaryLogCommand;
+import
com.github.shyiko.mysql.binlog.network.protocol.command.DumpBinaryLogGtidCommand;
+import com.github.shyiko.mysql.binlog.network.protocol.command.PingCommand;
+import com.github.shyiko.mysql.binlog.network.protocol.command.QueryCommand;
+import
com.github.shyiko.mysql.binlog.network.protocol.command.SSLRequestCommand;
+
+import javax.net.ssl.SSLContext;
+import javax.net.ssl.TrustManager;
+import javax.net.ssl.X509TrustManager;
+import java.io.EOFException;
+import java.io.IOException;
+import java.net.InetSocketAddress;
+import java.net.Socket;
+import java.net.SocketException;
+import java.security.GeneralSecurityException;
+import java.security.cert.CertificateException;
+import java.security.cert.X509Certificate;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.LinkedList;
+import java.util.List;
+import java.util.Locale;
+import java.util.concurrent.Callable;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ThreadFactory;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.concurrent.locks.Lock;
+import java.util.concurrent.locks.ReentrantLock;
+import java.util.logging.Level;
+import java.util.logging.Logger;
+
+/**
+ * MySQL replication stream client.
+ *
+ * <p>Copied from com.zendesk:mysql-binlog-connector-java:0.27.2.
+ *
+ * <p>Lines 149-150, 881, and 1075-1189: replay incomplete non-GTID
transactions from their start
+ * after a keepalive reconnect. See debezium/dbz#2359.
+ *
+ * @author <a href="mailto:[email protected]">Stanley Shyiko</a>
+ */
+public class BinaryLogClient implements BinaryLogClientMXBean {
+
+ private static final SSLSocketFactory
DEFAULT_REQUIRED_SSL_MODE_SOCKET_FACTORY = new DefaultSSLSocketFactory() {
+
+ @Override
+ protected void initSSLContext(SSLContext sc) throws
GeneralSecurityException {
+ sc.init(null, new TrustManager[]{
+ new X509TrustManager() {
+
+ @Override
+ public void checkClientTrusted(X509Certificate[]
x509Certificates, String s)
+ throws CertificateException { }
+
+ @Override
+ public void checkServerTrusted(X509Certificate[]
x509Certificates, String s)
+ throws CertificateException { }
+
+ @Override
+ public X509Certificate[] getAcceptedIssuers() {
+ return new X509Certificate[0];
+ }
+ }
+ }, null);
+ }
+ };
+ private static final SSLSocketFactory
DEFAULT_VERIFY_CA_SSL_MODE_SOCKET_FACTORY = new DefaultSSLSocketFactory();
+
+ // https://dev.mysql.com/doc/internals/en/sending-more-than-16mbyte.html
+ private static final int MAX_PACKET_LENGTH = 16777215;
+
+ private final Logger logger = Logger.getLogger(getClass().getName());
+
+ private final String hostname;
+ private final int port;
+ private final String schema;
+ private final String username;
+ private final String password;
+
+ private boolean blocking = true;
+ private long serverId = 65535;
+ private volatile String binlogFilename;
+ private volatile long binlogPosition = 4;
+ private volatile long connectionId;
+ private SSLMode sslMode = SSLMode.DISABLED;
+
+ protected GtidSet gtidSet;
+ protected final Object gtidSetAccessLock = new Object();
+ private boolean gtidSetFallbackToPurged;
+ private boolean gtidEnabled = false;
+ private boolean useBinlogFilenamePositionInGtidMode;
+ protected String gtid;
+ private boolean tx;
+ private volatile String transactionStartFilename;
+ private volatile long transactionStartPosition;
+
+ private EventDeserializer eventDeserializer = new EventDeserializer();
+
+ private final List<EventListener> eventListeners = new
CopyOnWriteArrayList<EventListener>();
+ private final List<LifecycleListener> lifecycleListeners = new
CopyOnWriteArrayList<LifecycleListener>();
+
+ private SocketFactory socketFactory;
+ private SSLSocketFactory sslSocketFactory;
+
+ protected volatile PacketChannel channel;
+ private volatile boolean connected;
+ private volatile long masterServerId = -1;
+
+ private ThreadFactory threadFactory;
+
+ private boolean keepAlive = true;
+ private long keepAliveInterval = TimeUnit.MINUTES.toMillis(1);
+
+ private long heartbeatInterval;
+ private volatile long eventLastSeen;
+
+ private long connectTimeout = TimeUnit.SECONDS.toMillis(3);
+
+ private volatile ExecutorService keepAliveThreadExecutor;
+
+ private final Lock connectLock = new ReentrantLock();
+ private final Lock keepAliveThreadExecutorLock = new ReentrantLock();
+ private boolean useSendAnnotateRowsEvent;
+
+
+ private Boolean isMariaDB;
+
+ /**
+ * Alias for BinaryLogClient("localhost", 3306, <no schema> = null,
username, password).
+ * @see BinaryLogClient#BinaryLogClient(String, int, String, String,
String)
+ * @param username login name
+ * @param password password
+ */
+ public BinaryLogClient(String username, String password) {
+ this("localhost", 3306, null, username, password);
+ }
+
+ /**
+ * Alias for BinaryLogClient("localhost", 3306, schema, username,
password).
+ * @see BinaryLogClient#BinaryLogClient(String, int, String, String,
String)
+ * @param schema database name, nullable
+ * @param username login name
+ * @param password password
+ */
+ public BinaryLogClient(String schema, String username, String password) {
+ this("localhost", 3306, schema, username, password);
+ }
+
+ /**
+ * Alias for BinaryLogClient(hostname, port, <no schema> = null,
username, password).
+ * @see BinaryLogClient#BinaryLogClient(String, int, String, String,
String)
+ * @param hostname mysql server hostname
+ * @param port mysql server port
+ * @param username login name
+ * @param password password
+ */
+ public BinaryLogClient(String hostname, int port, String username, String
password) {
+ this(hostname, port, null, username, password);
+ }
+
+ /**
+ * @param hostname mysql server hostname
+ * @param port mysql server port
+ * @param schema database name, nullable. Note that this parameter has
nothing to do with event filtering. It's
+ * used only during the authentication.
+ * @param username login name
+ * @param password password
+ */
+ public BinaryLogClient(String hostname, int port, String schema, String
username, String password) {
+ this.hostname = hostname;
+ this.port = port;
+ this.schema = schema;
+ this.username = username;
+ this.password = password;
+ }
+
+ public boolean isBlocking() {
+ return blocking;
+ }
+
+ /**
+ * @param blocking blocking mode. If set to false - BinaryLogClient will
disconnect after the last event.
+ */
+ public void setBlocking(boolean blocking) {
+ this.blocking = blocking;
+ }
+
+ public SSLMode getSSLMode() {
+ return sslMode;
+ }
+
+ public void setSSLMode(SSLMode sslMode) {
+ if (sslMode == null) {
+ throw new IllegalArgumentException("SSL mode cannot be NULL");
+ }
+ this.sslMode = sslMode;
+ }
+
+ public long getMasterServerId() {
+ return this.masterServerId;
+ }
+
+ /**
+ * @return server id (65535 by default)
+ * @see #setServerId(long)
+ */
+ public long getServerId() {
+ return serverId;
+ }
+
+ /**
+ * @param serverId server id (in the range from 1 to 2^32 - 1). This value
MUST be unique across whole replication
+ * group (that is, different from any other server id being used by any
master or slave). Keep in mind that each
+ * binary log client (mysql-binlog-connector-java/BinaryLogClient,
mysqlbinlog, etc) should be treated as a
+ * simplified slave and thus MUST also use a different server id.
+ * @see #getServerId()
+ */
+ public void setServerId(long serverId) {
+ this.serverId = serverId;
+ }
+
+ /**
+ * @return binary log filename, nullable (and null be default). Note that
this value is automatically tracked by
+ * the client and thus is subject to change (in response to {@link
EventType#ROTATE}, for example).
+ * @see #setBinlogFilename(String)
+ */
+ public String getBinlogFilename() {
+ return binlogFilename;
+ }
+
+ /**
+ * @param binlogFilename binary log filename.
+ * Special values are:
+ * <ul>
+ * <li>null, which turns on automatic resolution (resulting in the last
known binlog and position). This is what
+ * happens by default when you don't specify binary log filename
explicitly.</li>
+ * <li>"" (empty string), which instructs server to stream events
starting from the oldest known binlog.</li>
+ * </ul>
+ * @see #getBinlogFilename()
+ */
+ public void setBinlogFilename(String binlogFilename) {
+ this.binlogFilename = binlogFilename;
+ }
+
+ /**
+ * @return binary log position of the next event, 4 by default (which is a
position of first event). Note that this
+ * value changes with each incoming event.
+ * @see #setBinlogPosition(long)
+ */
+ public long getBinlogPosition() {
+ return binlogPosition;
+ }
+
+ /**
+ * @param binlogPosition binary log position. Any value less than 4 gets
automatically adjusted to 4 on connect.
+ * @see #getBinlogPosition()
+ */
+ public void setBinlogPosition(long binlogPosition) {
+ this.binlogPosition = binlogPosition;
+ }
+
+ /**
+ * @return thread id
+ */
+ public long getConnectionId() {
+ return connectionId;
+ }
+
+ /**
+ * @return GTID set. Note that this value changes with each received GTID
event (provided client is in GTID mode).
+ * @see #setGtidSet(String)
+ */
+ public String getGtidSet() {
+ synchronized (gtidSetAccessLock) {
+ return gtidSet != null ? gtidSet.toString() : null;
+ }
+ }
+
+ /**
+ * @param gtidStr GTID set string (can be an empty string).
+ * <p>NOTE #1: Any value but null will switch BinaryLogClient into a GTID
mode (this will also set binlogFilename
+ * to "" (provided it's null) forcing MySQL to send events starting from
the oldest known binlog (keep in mind
+ * that connection will fail if gtid_purged is anything but empty (unless
+ * {@link #setGtidSetFallbackToPurged(boolean)} is set to true))).
+ * <p>NOTE #2: GTID set is automatically updated with each incoming GTID
event (provided GTID mode is on).
+ * @see #getGtidSet()
+ * @see #setGtidSetFallbackToPurged(boolean)
+ */
+ public void setGtidSet(String gtidStr) {
+ if ( gtidStr == null )
+ return;
+
+ this.gtidEnabled = true;
+
+ if (this.binlogFilename == null) {
+ this.binlogFilename = "";
+ }
+
+ synchronized (gtidSetAccessLock) {
+ if ( !gtidStr.equals("") ) {
+ if ( MariadbGtidSet.isMariaGtidSet(gtidStr) ) {
+ this.gtidSet = new MariadbGtidSet(gtidStr);
+ } else {
+ this.gtidSet = new GtidSet(gtidStr);
+ }
+ }
+ }
+ }
+
+ /**
+ * @see #setGtidSetFallbackToPurged(boolean)
+ * @return whether gtid_purged is used as a fallback
+ */
+ public boolean isGtidSetFallbackToPurged() {
+ return gtidSetFallbackToPurged;
+ }
+
+ /**
+ * @param gtidSetFallbackToPurged true if gtid_purged should be used as a
fallback when gtidSet is set to "" and
+ * MySQL server has purged some of the binary logs, false otherwise
(default).
+ */
+ public void setGtidSetFallbackToPurged(boolean gtidSetFallbackToPurged) {
+ this.gtidSetFallbackToPurged = gtidSetFallbackToPurged;
+ }
+
+ /**
+ * @see #setUseBinlogFilenamePositionInGtidMode(boolean)
+ * @return value of useBinlogFilenamePostionInGtidMode
+ */
+ public boolean isUseBinlogFilenamePositionInGtidMode() {
+ return useBinlogFilenamePositionInGtidMode;
+ }
+
+ /**
+ * @param useBinlogFilenamePositionInGtidMode true if MySQL server should
start streaming events from a given
+ * {@link #getBinlogFilename()} and {@link #getBinlogPosition()} instead
of "the oldest known binlog" when
+ * {@link #getGtidSet()} is set, false otherwise (default).
+ */
+ public void setUseBinlogFilenamePositionInGtidMode(boolean
useBinlogFilenamePositionInGtidMode) {
+ this.useBinlogFilenamePositionInGtidMode =
useBinlogFilenamePositionInGtidMode;
+ }
+
+ /**
+ * @return true if "keep alive" thread should be automatically started
(default), false otherwise.
+ * @see #setKeepAlive(boolean)
+ */
+ public boolean isKeepAlive() {
+ return keepAlive;
+ }
+
+ /**
+ * @param keepAlive true if "keep alive" thread should be automatically
started (recommended and true by default),
+ * false otherwise.
+ * @see #isKeepAlive()
+ * @see #setKeepAliveInterval(long)
+ */
+ public void setKeepAlive(boolean keepAlive) {
+ this.keepAlive = keepAlive;
+ }
+
+ /**
+ * @return "keep alive" interval in milliseconds, 1 minute by default.
+ * @see #setKeepAliveInterval(long)
+ */
+ public long getKeepAliveInterval() {
+ return keepAliveInterval;
+ }
+
+ /**
+ * @param keepAliveInterval "keep alive" interval in milliseconds.
+ * @see #getKeepAliveInterval()
+ * @see #setHeartbeatInterval(long)
+ */
+ public void setKeepAliveInterval(long keepAliveInterval) {
+ this.keepAliveInterval = keepAliveInterval;
+ }
+
+ /**
+ * @return "keep alive" connect timeout in milliseconds.
+ * @see #setKeepAliveConnectTimeout(long)
+ *
+ * @deprecated in favour of {@link #getConnectTimeout()}
+ */
+ public long getKeepAliveConnectTimeout() {
+ return connectTimeout;
+ }
+
+ /**
+ * @param connectTimeout "keep alive" connect timeout in milliseconds.
+ * @see #getKeepAliveConnectTimeout()
+ *
+ * @deprecated in favour of {@link #setConnectTimeout(long)}
+ */
+ public void setKeepAliveConnectTimeout(long connectTimeout) {
+ this.connectTimeout = connectTimeout;
+ }
+
+ /**
+ * @return heartbeat period in milliseconds (0 if not set (default)).
+ * @see #setHeartbeatInterval(long)
+ */
+ public long getHeartbeatInterval() {
+ return heartbeatInterval;
+ }
+
+ /**
+ * @param heartbeatInterval heartbeat period in milliseconds.
+ * <p>
+ * If set (recommended)
+ * <ul>
+ * <li> HEARTBEAT event will be emitted every "heartbeatInterval".
+ * <li> if {@link #setKeepAlive(boolean)} is on then keepAlive thread will
attempt to reconnect if no
+ * HEARTBEAT events were received within {@link
#setKeepAliveInterval(long)} (instead of trying to send
+ * PING every {@link #setKeepAliveInterval(long)}, which is
fundamentally flawed -
+ * https://github.com/shyiko/mysql-binlog-connector-java/issues/118).
+ * </ul>
+ * Note that when used together with keepAlive heartbeatInterval MUST be
set less than keepAliveInterval.
+ *
+ * @see #getHeartbeatInterval()
+ */
+ public void setHeartbeatInterval(long heartbeatInterval) {
+ this.heartbeatInterval = heartbeatInterval;
+ }
+
+ /**
+ * @return connect timeout in milliseconds, 3 seconds by default.
+ * @see #setConnectTimeout(long)
+ */
+ public long getConnectTimeout() {
+ return connectTimeout;
+ }
+
+ /**
+ * @param connectTimeout connect timeout in milliseconds.
+ * @see #getConnectTimeout()
+ */
+ public void setConnectTimeout(long connectTimeout) {
+ this.connectTimeout = connectTimeout;
+ }
+
+ /**
+ * @param eventDeserializer custom event deserializer
+ */
+ public void setEventDeserializer(EventDeserializer eventDeserializer) {
+ if (eventDeserializer == null) {
+ throw new IllegalArgumentException("Event deserializer cannot be
NULL");
+ }
+ this.eventDeserializer = eventDeserializer;
+ }
+
+ /**
+ * @param socketFactory custom socket factory. If not provided, socket
will be created with "new Socket()".
+ */
+ public void setSocketFactory(SocketFactory socketFactory) {
+ this.socketFactory = socketFactory;
+ }
+
+ /**
+ * @param sslSocketFactory custom ssl socket factory
+ */
+ public void setSslSocketFactory(SSLSocketFactory sslSocketFactory) {
+ this.sslSocketFactory = sslSocketFactory;
+ }
+
+ /**
+ * @param threadFactory custom thread factory. If not provided, threads
will be created using simple "new Thread()".
+ */
+ public void setThreadFactory(ThreadFactory threadFactory) {
+ this.threadFactory = threadFactory;
+ }
+
+
+ /**
+ * @return true/false depending on whether we've connected to MariaDB.
NULL if not connected.
+ */
+ public Boolean getMariaDB() {
+ return isMariaDB;
+ }
+
+ public boolean isUseSendAnnotateRowsEvent() {
+ return useSendAnnotateRowsEvent;
+ }
+
+ public void setUseSendAnnotateRowsEvent(boolean useSendAnnotateRowsEvent) {
+ this.useSendAnnotateRowsEvent = useSendAnnotateRowsEvent;
+ }
+ /**
+ * Connect to the replication stream. Note that this method blocks until
disconnected.
+ * @throws AuthenticationException if authentication fails
+ * @throws ServerException if MySQL server responds with an error
+ * @throws IOException if anything goes wrong while trying to connect
+ * @throws IllegalStateException if binary log client is already connected
+ */
+ public void connect() throws IOException, IllegalStateException {
+ if (!connectLock.tryLock()) {
+ throw new IllegalStateException("BinaryLogClient is already
connected");
+ }
+ boolean notifyWhenDisconnected = false;
+ try {
+ Callable cancelDisconnect = null;
+ try {
+ try {
+ long start = System.currentTimeMillis();
+ channel = openChannel();
+ if (connectTimeout > 0 && !isKeepAliveThreadRunning()) {
+ cancelDisconnect = scheduleDisconnectIn(connectTimeout
-
+ (System.currentTimeMillis() - start));
+ }
+ if (channel.getInputStream().peek() == -1) {
+ throw new EOFException();
+ }
+ } catch (IOException e) {
+ throw new IOException("Failed to connect to MySQL on " +
hostname + ":" + port +
+ ". Please make sure it's running.", e);
+ }
+ GreetingPacket greetingPacket = receiveGreeting();
+
+ detectMariaDB(greetingPacket);
+ tryUpgradeToSSL(greetingPacket);
+
+ new Authenticator(greetingPacket, channel, schema, username,
password).authenticate();
+ channel.authenticationComplete();
+
+ connectionId = greetingPacket.getThreadId();
+ if ("".equals(binlogFilename)) {
+ setupGtidSet();
+ }
+ if (binlogFilename == null) {
+ fetchBinlogFilenameAndPosition();
+ }
+ if (binlogPosition < 4) {
+ if (logger.isLoggable(Level.WARNING)) {
+ logger.warning("Binary log position adjusted from " +
binlogPosition + " to " + 4);
+ }
+ binlogPosition = 4;
+ }
+ setupConnection();
+ gtid = null;
+ tx = false;
+ requestBinaryLogStream();
+ } catch (IOException e) {
+ disconnectChannel();
+ throw e;
+ } finally {
+ if (cancelDisconnect != null) {
+ try {
+ cancelDisconnect.call();
+ } catch (Exception e) {
+ if (logger.isLoggable(Level.WARNING)) {
+ logger.warning("\"" + e.getMessage() +
+ "\" was thrown while canceling scheduled
disconnect call");
+ }
+ }
+ }
+ }
+ connected = true;
+ notifyWhenDisconnected = true;
+ if (logger.isLoggable(Level.INFO)) {
+ String position;
+ synchronized (gtidSetAccessLock) {
+ position = gtidSet != null ? gtidSet.toString() :
binlogFilename + "/" + binlogPosition;
+ }
+ logger.info("Connected to " + hostname + ":" + port + " at " +
position +
+ " (" + (blocking ? "sid:" + serverId + ", " : "") + "cid:"
+ connectionId + ")");
+ }
+ for (LifecycleListener lifecycleListener : lifecycleListeners) {
+ lifecycleListener.onConnect(this);
+ }
+ if (keepAlive && !isKeepAliveThreadRunning()) {
+ spawnKeepAliveThread();
+ }
+ ensureEventDataDeserializer(EventType.ROTATE,
RotateEventDataDeserializer.class);
+ ensureEventDataDeserializer(EventType.QUERY,
QueryEventDataDeserializer.class);
+ synchronized (gtidSetAccessLock) {
+ if (this.gtidEnabled) {
+ ensureGtidEventDataDeserializer();
+ }
+ }
+ listenForEventPackets();
+ } finally {
+ connectLock.unlock();
+ if (notifyWhenDisconnected) {
+ for (LifecycleListener lifecycleListener : lifecycleListeners)
{
+ lifecycleListener.onDisconnect(this);
+ }
+ }
+ }
+ }
+
+ private void detectMariaDB(GreetingPacket packet) {
+ String serverVersion = packet.getServerVersion();
+ if ( serverVersion == null )
+ return;
+
+ this.isMariaDB = serverVersion.toLowerCase().contains("mariadb");
+ }
+ /**
+ * Apply additional options for connection before requesting binlog stream.
+ */
+ protected void setupConnection() throws IOException {
+ ChecksumType checksumType = fetchBinlogChecksum();
+ if (checksumType != ChecksumType.NONE) {
+ confirmSupportOfChecksum(checksumType);
+ }
+ setMasterServerId();
+ if (heartbeatInterval > 0) {
+ enableHeartbeat();
+ }
+ }
+
+ private PacketChannel openChannel() throws IOException {
+ Socket socket = socketFactory != null ? socketFactory.createSocket() :
new Socket();
+ socket.connect(new InetSocketAddress(hostname, port), (int)
connectTimeout);
+ return new PacketChannel(socket);
+ }
+
+ private Callable scheduleDisconnectIn(final long timeout) {
+ final BinaryLogClient self = this;
+ final CountDownLatch connectLatch = new CountDownLatch(1);
+ final Thread thread = newNamedThread(new Runnable() {
+ @Override
+ public void run() {
+ try {
+ connectLatch.await(timeout, TimeUnit.MILLISECONDS);
+ } catch (InterruptedException e) {
+ if (logger.isLoggable(Level.WARNING)) {
+ logger.log(Level.WARNING, e.getMessage());
+ }
+ }
+ if (connectLatch.getCount() != 0) {
+ if (logger.isLoggable(Level.WARNING)) {
+ logger.warning("Failed to establish connection in " +
timeout + "ms. " +
+ "Forcing disconnect.");
+ }
+ try {
+ self.disconnectChannel();
+ } catch (IOException e) {
+ if (logger.isLoggable(Level.WARNING)) {
+ logger.log(Level.WARNING, e.getMessage());
+ }
+ }
+ }
+ }
+ }, "blc-disconnect-" + hostname + ":" + port);
+ thread.start();
+ return new Callable() {
+
+ public Object call() throws Exception {
+ connectLatch.countDown();
+ thread.join();
+ return null;
+ }
+ };
+ }
+
+ protected void checkError(byte[] packet) throws IOException {
+ if (packet[0] == (byte) 0xFF /* error */) {
+ byte[] bytes = Arrays.copyOfRange(packet, 1, packet.length);
+ ErrorPacket errorPacket = new ErrorPacket(bytes);
+ throw new ServerException(errorPacket.getErrorMessage(),
errorPacket.getErrorCode(),
+ errorPacket.getSqlState());
+ }
+ }
+
+ private GreetingPacket receiveGreeting() throws IOException {
+ byte[] initialHandshakePacket = channel.read();
+ checkError(initialHandshakePacket);
+
+ return new GreetingPacket(initialHandshakePacket);
+ }
+
+ private boolean tryUpgradeToSSL(GreetingPacket greetingPacket) throws
IOException {
+ int collation = greetingPacket.getServerCollation();
+
+ if (sslMode != SSLMode.DISABLED) {
+ boolean serverSupportsSSL =
(greetingPacket.getServerCapabilities() & ClientCapabilities.SSL) != 0;
+ if (!serverSupportsSSL && (sslMode == SSLMode.REQUIRED || sslMode
== SSLMode.VERIFY_CA ||
+ sslMode == SSLMode.VERIFY_IDENTITY)) {
+ throw new IOException("MySQL server does not support SSL");
+ }
+ if (serverSupportsSSL) {
+ SSLRequestCommand sslRequestCommand = new SSLRequestCommand();
+ sslRequestCommand.setCollation(collation);
+ channel.write(sslRequestCommand);
+ SSLSocketFactory sslSocketFactory =
+ this.sslSocketFactory != null ?
+ this.sslSocketFactory :
+ sslMode == SSLMode.REQUIRED || sslMode ==
SSLMode.PREFERRED ?
+ DEFAULT_REQUIRED_SSL_MODE_SOCKET_FACTORY :
+ DEFAULT_VERIFY_CA_SSL_MODE_SOCKET_FACTORY;
+ channel.upgradeToSSL(sslSocketFactory,
+ sslMode == SSLMode.VERIFY_IDENTITY ? new
TLSHostnameVerifier() : null);
+ logger.info("SSL enabled");
+ return true;
+ }
+ }
+ return false;
+ }
+
+ private void enableHeartbeat() throws IOException {
+ channel.write(new QueryCommand("set @master_heartbeat_period=" +
heartbeatInterval * 1000000));
+ byte[] statementResult = channel.read();
+ checkError(statementResult);
+ }
+
+ private void setMasterServerId() throws IOException {
+ channel.write(new QueryCommand("select @@server_id"));
+ ResultSetRowPacket[] resultSet = readResultSet();
+ if (resultSet.length >= 0) {
+ this.masterServerId = Long.parseLong(resultSet[0].getValue(0));
+ }
+ }
+
+ protected void requestBinaryLogStream() throws IOException {
+ long serverId = blocking ? this.serverId : 0; //
http://bugs.mysql.com/bug.php?id=71178
+ if ( this.isMariaDB )
+ requestBinaryLogStreamMaria(serverId);
+ else
+ requestBinaryLogStreamMysql(serverId);
+ }
+
+ private void requestBinaryLogStreamMysql(long serverId) throws IOException
{
+ Command dumpBinaryLogCommand;
+ synchronized (gtidSetAccessLock) {
+ if (this.gtidEnabled) {
+ dumpBinaryLogCommand = new DumpBinaryLogGtidCommand(serverId,
+ useBinlogFilenamePositionInGtidMode ? binlogFilename : "",
+ useBinlogFilenamePositionInGtidMode ? binlogPosition : 4,
+ gtidSet);
+ } else {
+ dumpBinaryLogCommand = new DumpBinaryLogCommand(serverId,
binlogFilename, binlogPosition);
+ }
+ }
+ channel.write(dumpBinaryLogCommand);
+ }
+
+ protected void requestBinaryLogStreamMaria(long serverId) throws
IOException {
+ Command dumpBinaryLogCommand;
+
+ /*
+ https://jira.mariadb.org/browse/MDEV-225
+ */
+ channel.write(new QueryCommand("SET @mariadb_slave_capability=1"));
+ checkError(channel.read());
+
+ synchronized (gtidSetAccessLock) {
+ if (this.gtidEnabled) {
+ logger.info(gtidSet.toString());
+ channel.write(new QueryCommand("SET @slave_connect_state = '"
+ gtidSet.toString() + "'"));
+ checkError(channel.read());
+ channel.write(new QueryCommand("SET @slave_gtid_strict_mode =
0"));
+ checkError(channel.read());
+ channel.write(new QueryCommand("SET
@slave_gtid_ignore_duplicates = 0"));
+ checkError(channel.read());
+ dumpBinaryLogCommand = new DumpBinaryLogCommand(serverId, "",
0L, isUseSendAnnotateRowsEvent());
+ } else {
+ dumpBinaryLogCommand = new DumpBinaryLogCommand(serverId,
binlogFilename, binlogPosition);
+ }
+ }
+ channel.write(dumpBinaryLogCommand);
+ }
+
+ protected void ensureEventDataDeserializer(EventType eventType,
+ Class<? extends EventDataDeserializer>
eventDataDeserializerClass) {
+ EventDataDeserializer eventDataDeserializer =
eventDeserializer.getEventDataDeserializer(eventType);
+ if (eventDataDeserializer.getClass() != eventDataDeserializerClass &&
+ eventDataDeserializer.getClass() !=
EventDataWrapper.Deserializer.class) {
+ EventDataDeserializer internalEventDataDeserializer;
+ try {
+ internalEventDataDeserializer =
eventDataDeserializerClass.newInstance();
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ eventDeserializer.setEventDataDeserializer(eventType,
+ new
EventDataWrapper.Deserializer(internalEventDataDeserializer,
+ eventDataDeserializer));
+ }
+ }
+
+ protected void ensureGtidEventDataDeserializer() {
+ ensureEventDataDeserializer(EventType.GTID,
GtidEventDataDeserializer.class);
+ ensureEventDataDeserializer(EventType.QUERY,
QueryEventDataDeserializer.class);
+ ensureEventDataDeserializer(EventType.ANNOTATE_ROWS,
AnnotateRowsEventDataDeserializer.class);
+ ensureEventDataDeserializer(EventType.MARIADB_GTID,
MariadbGtidEventDataDeserializer.class);
+ ensureEventDataDeserializer(EventType.MARIADB_GTID_LIST,
MariadbGtidListEventDataDeserializer.class);
+ }
+
+ private void spawnKeepAliveThread() {
+ final ExecutorService threadExecutor =
+ Executors.newSingleThreadExecutor(new ThreadFactory() {
+
+ @Override
+ public Thread newThread(Runnable runnable) {
+ return newNamedThread(runnable, "blc-keepalive-" +
hostname + ":" + port);
+ }
+ });
+ try {
+ keepAliveThreadExecutorLock.lock();
+ threadExecutor.submit(new Runnable() {
+ @Override
+ public void run() {
+ while (!threadExecutor.isShutdown()) {
+ try {
+ Thread.sleep(keepAliveInterval);
+ } catch (InterruptedException e) {
+ // expected in case of disconnect
+ }
+ if (threadExecutor.isShutdown()) {
+ logger.info("threadExecutor is shut down,
terminating keepalive thread");
+ return;
+ }
+ boolean connectionLost = false;
+ if (heartbeatInterval > 0) {
+ connectionLost = System.currentTimeMillis() -
eventLastSeen > keepAliveInterval;
+ } else {
+ try {
+ channel.write(new PingCommand());
+ } catch (IOException e) {
+ connectionLost = true;
+ }
+ }
+ if (connectionLost) {
+ logger.info("Keepalive: Trying to restore lost
connection to " + hostname + ":" + port);
+ try {
+ terminateConnect();
+ rewindToTransactionStartIfNeeded();
+ connect(connectTimeout);
+ } catch (Exception ce) {
+ logger.warning("keepalive: Failed to restore
connection to " + hostname + ":" + port +
+ ". Next attempt in " + keepAliveInterval +
"ms");
+ }
+ }
+ }
+ }
+ });
+ keepAliveThreadExecutor = threadExecutor;
+ } finally {
+ keepAliveThreadExecutorLock.unlock();
+ }
+ }
+
+ private Thread newNamedThread(Runnable runnable, String threadName) {
+ Thread thread = threadFactory == null ? new Thread(runnable) :
threadFactory.newThread(runnable);
+ thread.setName(threadName);
+ return thread;
+ }
+
+ boolean isKeepAliveThreadRunning() {
+ try {
+ keepAliveThreadExecutorLock.lock();
+ return keepAliveThreadExecutor != null &&
!keepAliveThreadExecutor.isShutdown();
+ } finally {
+ keepAliveThreadExecutorLock.unlock();
+ }
+ }
+
+ /**
+ * Connect to the replication stream in a separate thread.
+ * @param timeout timeout in milliseconds
+ * @throws AuthenticationException if authentication fails
+ * @throws ServerException if MySQL server responds with an error
+ * @throws IOException if anything goes wrong while trying to connect
+ * @throws TimeoutException if client was unable to connect within given
time limit
+ */
+ public void connect(final long timeout) throws IOException,
TimeoutException {
+ final CountDownLatch countDownLatch = new CountDownLatch(1);
+ AbstractLifecycleListener connectListener = new
AbstractLifecycleListener() {
+ @Override
+ public void onConnect(BinaryLogClient client) {
+ countDownLatch.countDown();
+ }
+ };
+ registerLifecycleListener(connectListener);
+ final AtomicReference<IOException> exceptionReference = new
AtomicReference<IOException>();
+ Runnable runnable = new Runnable() {
+
+ @Override
+ public void run() {
+ try {
+ setConnectTimeout(timeout);
+ connect();
+ } catch (IOException e) {
+ exceptionReference.set(e);
+ countDownLatch.countDown(); // making sure we don't end up
waiting whole "timeout"
+ } catch (Exception e) {
+ exceptionReference.set(new IOException(e)); // method is
asynchronous, catch all exceptions so that they are not lost
+ countDownLatch.countDown(); // making sure we don't end up
waiting whole "timeout"
+ }
+ }
+ };
+ newNamedThread(runnable, "blc-" + hostname + ":" + port).start();
+ boolean started = false;
+ try {
+ started = countDownLatch.await(timeout, TimeUnit.MILLISECONDS);
+ } catch (InterruptedException e) {
+ if (logger.isLoggable(Level.WARNING)) {
+ logger.log(Level.WARNING, e.getMessage());
+ }
+ }
+ unregisterLifecycleListener(connectListener);
+ if (exceptionReference.get() != null) {
+ throw exceptionReference.get();
+ }
+ if (!started) {
+ try {
+ terminateConnect();
+ } finally {
+ throw new TimeoutException("BinaryLogClient was unable to
connect in " + timeout + "ms");
+ }
+ }
+ }
+
+ /**
+ * @return true if client is connected, false otherwise
+ */
+ public boolean isConnected() {
+ return connected;
+ }
+
+ private String fetchGtidPurged() throws IOException {
+ channel.write(new QueryCommand("show global variables like
'gtid_purged'"));
+ ResultSetRowPacket[] resultSet = readResultSet();
+ if (resultSet.length != 0) {
+ return resultSet[0].getValue(1).toUpperCase();
+ }
+ return "";
+ }
+
+ protected void setupGtidSet() throws IOException{
+ if (!this.gtidEnabled)
+ return;
+
+ synchronized (gtidSetAccessLock) {
+ if ( this.isMariaDB ) {
+ if ( gtidSet == null ) {
+ gtidSet = new MariadbGtidSet("");
+ } else if ( !(gtidSet instanceof MariadbGtidSet) ) {
+ throw new RuntimeException("Connected to MariaDB but given
a mysql GTID set!");
+ }
+ } else {
+ if ( gtidSet == null && gtidSetFallbackToPurged ) {
+ gtidSet = new GtidSet(fetchGtidPurged());
+ } else if ( gtidSet == null ){
+ gtidSet = new GtidSet("");
+ } else if ( gtidSet instanceof MariadbGtidSet ) {
+ throw new RuntimeException("Connected to Mysql but given a
MariaDB GTID set!");
+ }
+ }
+ }
+
+ }
+
+ private void fetchBinlogFilenameAndPosition() throws IOException {
+ ResultSetRowPacket[] resultSet;
+ channel.write(new QueryCommand("show master status"));
+ resultSet = readResultSet();
+ if (resultSet.length == 0) {
+ throw new IOException("Failed to determine binlog
filename/position");
+ }
+ ResultSetRowPacket resultSetRow = resultSet[0];
+ binlogFilename = resultSetRow.getValue(0);
+ binlogPosition = Long.parseLong(resultSetRow.getValue(1));
+ }
+
+ private ChecksumType fetchBinlogChecksum() throws IOException {
+ channel.write(new QueryCommand("show global variables like
'binlog_checksum'"));
+ ResultSetRowPacket[] resultSet = readResultSet();
+ if (resultSet.length == 0) {
+ return ChecksumType.NONE;
+ }
+ return ChecksumType.valueOf(resultSet[0].getValue(1).toUpperCase());
+ }
+
+ private void confirmSupportOfChecksum(ChecksumType checksumType) throws
IOException {
+ channel.write(new QueryCommand("set @master_binlog_checksum=
@@global.binlog_checksum"));
+ byte[] statementResult = channel.read();
+ checkError(statementResult);
+ eventDeserializer.setChecksumType(checksumType);
+ }
+
+ private void listenForEventPackets() throws IOException {
+ ByteArrayInputStream inputStream = channel.getInputStream();
+ boolean completeShutdown = false;
+ try {
+ while (inputStream.peek() != -1) {
+ int packetLength = inputStream.readInteger(3);
+ inputStream.skip(1); // 1 byte for sequence
+ int marker = inputStream.read();
+ if (marker == 0xFF) {
+ ErrorPacket errorPacket = new
ErrorPacket(inputStream.read(packetLength - 1));
+ throw new ServerException(errorPacket.getErrorMessage(),
errorPacket.getErrorCode(),
+ errorPacket.getSqlState());
+ }
+ if (marker == 0xFE && !blocking) {
+ completeShutdown = true;
+ break;
+ }
+ Event event;
+ try {
+ event = eventDeserializer.nextEvent(packetLength ==
MAX_PACKET_LENGTH ?
+ new
ByteArrayInputStream(readPacketSplitInChunks(inputStream, packetLength - 1)) :
+ inputStream);
+ if (event == null) {
+ throw new EOFException();
+ }
+ } catch (Exception e) {
+ Throwable cause = e instanceof
EventDataDeserializationException ? e.getCause() : e;
+ if (cause instanceof EOFException || cause instanceof
SocketException) {
+ throw e;
+ }
+ if (isConnected()) {
+ for (LifecycleListener lifecycleListener :
lifecycleListeners) {
+
lifecycleListener.onEventDeserializationFailure(this, e);
+ }
+ }
+ continue;
+ }
+ if (isConnected()) {
+ eventLastSeen = System.currentTimeMillis();
+ updateNonGtidTransactionStateBeforeEvent(event);
+ updateGtidSet(event);
+ notifyEventListeners(event);
+ updateClientBinlogFilenameAndPosition(event);
+ updateNonGtidTransactionStateAfterEvent(event);
+ }
+ }
+ } catch (Exception e) {
+ if (isConnected()) {
+ for (LifecycleListener lifecycleListener : lifecycleListeners)
{
+ lifecycleListener.onCommunicationFailure(this, e);
+ }
+ }
+ } finally {
+ if (isConnected()) {
+ if (completeShutdown) {
+ disconnect(); // initiate complete shutdown sequence
(which includes keep alive thread)
+ } else {
+ disconnectChannel();
+ }
+ }
+ }
+ }
+
+ private byte[] readPacketSplitInChunks(ByteArrayInputStream inputStream,
int packetLength) throws IOException {
+ byte[] result = inputStream.read(packetLength);
+ int chunkLength;
+ do {
+ chunkLength = inputStream.readInteger(3);
+ inputStream.skip(1); // 1 byte for sequence
+ result = Arrays.copyOf(result, result.length + chunkLength);
+ inputStream.fill(result, result.length - chunkLength, chunkLength);
+ } while (chunkLength == Packet.MAX_LENGTH);
+ return result;
+ }
+
+ private void updateClientBinlogFilenameAndPosition(Event event) {
+ EventHeader eventHeader = event.getHeader();
+ EventType eventType = eventHeader.getEventType();
+ if (eventType == EventType.ROTATE) {
+ RotateEventData rotateEventData = (RotateEventData)
EventDataWrapper.internal(event.getData());
+ binlogFilename = rotateEventData.getBinlogFilename();
+ binlogPosition = rotateEventData.getBinlogPosition();
+ } else
+ // do not update binlogPosition on TABLE_MAP so that in case of
reconnect (using a different instance of
+ // client) table mapping cache could be reconstructed before hitting
row mutation event
+ if (eventType != EventType.TABLE_MAP && eventHeader instanceof
EventHeaderV4) {
+ EventHeaderV4 trackableEventHeader = (EventHeaderV4) eventHeader;
+ long nextBinlogPosition = trackableEventHeader.getNextPosition();
+ if (nextBinlogPosition > 0) {
+ binlogPosition = nextBinlogPosition;
+ }
+ }
+ }
+
+ // visible for testing
+ void updateNonGtidTransactionStateBeforeEvent(Event event) {
+ synchronized (gtidSetAccessLock) {
+ if (gtidEnabled) {
+ return;
+ }
+ }
+ if (!(event.getHeader() instanceof EventHeaderV4)) {
+ return;
+ }
+ EventType eventType = event.getHeader().getEventType();
+ if (eventType == EventType.ANONYMOUS_GTID || eventType ==
EventType.MARIADB_GTID) {
+ EventHeaderV4 eventHeader = (EventHeaderV4) event.getHeader();
+ transactionStartPosition = eventHeader.getPosition();
+ transactionStartFilename = binlogFilename;
+ } else if (eventType == EventType.QUERY) {
+ QueryEventData queryEventData = (QueryEventData)
EventDataWrapper.internal(event.getData());
+ if ("BEGIN".equals(queryEventData.getSql())) {
+ tx = true;
+ if (transactionStartFilename == null) {
+ EventHeaderV4 eventHeader = (EventHeaderV4)
event.getHeader();
+ transactionStartPosition = eventHeader.getPosition();
+ transactionStartFilename = binlogFilename;
+ }
+ }
+ }
+ }
+
+ // visible for testing
+ void rewindToTransactionStartIfNeeded() {
+ String filename = transactionStartFilename;
+ if (filename != null) {
+ long position = transactionStartPosition;
+ logger.info("Keepalive: Replaying incomplete transaction from " +
+ filename + "/" + position);
+ binlogFilename = filename;
+ binlogPosition = position;
+ }
+ }
+
+ // visible for testing
+ void updateNonGtidTransactionStateAfterEvent(Event event) {
+ EventType eventType = event.getHeader().getEventType();
+ if (eventType == EventType.XID || eventType ==
EventType.TRANSACTION_PAYLOAD) {
+ clearNonGtidTransactionState();
+ } else if (eventType == EventType.QUERY) {
+ QueryEventData queryEventData = (QueryEventData)
EventDataWrapper.internal(event.getData());
+ String sql = queryEventData.getSql();
+ if ("COMMIT".equals(sql) || "ROLLBACK".equals(sql) ||
+ (!"BEGIN".equals(sql) && !tx)) {
+ clearNonGtidTransactionState();
+ }
+ }
+ }
+
+ private void clearNonGtidTransactionState() {
+ tx = false;
+ transactionStartFilename = null;
+ transactionStartPosition = 0;
+ }
+
+ protected void updateGtidSet(Event event) {
+ synchronized (gtidSetAccessLock) {
+ if (gtidSet == null) {
+ return;
+ }
+ }
+ EventHeader eventHeader = event.getHeader();
+ switch(eventHeader.getEventType()) {
+ case GTID:
+ GtidEventData gtidEventData = (GtidEventData)
EventDataWrapper.internal(event.getData());
+ gtid = gtidEventData.getGtid();
+ break;
+ case XID:
+ commitGtid();
+ tx = false;
+ break;
+ case QUERY:
+ QueryEventData queryEventData = (QueryEventData)
EventDataWrapper.internal(event.getData());
+ String sql = queryEventData.getSql();
+ if (sql == null) {
+ break;
+ }
+ commitGtid(sql);
+ break;
+ case ANNOTATE_ROWS:
+ AnnotateRowsEventData annotateRowsEventData =
(AnnotateRowsEventData)
EventDeserializer.EventDataWrapper.internal(event.getData());
+ sql = annotateRowsEventData.getRowsQuery();
+ if (sql == null) {
+ break;
+ }
+ commitGtid(sql);
+ break;
+ case MARIADB_GTID:
+ MariadbGtidEventData mariadbGtidEventData =
(MariadbGtidEventData)
EventDeserializer.EventDataWrapper.internal(event.getData());
+ mariadbGtidEventData.setServerId(eventHeader.getServerId());
+ gtid = mariadbGtidEventData.toString();
+ break;
+ case MARIADB_GTID_LIST:
+ MariadbGtidListEventData mariadbGtidListEventData =
(MariadbGtidListEventData)
EventDeserializer.EventDataWrapper.internal(event.getData());
+ gtid = mariadbGtidListEventData.getMariaGTIDSet().toString();
+ break;
+ default:
+ }
+ }
+
+ protected void commitGtid(String sql) {
+ if ("BEGIN".equals(sql)) {
+ tx = true;
+ } else
+ if ("COMMIT".equals(sql) || "ROLLBACK".equals(sql)) {
+ commitGtid();
+ tx = false;
+ } else
+ if (!tx) {
+ // auto-commit query, likely DDL
+ commitGtid();
+ }
+ }
+
+ private void commitGtid() {
+ if (gtid != null) {
+ synchronized (gtidSetAccessLock) {
+ gtidSet.add(gtid);
+ }
+ }
+ }
+
+ private ResultSetRowPacket[] readResultSet() throws IOException {
+ List<ResultSetRowPacket> resultSet = new LinkedList<>();
+ byte[] statementResult = channel.read();
+ checkError(statementResult);
+
+ while ((channel.read())[0] != (byte) 0xFE /* eof */) { /* skip */ }
+ for (byte[] bytes; (bytes = channel.read())[0] != (byte) 0xFE /* eof
*/; ) {
+ checkError(bytes);
+ resultSet.add(new ResultSetRowPacket(bytes));
+ }
+ return resultSet.toArray(new ResultSetRowPacket[resultSet.size()]);
+ }
+
+ /**
+ * @return registered event listeners
+ */
+ public List<EventListener> getEventListeners() {
+ return Collections.unmodifiableList(eventListeners);
+ }
+
+ /**
+ * Register event listener. Note that multiple event listeners will be
called in order they
+ * where registered.
+ * @param eventListener event listener
+ */
+ public void registerEventListener(EventListener eventListener) {
+ eventListeners.add(eventListener);
+ }
+
+ /**
+ * Unregister all event listener of specific type.
+ * @param listenerClass event listener class to unregister
+ */
+ public void unregisterEventListener(Class<? extends EventListener>
listenerClass) {
+ for (EventListener eventListener: eventListeners) {
+ if (listenerClass.isInstance(eventListener)) {
+ eventListeners.remove(eventListener);
+ }
+ }
+ }
+
+ /**
+ * Unregister single event listener.
+ * @param eventListener event listener to unregister
+ */
+ public void unregisterEventListener(EventListener eventListener) {
+ eventListeners.remove(eventListener);
+ }
+
+ private void notifyEventListeners(Event event) {
+ if (event.getData() instanceof EventDataWrapper) {
+ event = new Event(event.getHeader(), ((EventDataWrapper)
event.getData()).getExternal());
+ }
+ for (EventListener eventListener : eventListeners) {
+ try {
+ eventListener.onEvent(event);
+ } catch (Exception e) {
+ if (logger.isLoggable(Level.WARNING)) {
+ logger.log(Level.WARNING, eventListener + " choked on " +
event, e);
+ }
+ }
+ }
+ }
+
+ /**
+ * @return registered lifecycle listeners
+ */
+ public List<LifecycleListener> getLifecycleListeners() {
+ return Collections.unmodifiableList(lifecycleListeners);
+ }
+
+ /**
+ * Register lifecycle listener. Note that multiple lifecycle listeners
will be called in order they
+ * where registered.
+ * @param lifecycleListener lifecycle listener to register
+ */
+ public void registerLifecycleListener(LifecycleListener lifecycleListener)
{
+ lifecycleListeners.add(lifecycleListener);
+ }
+
+ /**
+ * Unregister all lifecycle listener of specific type.
+ * @param listenerClass lifecycle listener class to unregister
+ */
+ public void unregisterLifecycleListener(Class<? extends LifecycleListener>
listenerClass) {
+ for (LifecycleListener lifecycleListener : lifecycleListeners) {
+ if (listenerClass.isInstance(lifecycleListener)) {
+ lifecycleListeners.remove(lifecycleListener);
+ }
+ }
+ }
+
+ /**
+ * Unregister single lifecycle listener.
+ * @param eventListener lifecycle listener to unregister
+ */
+ public void unregisterLifecycleListener(LifecycleListener eventListener) {
+ lifecycleListeners.remove(eventListener);
+ }
+
+ /**
+ * Disconnect from the replication stream.
+ * Note that this does not cause binlogFilename/binlogPosition to be
cleared out.
+ * As the result following {@link #connect()} resumes client from where it
left off.
+ */
+ public void disconnect() throws IOException {
+ terminateKeepAliveThread();
+ terminateConnect();
+ }
+
+ private void terminateKeepAliveThread() {
+ try {
+ keepAliveThreadExecutorLock.lock();
+ ExecutorService keepAliveThreadExecutor =
this.keepAliveThreadExecutor;
+ if ( keepAliveThreadExecutor == null ) {
+ return;
+ }
+ keepAliveThreadExecutor.shutdownNow();
+ } finally {
+ keepAliveThreadExecutorLock.unlock();
+ }
+ while (!awaitTerminationInterruptibly(keepAliveThreadExecutor,
+ Long.MAX_VALUE, TimeUnit.NANOSECONDS)) {
+ // ignore
+ }
+ }
+
+ private static boolean awaitTerminationInterruptibly(ExecutorService
executorService, long timeout, TimeUnit unit) {
+ try {
+ return executorService.awaitTermination(timeout, unit);
+ } catch (InterruptedException e) {
+ return false;
+ }
+ }
+
+ private void terminateConnect() throws IOException {
+ do {
+ disconnectChannel();
+ } while (!tryLockInterruptibly(connectLock, 1000,
TimeUnit.MILLISECONDS));
+ connectLock.unlock();
+ }
+
+ private static boolean tryLockInterruptibly(Lock lock, long time, TimeUnit
unit) {
+ try {
+ return lock.tryLock(time, unit);
+ } catch (InterruptedException e) {
+ return false;
+ }
+ }
+
+ private void disconnectChannel() throws IOException {
+ connected = false;
+ if (channel != null && channel.isOpen()) {
+ channel.close();
+ }
+ }
+
+ /**
+ * {@link BinaryLogClient}'s event listener.
+ */
+ public interface EventListener {
+
+ void onEvent(Event event);
+ }
+
+ /**
+ * {@link BinaryLogClient}'s lifecycle listener.
+ */
+ public interface LifecycleListener {
+
+ /**
+ * Called once client has successfully logged in but before started to
receive binlog events.
+ * @param client the client that logged in
+ */
+ void onConnect(BinaryLogClient client);
+
+ /**
+ * It's guarantied to be called before {@link
#onDisconnect(BinaryLogClient)}) in case of
+ * communication failure.
+ * @param client the client that triggered the communication
failure
+ * @param ex The exception that triggered the communication
failutre
+ */
+ void onCommunicationFailure(BinaryLogClient client, Exception ex);
+
+ /**
+ * Called in case of failed event deserialization. Note this type of
error does NOT cause client to
+ * disconnect. If you wish to stop receiving events you'll need to
fire client.disconnect() manually.
+ * @param client the client that failed event deserialization
+ * @param ex The exception that triggered the failutre
+ */
+ void onEventDeserializationFailure(BinaryLogClient client, Exception
ex);
+
+ /**
+ * Called upon disconnect (regardless of the reason).
+ * @param client the client that disconnected
+ */
+ void onDisconnect(BinaryLogClient client);
+ }
+
+ /**
+ * Default (no-op) implementation of {@link LifecycleListener}.
+ */
+ public static abstract class AbstractLifecycleListener implements
LifecycleListener {
+
+ public void onConnect(BinaryLogClient client) { }
+
+ public void onCommunicationFailure(BinaryLogClient client, Exception
ex) { }
+
+ public void onEventDeserializationFailure(BinaryLogClient client,
Exception ex) { }
+
+ public void onDisconnect(BinaryLogClient client) { }
+
+ }
+
+}
diff --git
a/fs_brokers/cdc_client/src/main/java/io/debezium/connector/mysql/MySqlStreamingChangeEventSource.java
b/fs_brokers/cdc_client/src/main/java/io/debezium/connector/mysql/MySqlStreamingChangeEventSource.java
index 275ec709211..a945e92de68 100644
---
a/fs_brokers/cdc_client/src/main/java/io/debezium/connector/mysql/MySqlStreamingChangeEventSource.java
+++
b/fs_brokers/cdc_client/src/main/java/io/debezium/connector/mysql/MySqlStreamingChangeEventSource.java
@@ -89,6 +89,9 @@ import static io.debezium.util.Strings.isNullOrEmpty;
* <p>Line 940 : change Log Level info to debug.
*
* <p>Line 420 : exclude OceanBase heartbeat events from restart event
counting.
+ *
+ * <p>Line 245 : use the Debezium progress heartbeat for the MySQL protocol
heartbeat, capped by
+ * the keepalive-safe interval.
*/
public class MySqlStreamingChangeEventSource
implements StreamingChangeEventSource<MySqlPartition,
MySqlOffsetContext> {
@@ -238,13 +241,13 @@ public class MySqlStreamingChangeEventSource
final long keepAliveInterval =
configuration.getLong(MySqlConnectorConfig.KEEP_ALIVE_INTERVAL_MS);
client.setKeepAliveInterval(keepAliveInterval);
- // Considering heartbeatInterval should be less than
keepAliveInterval, we use the
- // heartbeatIntervalFactor
- // multiply by keepAliveInterval and set the result value to
heartbeatInterval.The default
- // value of heartbeatIntervalFactor
- // is 0.8, and we believe the left time (0.2 * keepAliveInterval) is
enough to process the
- // packet received from the MySQL server.
- client.setHeartbeatInterval((long) (keepAliveInterval *
heartbeatIntervalFactor));
+ final long maxHeartbeatInterval =
+ (long) (keepAliveInterval * heartbeatIntervalFactor);
+ final long heartbeatInterval =
connectorConfig.getHeartbeatInterval().toMillis();
+ client.setHeartbeatInterval(
+ heartbeatInterval > 0
+ ? Math.min(heartbeatInterval, maxHeartbeatInterval)
+ : maxHeartbeatInterval);
boolean filterDmlEventsByGtidSource =
configuration.getBoolean(MySqlConnectorConfig.GTID_SOURCE_FILTER_DML_EVENTS);
diff --git
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/common/Constants.java
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/common/Constants.java
index a9eea173d4d..93aa72c4249 100644
---
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/common/Constants.java
+++
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/common/Constants.java
@@ -22,7 +22,7 @@ public class Constants {
public static final long POLL_SPLIT_RECORDS_TIMEOUTS = 15000L;
// Debezium default properties
- public static final long DEBEZIUM_HEARTBEAT_INTERVAL_MS = 3000L;
+ public static final long DEBEZIUM_HEARTBEAT_INTERVAL_MS = 5_000L;
public static final String DORIS_TARGET_DB = "doris_target_db";
diff --git
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlSourceReader.java
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlSourceReader.java
index 586fac39abc..19380dfe421 100644
---
a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlSourceReader.java
+++
b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlSourceReader.java
@@ -1000,9 +1000,8 @@ public class MySqlSourceReader extends
AbstractCdcSourceReader {
configFactory.jdbcProperties(jdbcProperteis);
Properties dbzProps = ConfigUtil.getDefaultDebeziumProps();
- dbzProps.setProperty(
- MySqlConnectorConfig.KEEP_ALIVE_INTERVAL_MS.name(),
- DEBEZIUM_HEARTBEAT_INTERVAL_MS + "");
+ // Do not override KEEP_ALIVE_INTERVAL_MS: connection liveness is
independent from CDC
+ // progress heartbeats.
dbzProps.setProperty(
EXCLUDE_HEARTBEAT_FROM_EVENT_COUNT,
Boolean.toString(excludeHeartbeatFromEventCount()));
diff --git
a/fs_brokers/cdc_client/src/test/java/com/github/shyiko/mysql/binlog/BinaryLogClientTransactionReplayTest.java
b/fs_brokers/cdc_client/src/test/java/com/github/shyiko/mysql/binlog/BinaryLogClientTransactionReplayTest.java
new file mode 100644
index 00000000000..008b19e5c66
--- /dev/null
+++
b/fs_brokers/cdc_client/src/test/java/com/github/shyiko/mysql/binlog/BinaryLogClientTransactionReplayTest.java
@@ -0,0 +1,135 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package com.github.shyiko.mysql.binlog;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+import com.github.shyiko.mysql.binlog.event.Event;
+import com.github.shyiko.mysql.binlog.event.EventHeaderV4;
+import com.github.shyiko.mysql.binlog.event.EventType;
+import com.github.shyiko.mysql.binlog.event.QueryEventData;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.EnumSource;
+import org.junit.jupiter.params.provider.ValueSource;
+
+class BinaryLogClientTransactionReplayTest {
+
+ private static final String BINLOG_FILE = "mysql-bin.000001";
+ private static final long TRANSACTION_START = 100L;
+
+ @Test
+ void rewindsIncompleteNonGtidTransaction() {
+ BinaryLogClient client = clientAtTransactionStart();
+ client.updateNonGtidTransactionStateBeforeEvent(
+ queryEvent("BEGIN", TRANSACTION_START, 150L));
+
+ client.setBinlogPosition(300L);
+ client.rewindToTransactionStartIfNeeded();
+
+ assertThat(client.getBinlogFilename()).isEqualTo(BINLOG_FILE);
+ assertThat(client.getBinlogPosition()).isEqualTo(TRANSACTION_START);
+ }
+
+ @ParameterizedTest
+ @EnumSource(value = EventType.class, names = {"ANONYMOUS_GTID",
"MARIADB_GTID"})
+ void tracksNonGtidTransactionStartAtGtidMarker(EventType transactionStart)
{
+ BinaryLogClient client = clientAtTransactionStart();
+ client.updateNonGtidTransactionStateBeforeEvent(
+ event(transactionStart, TRANSACTION_START, 150L));
+
+ client.setBinlogPosition(300L);
+ client.rewindToTransactionStartIfNeeded();
+
+ assertThat(client.getBinlogPosition()).isEqualTo(TRANSACTION_START);
+ }
+
+ @ParameterizedTest
+ @EnumSource(value = EventType.class, names = {"XID",
"TRANSACTION_PAYLOAD"})
+ void doesNotRewindCompletedNonGtidTransaction(EventType transactionEnd) {
+ BinaryLogClient client = clientAtTransactionStart();
+ client.updateNonGtidTransactionStateBeforeEvent(
+ queryEvent("BEGIN", TRANSACTION_START, 150L));
+ client.updateNonGtidTransactionStateAfterEvent(event(transactionEnd));
+
+ client.setBinlogPosition(300L);
+ client.rewindToTransactionStartIfNeeded();
+
+ assertThat(client.getBinlogPosition()).isEqualTo(300L);
+ }
+
+ @ParameterizedTest
+ @ValueSource(strings = {"COMMIT", "ROLLBACK"})
+ void doesNotRewindCompletedQueryTransaction(String transactionEnd) {
+ BinaryLogClient client = clientAtTransactionStart();
+ client.updateNonGtidTransactionStateBeforeEvent(
+ queryEvent("BEGIN", TRANSACTION_START, 150L));
+
client.updateNonGtidTransactionStateAfterEvent(queryEvent(transactionEnd, 300L,
350L));
+
+ client.setBinlogPosition(350L);
+ client.rewindToTransactionStartIfNeeded();
+
+ assertThat(client.getBinlogPosition()).isEqualTo(350L);
+ }
+
+ @Test
+ void keepsGtidReconnectBehaviorUnchanged() {
+ BinaryLogClient client = clientAtTransactionStart();
+ client.setGtidSet("");
+ client.updateNonGtidTransactionStateBeforeEvent(
+ queryEvent("BEGIN", TRANSACTION_START, 150L));
+
+ client.setBinlogPosition(300L);
+ client.rewindToTransactionStartIfNeeded();
+
+ assertThat(client.getBinlogPosition()).isEqualTo(300L);
+ }
+
+ private static BinaryLogClient clientAtTransactionStart() {
+ BinaryLogClient client = new BinaryLogClient("localhost", 3306,
"root", "password");
+ client.setBinlogFilename(BINLOG_FILE);
+ client.setBinlogPosition(TRANSACTION_START);
+ return client;
+ }
+
+ private static Event queryEvent(String sql, long position, long
nextPosition) {
+ QueryEventData data = new QueryEventData();
+ data.setSql(sql);
+ EventHeaderV4 header = eventHeader(EventType.QUERY);
+ header.setEventLength(nextPosition - position);
+ header.setNextPosition(nextPosition);
+ return new Event(header, data);
+ }
+
+ private static Event event(EventType eventType) {
+ return new Event(eventHeader(eventType), null);
+ }
+
+ private static Event event(EventType eventType, long position, long
nextPosition) {
+ EventHeaderV4 header = eventHeader(eventType);
+ header.setEventLength(nextPosition - position);
+ header.setNextPosition(nextPosition);
+ return new Event(header, null);
+ }
+
+ private static EventHeaderV4 eventHeader(EventType eventType) {
+ EventHeaderV4 header = new EventHeaderV4();
+ header.setEventType(eventType);
+ return header;
+ }
+}
diff --git
a/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/MySqlBinaryLogClientKeepAliveITCase.java
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/MySqlBinaryLogClientKeepAliveITCase.java
new file mode 100644
index 00000000000..0ae7bd4d98c
--- /dev/null
+++
b/fs_brokers/cdc_client/src/test/java/org/apache/doris/cdcclient/itcase/MySqlBinaryLogClientKeepAliveITCase.java
@@ -0,0 +1,208 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.cdcclient.itcase;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+import com.github.shyiko.mysql.binlog.BinaryLogClient;
+import com.github.shyiko.mysql.binlog.event.WriteRowsEventData;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.Timeout;
+import org.testcontainers.containers.MySQLContainer;
+import org.testcontainers.junit.jupiter.Container;
+import org.testcontainers.junit.jupiter.Testcontainers;
+import org.testcontainers.utility.DockerImageName;
+
+import java.io.Serializable;
+import java.sql.Connection;
+import java.sql.DriverManager;
+import java.sql.ResultSet;
+import java.sql.Statement;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.concurrent.locks.LockSupport;
+import java.util.stream.Collectors;
+import java.util.stream.IntStream;
+
+@Testcontainers
+class MySqlBinaryLogClientKeepAliveITCase {
+
+ private static final String ROOT_USER = "root";
+ private static final String ROOT_PASSWORD = "123456";
+ private static final String TABLE = "keepalive_replay";
+ private static final int ROW_COUNT = 500;
+
+ @Container
+ static final MySQLContainer<?> MYSQL =
+ new MySQLContainer<>(DockerImageName.parse("mysql:8.0"))
+ .withDatabaseName("cdc_test")
+ .withUsername("cdc")
+ .withPassword(ROOT_PASSWORD)
+ .withEnv("MYSQL_ROOT_PASSWORD", ROOT_PASSWORD);
+
+ @Test
+ @Timeout(value = 30, unit = TimeUnit.SECONDS)
+ void keepAliveReconnectReplaysIncompleteNonGtidTransaction() throws
Exception {
+ BinlogPosition startPosition;
+ try (Connection connection = rootConnection();
+ Statement statement = connection.createStatement()) {
+ statement.execute("DROP TABLE IF EXISTS " + TABLE);
+ statement.execute("CREATE TABLE " + TABLE + " (id INT PRIMARY
KEY)");
+ startPosition = currentBinlogPosition(statement);
+ }
+
+ BinaryLogClient client =
+ new BinaryLogClient(
+ MYSQL.getHost(),
+ MYSQL.getMappedPort(MySQLContainer.MYSQL_PORT),
+ ROOT_USER,
+ ROOT_PASSWORD);
+ client.setBinlogFilename(startPosition.filename);
+ client.setBinlogPosition(startPosition.position);
+ client.setHeartbeatInterval(100L);
+ client.setKeepAliveInterval(300L);
+ client.setConnectTimeout(3_000L);
+
+ Set<Integer> receivedIds = ConcurrentHashMap.newKeySet();
+ AtomicBoolean interruptFirstRowsEvent = new AtomicBoolean(true);
+ AtomicInteger connectionCount = new AtomicInteger();
+ AtomicReference<Throwable> listenerFailure = new AtomicReference<>();
+ CountDownLatch firstRowsEventInterrupted = new CountDownLatch(1);
+ CountDownLatch reconnected = new CountDownLatch(1);
+ CountDownLatch allRowsReceived = new CountDownLatch(1);
+
+ client.registerLifecycleListener(
+ new BinaryLogClient.AbstractLifecycleListener() {
+ @Override
+ public void onConnect(BinaryLogClient connectedClient) {
+ if (connectionCount.incrementAndGet() > 1) {
+ reconnected.countDown();
+ }
+ }
+ });
+ client.registerEventListener(
+ event -> {
+ if (!(event.getData() instanceof WriteRowsEventData)) {
+ return;
+ }
+ List<Serializable[]> rows =
+ ((WriteRowsEventData) event.getData()).getRows();
+ if (interruptFirstRowsEvent.compareAndSet(true, false)) {
+ if (rows.size() <= 20) {
+ listenerFailure.set(
+ new AssertionError(
+ "Expected one multi-row event, but
received "
+ + rows.size()
+ + " rows"));
+ firstRowsEventInterrupted.countDown();
+ return;
+ }
+ addRows(receivedIds, rows.subList(0, 20));
+ firstRowsEventInterrupted.countDown();
+
+ long deadline = System.nanoTime() +
TimeUnit.SECONDS.toNanos(5);
+ while (client.isConnected() && System.nanoTime() <
deadline) {
+
LockSupport.parkNanos(TimeUnit.MILLISECONDS.toNanos(1));
+ }
+ if (client.isConnected()) {
+ listenerFailure.set(
+ new AssertionError(
+ "Keepalive did not disconnect the
blocked listener"));
+ }
+ return;
+ }
+
+ addRows(receivedIds, rows);
+ if (receivedIds.size() >= ROW_COUNT) {
+ allRowsReceived.countDown();
+ }
+ });
+
+ try {
+ client.connect(5_000L);
+ try (Connection connection = rootConnection();
+ Statement statement = connection.createStatement()) {
+ statement.execute(insertRowsSql());
+ }
+
+ assertThat(firstRowsEventInterrupted.await(5,
TimeUnit.SECONDS)).isTrue();
+ assertThat(listenerFailure.get()).isNull();
+ assertThat(reconnected.await(10, TimeUnit.SECONDS)).isTrue();
+ assertThat(allRowsReceived.await(10, TimeUnit.SECONDS)).isTrue();
+ assertThat(listenerFailure.get()).isNull();
+ assertThat(connectionCount.get()).isGreaterThanOrEqualTo(2);
+
assertThat(receivedIds).containsExactlyInAnyOrderElementsOf(expectedIds());
+ } finally {
+ client.disconnect();
+ }
+ }
+
+ private static void addRows(Set<Integer> receivedIds, List<Serializable[]>
rows) {
+ for (Serializable[] row : rows) {
+ receivedIds.add(((Number) row[0]).intValue());
+ }
+ }
+
+ private static List<Integer> expectedIds() {
+ return IntStream.rangeClosed(1,
ROW_COUNT).boxed().collect(Collectors.toList());
+ }
+
+ private static String insertRowsSql() {
+ List<String> values = new ArrayList<>(ROW_COUNT);
+ for (int id = 1; id <= ROW_COUNT; id++) {
+ values.add("(" + id + ")");
+ }
+ return "INSERT INTO " + TABLE + " VALUES " + String.join(",", values);
+ }
+
+ private static BinlogPosition currentBinlogPosition(Statement statement)
throws Exception {
+ try (ResultSet resultSet = statement.executeQuery("SHOW MASTER
STATUS")) {
+ assertThat(resultSet.next()).isTrue();
+ return new BinlogPosition(resultSet.getString("File"),
resultSet.getLong("Position"));
+ }
+ }
+
+ private static Connection rootConnection() throws Exception {
+ String url =
+ "jdbc:mysql://"
+ + MYSQL.getHost()
+ + ":"
+ + MYSQL.getMappedPort(MySQLContainer.MYSQL_PORT)
+ + "/"
+ + MYSQL.getDatabaseName()
+ + "?serverTimezone=UTC";
+ return DriverManager.getConnection(url, ROOT_USER, ROOT_PASSWORD);
+ }
+
+ private static final class BinlogPosition {
+ private final String filename;
+ private final long position;
+
+ private BinlogPosition(String filename, long position) {
+ this.filename = filename;
+ this.position = position;
+ }
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]