Skip to content

Commit

Permalink
[Bugfix][connector-cdc-mysql] Fix listener not released when BinlogCl…
Browse files Browse the repository at this point in the history
…ient reuse (apache#5011)
  • Loading branch information
happyboy1024 authored and Jarvis committed Jul 13, 2023
1 parent 124f604 commit 9991aa1
Showing 1 changed file with 17 additions and 1 deletion.
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
package org.apache.seatunnel.connectors.seatunnel.cdc.mysql.source.reader.fetch;

import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
import org.apache.seatunnel.common.utils.ReflectionUtils;
import org.apache.seatunnel.connectors.cdc.base.config.JdbcSourceConfig;
import org.apache.seatunnel.connectors.cdc.base.dialect.JdbcDataSourceDialect;
import org.apache.seatunnel.connectors.cdc.base.relational.JdbcSourceEventDispatcher;
Expand Down Expand Up @@ -70,6 +71,7 @@
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.Optional;

import static org.apache.seatunnel.connectors.seatunnel.cdc.mysql.source.offset.BinlogOffset.BINLOG_FILENAME_OFFSET_KEY;
import static org.apache.seatunnel.connectors.seatunnel.cdc.mysql.utils.MySqlConnectionUtils.createBinaryClient;
Expand Down Expand Up @@ -326,13 +328,27 @@ public MySqlTaskContextImpl(
MySqlDatabaseSchema schema,
BinaryLogClient reusedBinaryLogClient) {
super(config, schema);
this.reusedBinaryLogClient = reusedBinaryLogClient;
this.reusedBinaryLogClient = resetBinaryLogClient(reusedBinaryLogClient);
}

@Override
public BinaryLogClient getBinaryLogClient() {
return reusedBinaryLogClient;
}

/** reset the listener of binaryLogClient before fetch task start. */
private BinaryLogClient resetBinaryLogClient(BinaryLogClient binaryLogClient) {
Optional<Object> eventListenersField =
ReflectionUtils.getField(
binaryLogClient, BinaryLogClient.class, "eventListeners");
eventListenersField.ifPresent(o -> ((List<BinaryLogClient.EventListener>) o).clear());
Optional<Object> lifecycleListeners =
ReflectionUtils.getField(
binaryLogClient, BinaryLogClient.class, "lifecycleListeners");
lifecycleListeners.ifPresent(
o -> ((List<BinaryLogClient.LifecycleListener>) o).clear());
return binaryLogClient;
}
}

/** Copied from debezium for accessing here. */
Expand Down

0 comments on commit 9991aa1

Please sign in to comment.