已合并
【bugfix】优化增量迁移性能 #170
lvlintao666创建于 2023年9月18日
【bugfix】优化增量迁移性能 #170
已合并
从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 | * txn the replay transaction | 249 | * 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 config | 163 | * 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 | - * 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 @@ | |||
| 23 | import java.util.UUID; | 23 | import java.util.UUID; |
| 24 | import java.util.concurrent.BlockingQueue; | 24 | import java.util.concurrent.BlockingQueue; |
| 25 | import java.util.concurrent.ExecutionException; | 25 | import java.util.concurrent.ExecutionException; |
| 26 | -import java.util.concurrent.Future; | ||
| 27 | import java.util.concurrent.LinkedBlockingQueue; | 26 | import java.util.concurrent.LinkedBlockingQueue; |
| 28 | import java.util.concurrent.PriorityBlockingQueue; | 27 | import java.util.concurrent.PriorityBlockingQueue; |
| 29 | import java.util.concurrent.ThreadPoolExecutor; | 28 | import java.util.concurrent.ThreadPoolExecutor; |
| @@ -45,7 +44,6 @@ | |||
| 45 | import org.apache.kafka.clients.producer.KafkaProducer; | 44 | import org.apache.kafka.clients.producer.KafkaProducer; |
| 46 | import org.apache.kafka.clients.producer.ProducerConfig; | 45 | import org.apache.kafka.clients.producer.ProducerConfig; |
| 47 | import org.apache.kafka.clients.producer.ProducerRecord; | 46 | import org.apache.kafka.clients.producer.ProducerRecord; |
| 48 | -import org.apache.kafka.clients.producer.RecordMetadata; | ||
| 49 | import org.apache.kafka.common.KafkaFuture; | 47 | import org.apache.kafka.common.KafkaFuture; |
| 50 | import org.apache.kafka.common.Node; | 48 | import org.apache.kafka.common.Node; |
| 51 | import org.apache.kafka.common.TopicPartition; | 49 | import 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 | - * 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 | - * boolean record breakpoint switch | ||
| 252 | - */ | ||
| 253 | - public boolean getIsBpSwitch() { | ||
| 254 | - return isBpSwitch; | ||
| 255 | - } | ||
| 256 | - | ||
| 257 | /** | 236 | /** |
| 258 | * get the breakpoint to delete last offset | 237 | * 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 exceeded | 779 | // Additional thread monitor the queue size and delete them if the limit is exceeded |
| 815 | if (totalMessageCount >= bpQueueSizeLimit) { | 780 | if (totalMessageCount >= bpQueueSizeLimit) { |