已合并
【bugfix】优化增量迁移性能 #170
lvlintao666创建于 2023年9月18日
【bugfix】优化增量迁移性能 #170
已合并
lvlintao666创建于 2023年9月18日
从refs/pull/170/head合入到master
共 6 个文件变更+13-74
@@ -206,7 +206,6 @@ private void initRecordBreakpoint(MySqlSinkConnectorConfig config) {
206 toDeleteOffsets = breakPointRecord.getToDeleteOffsets();206 toDeleteOffsets = breakPointRecord.getToDeleteOffsets();
207 breakPointRecord.setBpQueueTimeLimit(config.getBpQueueTimeLimit());207 breakPointRecord.setBpQueueTimeLimit(config.getBpQueueTimeLimit());
208 breakPointRecord.setBpQueueSizeLimit(config.getBpQueueSizeLimit());208 breakPointRecord.setBpQueueSizeLimit(config.getBpQueueSizeLimit());
209- breakPointRecord.setIsBpSwitch(config.getBpSwitch());
210 breakPointRecord.start();209 breakPointRecord.start();
211 if (!breakPointRecord.isTopicExist()) {210 if (!breakPointRecord.isTopicExist()) {
212 breakPointRecord.initializeStorage();211 breakPointRecord.initializeStorage();
@@ -47,7 +47,6 @@ public class WorkThread extends Thread {
47 private BreakPointRecord breakPointRecord;47 private BreakPointRecord breakPointRecord;
48 private PriorityBlockingQueue<Long> replayedOffsets;48 private PriorityBlockingQueue<Long> replayedOffsets;
49 private boolean isTransaction;49 private boolean isTransaction;
50- private boolean isBpSwitch;
51 private boolean isConnection = true;50 private boolean isConnection = true;
52 private boolean isAlive = true;51 private boolean isAlive = true;
53 52 
@@ -66,7 +65,6 @@ public WorkThread(ConnectionInfo connectionInfo, BlockingQueue<String> feedBackQ
66 this.feedBackQueue = feedBackQueue;65 this.feedBackQueue = feedBackQueue;
67 this.breakPointRecord = breakPointRecord;66 this.breakPointRecord = breakPointRecord;
68 this.replayedOffsets = breakPointRecord.getReplayedOffset();67 this.replayedOffsets = breakPointRecord.getReplayedOffset();
69- this.isBpSwitch = breakPointRecord.getIsBpSwitch();
70 this.isTransaction = true;68 this.isTransaction = true;
71 }69 }
72 70 
@@ -251,16 +249,14 @@ public List<String> getFailSqlList() {
251 * @param txn the replay transaction249 * @param txn the replay transaction
252 */250 */
253 private void savedBreakPointInfo(Transaction txn) {251 private void savedBreakPointInfo(Transaction txn) {
254- if (isBpSwitch) {252+ BreakPointObject txnBpObject = new BreakPointObject();
255- BreakPointObject txnBpObject = new BreakPointObject();253+ txnBpObject.setBeginOffset(txn.getTxnBeginOffset());
256- txnBpObject.setBeginOffset(txn.getTxnBeginOffset());254+ txnBpObject.setEndOffset(txn.getTxnEndOffset());
257- txnBpObject.setEndOffset(txn.getTxnEndOffset());255+ txnBpObject.setTimeStamp(LocalDateTime.now().toString());
258- txnBpObject.setTimeStamp(LocalDateTime.now().toString());256+ if (!txn.getSourceField().getGtid().isEmpty()) {
259- if (!txn.getSourceField().getGtid().isEmpty()) {257+ txnBpObject.setGtid(txn.getSourceField().getGtid());
260- txnBpObject.setGtid(txn.getSourceField().getGtid());
261- }
262- breakPointRecord.storeRecord(txnBpObject, isTransaction);
263 }258 }
259+ breakPointRecord.storeRecord(txnBpObject, isTransaction);
264 }260 }
265 261 
266 /**262 /**
@@ -195,7 +195,6 @@ private void initRecordBreakpoint(OpengaussSinkConnectorConfig config) {
195 toDeleteOffsets = breakPointRecord.getToDeleteOffsets();195 toDeleteOffsets = breakPointRecord.getToDeleteOffsets();
196 breakPointRecord.setBpQueueTimeLimit(config.getBpQueueTimeLimit());196 breakPointRecord.setBpQueueTimeLimit(config.getBpQueueTimeLimit());
197 breakPointRecord.setBpQueueSizeLimit(config.getBpQueueSizeLimit());197 breakPointRecord.setBpQueueSizeLimit(config.getBpQueueSizeLimit());
198- breakPointRecord.setIsBpSwitch(config.getBpSwitch());
199 breakPointRecord.start();198 breakPointRecord.start();
200 if (!breakPointRecord.isTopicExist()) {199 if (!breakPointRecord.isTopicExist()) {
201 breakPointRecord.initializeStorage();200 breakPointRecord.initializeStorage();
@@ -88,7 +88,6 @@ public class WorkThread extends Thread {
88 private boolean isClearFile;88 private boolean isClearFile;
89 private boolean isTransaction;89 private boolean isTransaction;
90 private boolean isConnection = true;90 private boolean isConnection = true;
91- private boolean isBpSwitch;
92 private boolean isStop = false;91 private boolean isStop = false;
93 92 
94 /**93 /**
@@ -108,7 +107,6 @@ public WorkThread(Map<String, String> schemaMappingMap, ConnectionInfo connectio
108 this.sqlTools = sqlTools;107 this.sqlTools = sqlTools;
109 this.breakPointRecord = breakPointRecord;108 this.breakPointRecord = breakPointRecord;
110 this.replayedOffsets = breakPointRecord.getReplayedOffset();109 this.replayedOffsets = breakPointRecord.getReplayedOffset();
111- this.isBpSwitch = breakPointRecord.getIsBpSwitch();
112 this.isTransaction = false;110 this.isTransaction = false;
113 }111 }
114 112 
@@ -143,9 +141,7 @@ public void run() {
143 successCount++;141 successCount++;
144 threadSinkRecordObject = sinkRecordObject;142 threadSinkRecordObject = sinkRecordObject;
145 replayedOffsets.offer(sinkRecordObject.getKafkaOffset());143 replayedOffsets.offer(sinkRecordObject.getKafkaOffset());
146- if (isBpSwitch) {144+ savedBreakPointInfo(sinkRecordObject, false);
147- savedBreakPointInfo(sinkRecordObject, false);
148- }
149 } catch (CommunicationsException exp) {145 } catch (CommunicationsException exp) {
150 updateConnectionAndExecuteSql(sql, sinkRecordObject);146 updateConnectionAndExecuteSql(sql, sinkRecordObject);
151 } catch (SQLException exp) {147 } catch (SQLException exp) {
@@ -377,9 +373,7 @@ private void updateConnectionAndExecuteSql(String sql, SinkRecordObject sinkReco
377 connection = connectionInfo.createMysqlConnection();373 connection = connectionInfo.createMysqlConnection();
378 statement = connection.createStatement();374 statement = connection.createStatement();
379 statement.executeUpdate(sql);375 statement.executeUpdate(sql);
380- if (isBpSwitch) {376+ savedBreakPointInfo(sinkRecordObject, false);
381- savedBreakPointInfo(sinkRecordObject, false);
382- }
383 successCount++;377 successCount++;
384 } catch (SQLException exp) {378 } catch (SQLException exp) {
385 if (!connectionInfo.checkConnectionStatus(connection)) {379 if (!connectionInfo.checkConnectionStatus(connection)) {
@@ -132,7 +132,6 @@ public class SinkConnectorConfig extends AbstractConfig {
132 .define(FILE_SIZE_LIMIT, ConfigDef.Type.STRING, "10", ConfigDef.Importance.HIGH, "file size limit")132 .define(FILE_SIZE_LIMIT, ConfigDef.Type.STRING, "10", ConfigDef.Importance.HIGH, "file size limit")
133 .define(BP_BOOTSTRAP_SERVERS, ConfigDef.Type.STRING, "localhost:9092",133 .define(BP_BOOTSTRAP_SERVERS, ConfigDef.Type.STRING, "localhost:9092",
134 ConfigDef.Importance.HIGH, "breakpoint kafka server")134 ConfigDef.Importance.HIGH, "breakpoint kafka server")
135- .define(BP_SWITCH, ConfigDef.Type.STRING, "false", ConfigDef.Importance.HIGH, "breakpoint switch")
136 .define(BP_TOPIC, ConfigDef.Type.STRING, "bp_topic", ConfigDef.Importance.HIGH, "breakpoint topic")135 .define(BP_TOPIC, ConfigDef.Type.STRING, "bp_topic", ConfigDef.Importance.HIGH, "breakpoint topic")
137 .define(BP_ATTEMPTS, ConfigDef.Type.STRING, "3", ConfigDef.Importance.HIGH, "breakpoint attempts")136 .define(BP_ATTEMPTS, ConfigDef.Type.STRING, "3", ConfigDef.Importance.HIGH, "breakpoint attempts")
138 .define(BP_QUEUE_MAX_SIZE, ConfigDef.Type.STRING, "3000",137 .define(BP_QUEUE_MAX_SIZE, ConfigDef.Type.STRING, "3000",
@@ -163,7 +162,6 @@ public class SinkConnectorConfig extends AbstractConfig {
163 /**162 /**
164 * breakpoint config163 * breakpoint config
165 */164 */
166- private boolean isBpSwitch = false;
167 private String bpTopic = "bp_topic";165 private String bpTopic = "bp_topic";
168 private String bootstrapServers = "localhost:9092";166 private String bootstrapServers = "localhost:9092";
169 private int bpMaxRetries = 3;167 private int bpMaxRetries = 3;
@@ -204,15 +202,6 @@ protected void logAll(Map<?, ?> props, String name) {
204 LOGGER.info(sb.toString());202 LOGGER.info(sb.toString());
205 }203 }
206 204 
207- /**
208- * Gets Breakpoint Switch
209- *
210- * @return the value of bpSwitch
211- */
212- public Boolean getBpSwitch() {
213- return isBpSwitch;
214- }
215- 
216 /**205 /**
217 * Gets TOPICS.206 * Gets TOPICS.
218 *207 *
@@ -503,8 +492,5 @@ private void rectifyParameter() {
503 if (isNumberValid(FILE_SIZE_LIMIT, fileSizeLimit)) {492 if (isNumberValid(FILE_SIZE_LIMIT, fileSizeLimit)) {
504 fileSizeLimit = Integer.parseInt(getString(FILE_SIZE_LIMIT));493 fileSizeLimit = Integer.parseInt(getString(FILE_SIZE_LIMIT));
505 }494 }
506- if (isBooleanValid(BP_SWITCH)) {
507- isBpSwitch = Boolean.parseBoolean(getString(BP_SWITCH));
508- }
509 }495 }
510}496}
@@ -23,7 +23,6 @@
23import java.util.UUID;23import java.util.UUID;
24import java.util.concurrent.BlockingQueue;24import java.util.concurrent.BlockingQueue;
25import java.util.concurrent.ExecutionException;25import java.util.concurrent.ExecutionException;
26-import java.util.concurrent.Future;
27import java.util.concurrent.LinkedBlockingQueue;26import java.util.concurrent.LinkedBlockingQueue;
28import java.util.concurrent.PriorityBlockingQueue;27import java.util.concurrent.PriorityBlockingQueue;
29import java.util.concurrent.ThreadPoolExecutor;28import java.util.concurrent.ThreadPoolExecutor;
@@ -45,7 +44,6 @@
45import org.apache.kafka.clients.producer.KafkaProducer;44import org.apache.kafka.clients.producer.KafkaProducer;
46import org.apache.kafka.clients.producer.ProducerConfig;45import org.apache.kafka.clients.producer.ProducerConfig;
47import org.apache.kafka.clients.producer.ProducerRecord;46import org.apache.kafka.clients.producer.ProducerRecord;
48-import org.apache.kafka.clients.producer.RecordMetadata;
49import org.apache.kafka.common.KafkaFuture;47import org.apache.kafka.common.KafkaFuture;
50import org.apache.kafka.common.Node;48import org.apache.kafka.common.Node;
51import org.apache.kafka.common.TopicPartition;49import org.apache.kafka.common.TopicPartition;
@@ -192,7 +190,6 @@ public class BreakPointRecord {
192 private Long totalMessageCount = 0L;190 private Long totalMessageCount = 0L;
193 private PriorityBlockingQueue<Long> replayedOffsets;191 private PriorityBlockingQueue<Long> replayedOffsets;
194 private boolean isGetBp;192 private boolean isGetBp;
195- private boolean isBpSwitch;
196 private Long breakpointEndOffset = UNLIMITED_VALUE;193 private Long breakpointEndOffset = UNLIMITED_VALUE;
197 194 
198 /**195 /**
@@ -236,24 +233,6 @@ public BreakPointRecord(Configuration config) {
236 }233 }
237 }234 }
238 235 
239- /**
240- * Sets the breakpoint switch
241- *
242- * @param isBpSwitch the record breakpoint switch
243- */
244- public void setIsBpSwitch(boolean isBpSwitch) {
245- this.isBpSwitch = isBpSwitch;
246- }
247- 
248- /**
249- * Gets the breakpoint switch
250- *
251- * @return boolean record breakpoint switch
252- */
253- public boolean getIsBpSwitch() {
254- return isBpSwitch;
255- }
256- 
257 /**236 /**
258 * get the breakpoint to delete last offset237 * get the breakpoint to delete last offset
259 *238 *
@@ -518,6 +497,7 @@ public void storeRecord(BreakPointObject record, Boolean isTransaction, boolean
518 497 
519 private void storeToKafkaThread() {498 private void storeToKafkaThread() {
520 threadPool.execute(() -> {499 threadPool.execute(() -> {
500+ Thread.currentThread().setName("send2kafka-thread");
521 List<BreakPointInfo> toStoreList = new ArrayList<>();501 List<BreakPointInfo> toStoreList = new ArrayList<>();
522 while (true) {502 while (true) {
523 try {503 try {
@@ -544,24 +524,8 @@ private void storeToKafkaThread() {
544 public void storeRecordToKafka(String key, String value) {524 public void storeRecordToKafka(String key, String value) {
545 ProducerRecord<String, String> produced = new ProducerRecord<>(bpRecordTopicName, PARTITION,525 ProducerRecord<String, String> produced = new ProducerRecord<>(bpRecordTopicName, PARTITION,
546 key, value);526 key, value);
547- Future<RecordMetadata> future = this.producer.send(produced);527+ this.producer.send(produced);
548- // Flush and then wait ...528+ totalMessageCount++;
549- RecordMetadata metadata = null; // block forever since we have to be sure this gets recorded
550- try {
551- metadata = future.get();
552- if (metadata != null) {
553- LOGGER.debug("Stored record in topic '{}' partition {} at offset {} ",
554- metadata.topic(), metadata.partition(), metadata.offset());
555- }
556- totalMessageCount++;
557- }
558- catch (InterruptedException e) {
559- LOGGER.trace("Interrupted before record was written into kafka file");
560- Thread.currentThread().interrupt();
561- }
562- catch (ExecutionException e) {
563- LOGGER.error("");
564- }
565 }529 }
566 530 
567 /**531 /**
@@ -810,6 +774,7 @@ public void run() {
810 774 
811 private void deleteBpBySizeTask() {775 private void deleteBpBySizeTask() {
812 threadPool.execute(() -> {776 threadPool.execute(() -> {
777+ Thread.currentThread().setName("delete-breakpoint-thread");
813 while (true) {778 while (true) {
814 // Additional thread monitor the queue size and delete them if the limit is exceeded779 // Additional thread monitor the queue size and delete them if the limit is exceeded
815 if (totalMessageCount >= bpQueueSizeLimit) {780 if (totalMessageCount >= bpQueueSizeLimit) {