回退版本
This commit is contained in:
parent
081c3cbbf3
commit
e406ecb782
|
|
@ -21,8 +21,8 @@ public class DataSenderFactory {
|
|||
|
||||
private static Logger logger = Logger.getLogger(DataSenderFactory.class.getName());
|
||||
|
||||
private static Set<SocketAddress> socketAddresses = new HashSet<SocketAddress>();
|
||||
private static Set<SocketAddress> unUsedSocketAddresses = new HashSet<SocketAddress>();
|
||||
private static List<SocketAddress> socketAddresses = new ArrayList<SocketAddress>();
|
||||
private static List<SocketAddress> unUsedSocketAddresses = new ArrayList<SocketAddress>();
|
||||
private static List<DataSender> availableSenders = new ArrayList<DataSender>();
|
||||
private static Object lock = new Object();
|
||||
|
||||
|
|
@ -31,13 +31,17 @@ public class DataSenderFactory {
|
|||
if (StringUtil.isEmpty(Config.Sender.SERVERS_ADDR)) {
|
||||
throw new IllegalArgumentException("Collection service configuration error.");
|
||||
}
|
||||
|
||||
//过滤重复地址
|
||||
Set<SocketAddress> tmpSocktAddress = new HashSet<SocketAddress>();
|
||||
for (String serverConfig : Config.Sender.SERVERS_ADDR.split(";")) {
|
||||
String[] server = serverConfig.split(":");
|
||||
if (server.length != 2)
|
||||
throw new IllegalArgumentException("Collection service configuration error.");
|
||||
socketAddresses.add(new InetSocketAddress(server[0], Integer.valueOf(server[1])));
|
||||
tmpSocktAddress.add(new InetSocketAddress(server[0], Integer.valueOf(server[1])));
|
||||
}
|
||||
|
||||
socketAddresses.addAll(tmpSocktAddress);
|
||||
|
||||
} catch (Exception e) {
|
||||
logger.log(Level.ALL, "Collection service configuration error.", e);
|
||||
System.exit(-1);
|
||||
|
|
@ -72,18 +76,18 @@ public class DataSenderFactory {
|
|||
// 初始化DataSender
|
||||
List<SocketAddress> usedSocketAddress = new ArrayList<SocketAddress>();
|
||||
|
||||
for (SocketAddress socketAddress : socketAddresses) {
|
||||
if (availableSenders.size() >= availableSize) {
|
||||
break;
|
||||
}
|
||||
int index;
|
||||
while (availableSenders.size() < availableSize) {
|
||||
// 随机获取服务器地址
|
||||
index = ThreadLocalRandom.current().nextInt(socketAddresses.size());
|
||||
try {
|
||||
availableSenders.add(new DataSender(socketAddress));
|
||||
usedSocketAddress.add(socketAddress);
|
||||
availableSenders.add(new DataSender(socketAddresses.get(index)));
|
||||
usedSocketAddress.add(socketAddresses.get(index));
|
||||
} catch (IOException e) {
|
||||
unUsedSocketAddresses.add(socketAddress);
|
||||
unUsedSocketAddresses.add(socketAddresses.get(index));
|
||||
}
|
||||
}
|
||||
unUsedSocketAddresses = new HashSet<SocketAddress>(socketAddresses);
|
||||
unUsedSocketAddresses = new ArrayList<SocketAddress>(socketAddresses);
|
||||
unUsedSocketAddresses.removeAll(usedSocketAddress);
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Reference in New Issue