From de95a6a1284d13ddbc8b48eb090b77b8d3e12c90 Mon Sep 17 00:00:00 2001 From: Jonathan Schneider Date: Sat, 9 Oct 2021 00:33:37 -0700 Subject: [PATCH] refactor: Use Java standard library instead of Guava (#3708) Co-authored-by: Moderne Co-authored-by: Moderne --- .../canal/client/kafka/KafkaOffsetCanalConnector.java | 9 ++++----- .../canal/client/rabbitmq/RabbitMQCanalConnector.java | 3 ++- .../canal/client/rocketmq/RocketMQCanalConnector.java | 3 ++- .../canal/connector/core/producer/MQMessageUtils.java | 10 +++++----- .../rabbitmq/consumer/CanalRabbitMQConsumer.java | 4 ++-- .../rocketmq/consumer/CanalRocketMQConsumer.java | 4 ++-- .../monitor/ManagerInstanceConfigMonitor.java | 7 ++++--- .../otter/canal/meta/FileMixedMetaManager.java | 2 +- .../alibaba/otter/canal/meta/MemoryMetaManager.java | 11 ++++------- .../inbound/mysql/tsdb/MemoryTableMeta_DDL_Test.java | 6 +++--- .../mysql/tsdb/MemoryTableMeta_Random_DDL_Test.java | 4 ++-- .../com/alibaba/otter/canal/protocol/FlatMessage.java | 5 ++--- 12 files changed, 33 insertions(+), 35 deletions(-) diff --git a/client/src/main/java/com/alibaba/otter/canal/client/kafka/KafkaOffsetCanalConnector.java b/client/src/main/java/com/alibaba/otter/canal/client/kafka/KafkaOffsetCanalConnector.java index 2bfa35f9..bd2f775d 100644 --- a/client/src/main/java/com/alibaba/otter/canal/client/kafka/KafkaOffsetCanalConnector.java +++ b/client/src/main/java/com/alibaba/otter/canal/client/kafka/KafkaOffsetCanalConnector.java @@ -6,7 +6,6 @@ import com.alibaba.otter.canal.client.kafka.protocol.KafkaMessage; import com.alibaba.otter.canal.protocol.FlatMessage; import com.alibaba.otter.canal.protocol.Message; import com.alibaba.otter.canal.protocol.exception.CanalClientException; -import com.google.common.collect.Lists; import org.apache.commons.lang3.StringUtils; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; @@ -43,7 +42,7 @@ public class KafkaOffsetCanalConnector extends KafkaCanalConnector { public List getListWithoutAck(Long timeout, TimeUnit unit, long offset) throws CanalClientException { waitClientRunning(); if (!running) { - return Lists.newArrayList(); + return new ArrayList<>(); } if (offset > -1) { @@ -61,7 +60,7 @@ public class KafkaOffsetCanalConnector extends KafkaCanalConnector { } return messages; } - return Lists.newArrayList(); + return new ArrayList<>(); } /** @@ -76,7 +75,7 @@ public class KafkaOffsetCanalConnector extends KafkaCanalConnector { public List getFlatListWithoutAck(Long timeout, TimeUnit unit, long offset) throws CanalClientException { waitClientRunning(); if (!running) { - return Lists.newArrayList(); + return new ArrayList<>(); } if (offset > -1) { @@ -96,7 +95,7 @@ public class KafkaOffsetCanalConnector extends KafkaCanalConnector { return flatMessages; } - return Lists.newArrayList(); + return new ArrayList<>(); } /** diff --git a/client/src/main/java/com/alibaba/otter/canal/client/rabbitmq/RabbitMQCanalConnector.java b/client/src/main/java/com/alibaba/otter/canal/client/rabbitmq/RabbitMQCanalConnector.java index bf43a4a8..b0bbc97c 100644 --- a/client/src/main/java/com/alibaba/otter/canal/client/rabbitmq/RabbitMQCanalConnector.java +++ b/client/src/main/java/com/alibaba/otter/canal/client/rabbitmq/RabbitMQCanalConnector.java @@ -13,6 +13,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.io.IOException; +import java.util.ArrayList; import java.util.List; import java.util.concurrent.BlockingQueue; import java.util.concurrent.LinkedBlockingQueue; @@ -239,7 +240,7 @@ public class RabbitMQCanalConnector implements CanalMQConnector { if (logger.isDebugEnabled()) { logger.debug("Get Message: {}", new String(messageData)); } - List messageList = Lists.newArrayList(); + List messageList = new ArrayList<>(); if (!flatMessage) { Message message = CanalMessageDeserializer.deserializer(messageData); messageList.add(message); diff --git a/client/src/main/java/com/alibaba/otter/canal/client/rocketmq/RocketMQCanalConnector.java b/client/src/main/java/com/alibaba/otter/canal/client/rocketmq/RocketMQCanalConnector.java index 71bc5537..ac0a6191 100644 --- a/client/src/main/java/com/alibaba/otter/canal/client/rocketmq/RocketMQCanalConnector.java +++ b/client/src/main/java/com/alibaba/otter/canal/client/rocketmq/RocketMQCanalConnector.java @@ -1,5 +1,6 @@ package com.alibaba.otter.canal.client.rocketmq; +import java.util.ArrayList; import java.util.List; import java.util.concurrent.BlockingQueue; import java.util.concurrent.LinkedBlockingQueue; @@ -165,7 +166,7 @@ public class RocketMQCanalConnector implements CanalMQConnector { if (logger.isDebugEnabled()) { logger.debug("Get Message: {}", messageExts); } - List messageList = Lists.newArrayList(); + List messageList = new ArrayList<>(); for (MessageExt messageExt : messageExts) { byte[] data = messageExt.getBody(); if (data != null) { diff --git a/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/producer/MQMessageUtils.java b/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/producer/MQMessageUtils.java index f53a05fc..79f50a81 100644 --- a/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/producer/MQMessageUtils.java +++ b/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/producer/MQMessageUtils.java @@ -36,7 +36,7 @@ public class MQMessageUtils { private static Map> partitionDatas = MigrateMap.makeComputingMap(CacheBuilder.newBuilder() .softValues(), pkHashConfigs -> { - List datas = Lists.newArrayList(); + List datas = new ArrayList<>(); String[] pkHashConfigArray = StringUtils.split(StringUtils.replace(pkHashConfigs, ",", @@ -75,7 +75,7 @@ public class MQMessageUtils { private static Map> dynamicTopicDatas = MigrateMap.makeComputingMap(CacheBuilder.newBuilder() .softValues(), pkHashConfigs -> { - List datas = Lists.newArrayList(); + List datas = new ArrayList<>(); String[] dynamicTopicArray = StringUtils.split(StringUtils.replace(pkHashConfigs, ",", ";"), @@ -102,7 +102,7 @@ public class MQMessageUtils { private static Map> topicPartitionDatas = MigrateMap.makeComputingMap(CacheBuilder.newBuilder() .softValues(), tPConfigs -> { - List datas = Lists.newArrayList(); + List datas = new ArrayList<>(); String[] tPArray = StringUtils.split(StringUtils.replace(tPConfigs, ",", ";"), @@ -258,7 +258,7 @@ public class MQMessageUtils { List[] partitionEntries = new List[partitionsNum]; for (int i = 0; i < partitionsNum; i++) { // 注意一下并发 - partitionEntries[i] = Collections.synchronizedList(Lists.newArrayList()); + partitionEntries[i] = Collections.synchronizedList(new ArrayList<>()); } for (EntryRowData data : datas) { @@ -681,7 +681,7 @@ public class MQMessageUtils { public boolean autoPkHash = false; public boolean tableHash = false; - public List pkNames = Lists.newArrayList(); + public List pkNames = new ArrayList<>(); } public static class DynamicTopicData { diff --git a/connector/rabbitmq-connector/src/main/java/com/alibaba/otter/canal/connector/rabbitmq/consumer/CanalRabbitMQConsumer.java b/connector/rabbitmq-connector/src/main/java/com/alibaba/otter/canal/connector/rabbitmq/consumer/CanalRabbitMQConsumer.java index b29836b5..ebd3fa03 100644 --- a/connector/rabbitmq-connector/src/main/java/com/alibaba/otter/canal/connector/rabbitmq/consumer/CanalRabbitMQConsumer.java +++ b/connector/rabbitmq-connector/src/main/java/com/alibaba/otter/canal/connector/rabbitmq/consumer/CanalRabbitMQConsumer.java @@ -1,6 +1,7 @@ package com.alibaba.otter.canal.connector.rabbitmq.consumer; import java.io.IOException; +import java.util.ArrayList; import java.util.List; import java.util.Properties; import java.util.concurrent.BlockingQueue; @@ -23,7 +24,6 @@ import com.alibaba.otter.canal.connector.rabbitmq.config.RabbitMQConstants; import com.alibaba.otter.canal.connector.rabbitmq.producer.AliyunCredentialsProvider; import com.alibaba.otter.canal.protocol.Message; import com.alibaba.otter.canal.protocol.exception.CanalClientException; -import com.google.common.collect.Lists; import com.rabbitmq.client.AMQP; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; @@ -137,7 +137,7 @@ public class CanalRabbitMQConsumer implements CanalMsgConsumer { if (logger.isDebugEnabled()) { logger.debug("Get Message: {}", new String(messageData)); } - List messageList = Lists.newArrayList(); + List messageList = new ArrayList<>(); if (!flatMessage) { Message message = CanalMessageSerializerUtil.deserializer(messageData); messageList.addAll(MessageUtil.convert(message)); diff --git a/connector/rocketmq-connector/src/main/java/com/alibaba/otter/canal/connector/rocketmq/consumer/CanalRocketMQConsumer.java b/connector/rocketmq-connector/src/main/java/com/alibaba/otter/canal/connector/rocketmq/consumer/CanalRocketMQConsumer.java index 3793ca73..e3f969cb 100644 --- a/connector/rocketmq-connector/src/main/java/com/alibaba/otter/canal/connector/rocketmq/consumer/CanalRocketMQConsumer.java +++ b/connector/rocketmq-connector/src/main/java/com/alibaba/otter/canal/connector/rocketmq/consumer/CanalRocketMQConsumer.java @@ -1,5 +1,6 @@ package com.alibaba.otter.canal.connector.rocketmq.consumer; +import java.util.ArrayList; import java.util.List; import java.util.Properties; import java.util.concurrent.BlockingQueue; @@ -30,7 +31,6 @@ import com.alibaba.otter.canal.connector.core.util.MessageUtil; import com.alibaba.otter.canal.connector.rocketmq.config.RocketMQConstants; import com.alibaba.otter.canal.protocol.Message; import com.alibaba.otter.canal.protocol.exception.CanalClientException; -import com.google.common.collect.Lists; /** * RocketMQ consumer SPI 实现 @@ -142,7 +142,7 @@ public class CanalRocketMQConsumer implements CanalMsgConsumer { if (logger.isDebugEnabled()) { logger.debug("Get Message: {}", messageExts); } - List messageList = Lists.newArrayList(); + List messageList = new ArrayList<>(); for (MessageExt messageExt : messageExts) { byte[] data = messageExt.getBody(); if (data != null) { diff --git a/deployer/src/main/java/com/alibaba/otter/canal/deployer/monitor/ManagerInstanceConfigMonitor.java b/deployer/src/main/java/com/alibaba/otter/canal/deployer/monitor/ManagerInstanceConfigMonitor.java index 8eaab861..cacdf425 100644 --- a/deployer/src/main/java/com/alibaba/otter/canal/deployer/monitor/ManagerInstanceConfigMonitor.java +++ b/deployer/src/main/java/com/alibaba/otter/canal/deployer/monitor/ManagerInstanceConfigMonitor.java @@ -1,5 +1,6 @@ package com.alibaba.otter.canal.deployer.monitor; +import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.concurrent.Executors; @@ -77,9 +78,9 @@ public class ManagerInstanceConfigMonitor extends AbstractCanalLifeCycle impleme } final List is = Lists.newArrayList(StringUtils.split(instances, ',')); - List start = Lists.newArrayList(); - List stop = Lists.newArrayList(); - List restart = Lists.newArrayList(); + List start = new ArrayList<>(); + List stop = new ArrayList<>(); + List restart = new ArrayList<>(); for (String instance : is) { if (!configs.containsKey(instance)) { PlainCanal newPlainCanal = configClient.findInstance(instance, null); diff --git a/meta/src/main/java/com/alibaba/otter/canal/meta/FileMixedMetaManager.java b/meta/src/main/java/com/alibaba/otter/canal/meta/FileMixedMetaManager.java index c0fefef3..53d92701 100644 --- a/meta/src/main/java/com/alibaba/otter/canal/meta/FileMixedMetaManager.java +++ b/meta/src/main/java/com/alibaba/otter/canal/meta/FileMixedMetaManager.java @@ -194,7 +194,7 @@ public class FileMixedMetaManager extends MemoryMetaManager implements CanalMeta synchronized (destination.intern()) { // 基于destination控制一下并发更新 data.setDestination(destination); - List clientDatas = Lists.newArrayList(); + List clientDatas = new ArrayList<>(); List clientIdentitys = destinations.get(destination); for (ClientIdentity clientIdentity : clientIdentitys) { FileMetaClientIdentityData clientData = new FileMetaClientIdentityData(); diff --git a/meta/src/main/java/com/alibaba/otter/canal/meta/MemoryMetaManager.java b/meta/src/main/java/com/alibaba/otter/canal/meta/MemoryMetaManager.java index 94ac0b7c..0883e447 100644 --- a/meta/src/main/java/com/alibaba/otter/canal/meta/MemoryMetaManager.java +++ b/meta/src/main/java/com/alibaba/otter/canal/meta/MemoryMetaManager.java @@ -1,9 +1,6 @@ package com.alibaba.otter.canal.meta; -import java.util.Collections; -import java.util.List; -import java.util.Map; -import java.util.Set; +import java.util.*; import java.util.concurrent.atomic.AtomicLong; import com.alibaba.otter.canal.common.AbstractCanalLifeCycle; @@ -35,7 +32,7 @@ public class MemoryMetaManager extends AbstractCanalLifeCycle implements CanalMe cursors = new MapMaker().makeMap(); - destinations = MigrateMap.makeComputingMap(destination -> Lists.newArrayList()); + destinations = MigrateMap.makeComputingMap(destination -> new ArrayList<>()); } public void stop() { @@ -186,8 +183,8 @@ public class MemoryMetaManager extends AbstractCanalLifeCycle implements CanalMe public synchronized Map listAllPositionRange() { Set batchIdSets = batches.keySet(); - List batchIds = Lists.newArrayList(batchIdSets); - Collections.sort(Lists.newArrayList(batchIds)); + List batchIds = new ArrayList<>(batchIdSets); + Collections.sort(new ArrayList<>(batchIds)); return Maps.newHashMap(batches); } diff --git a/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/MemoryTableMeta_DDL_Test.java b/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/MemoryTableMeta_DDL_Test.java index eb2d575d..6630666b 100644 --- a/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/MemoryTableMeta_DDL_Test.java +++ b/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/MemoryTableMeta_DDL_Test.java @@ -3,6 +3,7 @@ package com.alibaba.otter.canal.parse.inbound.mysql.tsdb; import java.io.File; import java.io.FileInputStream; import java.net.URL; +import java.util.ArrayList; import java.util.List; import org.apache.commons.io.IOUtils; @@ -15,7 +16,6 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; import com.alibaba.druid.sql.repository.Schema; import com.alibaba.otter.canal.parse.inbound.TableMeta; -import com.google.common.collect.Lists; /** * @author agapple 2017年8月1日 下午7:15:54 @@ -79,7 +79,7 @@ public class MemoryTableMeta_DDL_Test { String sql = StringUtils.join(IOUtils.readLines(new FileInputStream(create)), "\n"); memoryTableMeta.apply(null, "test", sql, null); - List tableNames = Lists.newArrayList(); + List tableNames = new ArrayList<>(); for (Schema schema : memoryTableMeta.getRepository().getSchemas()) { tableNames.addAll(schema.showTables()); } @@ -99,7 +99,7 @@ public class MemoryTableMeta_DDL_Test { String sql = StringUtils.join(IOUtils.readLines(new FileInputStream(create)), "\n"); memoryTableMeta.apply(null, "test", sql, null); - List tableNames = Lists.newArrayList(); + List tableNames = new ArrayList<>(); for (Schema schema : memoryTableMeta.getRepository().getSchemas()) { tableNames.addAll(schema.showTables()); } diff --git a/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/MemoryTableMeta_Random_DDL_Test.java b/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/MemoryTableMeta_Random_DDL_Test.java index bd2cf95c..f11eb27b 100644 --- a/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/MemoryTableMeta_Random_DDL_Test.java +++ b/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/MemoryTableMeta_Random_DDL_Test.java @@ -3,6 +3,7 @@ package com.alibaba.otter.canal.parse.inbound.mysql.tsdb; import java.io.File; import java.io.FileInputStream; import java.net.URL; +import java.util.ArrayList; import java.util.List; import org.apache.commons.io.IOUtils; @@ -15,7 +16,6 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; import com.alibaba.druid.sql.repository.Schema; import com.alibaba.otter.canal.parse.inbound.TableMeta; -import com.google.common.collect.Lists; /** * @author agapple 2017年8月1日 下午7:15:54 @@ -69,7 +69,7 @@ public class MemoryTableMeta_Random_DDL_Test { } private void compareTableMeta(int num, MemoryTableMeta source, MemoryTableMeta target) { - List tableNames = Lists.newArrayList(); + List tableNames = new ArrayList<>(); for (Schema schema : source.getRepository().getSchemas()) { tableNames.addAll(schema.showTables()); } diff --git a/protocol/src/main/java/com/alibaba/otter/canal/protocol/FlatMessage.java b/protocol/src/main/java/com/alibaba/otter/canal/protocol/FlatMessage.java index 2b958b1c..0957b759 100644 --- a/protocol/src/main/java/com/alibaba/otter/canal/protocol/FlatMessage.java +++ b/protocol/src/main/java/com/alibaba/otter/canal/protocol/FlatMessage.java @@ -1,11 +1,10 @@ package com.alibaba.otter.canal.protocol; import java.io.Serializable; +import java.util.ArrayList; import java.util.List; import java.util.Map; -import com.google.common.collect.Lists; - /** * @author machengyuan 2018-9-13 下午10:31:14 * @version 1.0.0 @@ -66,7 +65,7 @@ public class FlatMessage implements Serializable { public void addPkName(String pkName) { if (this.pkNames == null) { - this.pkNames = Lists.newArrayList(); + this.pkNames = new ArrayList<>(); } this.pkNames.add(pkName); }