From bb203bdc7163a05b4956ac7327c2f176758b3a3f Mon Sep 17 00:00:00 2001 From: YuQing <384681@qq.com> Date: Wed, 13 Dec 2023 10:06:39 +0800 Subject: [PATCH] connect to storage server failover with multi IPs --- .gitignore | 1 + HISTORY | 3 +- fdfs_client.conf | 8 +- .../org/csource/fastdfs/ClientGlobal.java | 99 ++++++++++++- .../java/org/csource/fastdfs/ProtoCommon.java | 41 ++++-- .../csource/fastdfs/StorageAddressMap.java | 62 ++++++++ .../org/csource/fastdfs/StorageClient.java | 1 - .../org/csource/fastdfs/StorageServer.java | 139 ++++++++++++++---- .../org/csource/fastdfs/TrackerClient.java | 79 +++++++++- .../org/csource/fastdfs/TrackerServer.java | 11 +- .../fastdfs-client.properties.sample | 8 + src/main/resources/fdfs_client.conf.sample | 8 + 12 files changed, 404 insertions(+), 56 deletions(-) create mode 100644 src/main/java/org/csource/fastdfs/StorageAddressMap.java diff --git a/.gitignore b/.gitignore index 76d9377..12cfc78 100644 --- a/.gitignore +++ b/.gitignore @@ -15,3 +15,4 @@ target *.conf *.PNG *.class +*.swp diff --git a/HISTORY b/HISTORY index 64bd5a0..11fde87 100644 --- a/HISTORY +++ b/HISTORY @@ -1,7 +1,8 @@ -Version 1.31 2023-12-07 +Version 1.31 2023-12-13 * adapt to FastDFS server V6.11 for IPv6 you must upgrade your FastDFS server to V6.11 or higher version + * connect to storage server failover with multi IPs Version 1.30 2023-01-29 * support tracker server fail over diff --git a/fdfs_client.conf b/fdfs_client.conf index 01debe5..b58567e 100644 --- a/fdfs_client.conf +++ b/fdfs_client.conf @@ -16,9 +16,13 @@ tracker_server = 10.0.11.247:22122 tracker_server = 10.0.11.248:22122 tracker_server = 10.0.11.249:22122 +# connect which ip address first for multi IPs of a storage server, value list: +## tracker: connect to the ip address return by tracker server first +## last-connected: connect to the ip address last connected first +# default value is tracker +connect_first_by = tracker + connection_pool.enabled = true connection_pool.max_count_per_entry = 500 connection_pool.max_idle_time = 3600 connection_pool.max_wait_time_in_ms = 1000 - -server_ipv6.enabled = false \ No newline at end of file diff --git a/src/main/java/org/csource/fastdfs/ClientGlobal.java b/src/main/java/org/csource/fastdfs/ClientGlobal.java index 0cccf01..d454e22 100644 --- a/src/main/java/org/csource/fastdfs/ClientGlobal.java +++ b/src/main/java/org/csource/fastdfs/ClientGlobal.java @@ -42,7 +42,7 @@ public class ClientGlobal { public static final String PROP_KEY_HTTP_SECRET_KEY = "fastdfs.http_secret_key"; public static final String PROP_KEY_HTTP_TRACKER_HTTP_PORT = "fastdfs.http_tracker_http_port"; public static final String PROP_KEY_TRACKER_SERVERS = "fastdfs.tracker_servers"; - + public static final String PROP_KEY_CONNECT_FIRST_BY = "fastdfs.connect_first_by"; public static final String PROP_KEY_CONNECTION_POOL_ENABLED = "fastdfs.connection_pool.enabled"; public static final String PROP_KEY_CONNECTION_POOL_MAX_COUNT_PER_ENTRY = "fastdfs.connection_pool.max_count_per_entry"; @@ -56,6 +56,9 @@ public class ClientGlobal { public static final String DEFAULT_HTTP_SECRET_KEY = "FastDFS1234567890"; public static final int DEFAULT_HTTP_TRACKER_HTTP_PORT = 80; + public static final int CONNECT_FIRST_BY_TRACKER = 0; + public static final int CONNECT_FIRST_BY_LAST_CONNECTED = 1; + public static final boolean DEFAULT_CONNECTION_POOL_ENABLED = true; public static final int DEFAULT_CONNECTION_POOL_MAX_COUNT_PER_ENTRY = 100; public static final int DEFAULT_CONNECTION_POOL_MAX_IDLE_TIME = 3600 ;//second @@ -67,6 +70,9 @@ public class ClientGlobal { public static boolean g_anti_steal_token = DEFAULT_HTTP_ANTI_STEAL_TOKEN; //if anti-steal token public static String g_secret_key = DEFAULT_HTTP_SECRET_KEY; //generage token secret key public static int g_tracker_http_port = DEFAULT_HTTP_TRACKER_HTTP_PORT; + public static int g_connect_first_by = CONNECT_FIRST_BY_TRACKER; + public static boolean g_multi_storage_ips = false; + public static StorageAddressMap g_storages_address_map; public static boolean g_connection_pool_enabled = DEFAULT_CONNECTION_POOL_ENABLED; public static int g_connection_pool_max_count_per_entry = DEFAULT_CONNECTION_POOL_MAX_COUNT_PER_ENTRY; @@ -78,6 +84,76 @@ public class ClientGlobal { private ClientGlobal() { } + private static void loadStoragsFromTracker() throws IOException, MyException { + TrackerClient tracker = new TrackerClient(); + StringBuilder builder = tracker.fetchStorageIds(); + if (builder.length() == 0) { + return; + } + + System.out.println(builder.toString()); + + boolean without_port = true; + int count = 0; + String[] lines = builder.toString().split("\n"); + String[] ipAddresses = new String[lines.length]; + for (String line : lines) { + String[] cols = line.split(" "); + if (cols.length != 3) { + throw new MyException("invalid line: " + line); + } + + String ipAddrs = cols[2]; + if (ipAddrs.indexOf(',') > 0) { + ipAddresses[count++] = ipAddrs; + } + } + + if (count == 0) { + return; + } + + int startIndex; + if (ipAddresses[0].charAt(0) == '[') { //IPv6 + if ((startIndex=ipAddresses[0].indexOf(']')) < 0) { + throw new MyException("invalid IPv6 address: " + ipAddresses[0]); + } + } else { + startIndex = 0; + } + if (ipAddresses[0].indexOf(':', startIndex) > 0) { + without_port = false; + } + + g_multi_storage_ips = true; + g_storages_address_map = new StorageAddressMap(without_port); + if (without_port) { + for (String ipAddr: ipAddresses) { + if (ipAddr.charAt(0) == '[') { //IPv6 + ipAddr = ipAddr.substring(1, ipAddr.length() - 1); + } + String[] cols = ipAddr.split(","); + g_storages_address_map.puts(cols[0], cols[1]); + } + } else { + for (String ipPort: ipAddresses) { + int colonIndex = ipPort.lastIndexOf(':'); + if (colonIndex < 0) { + throw new MyException("invalid ip and port: " + ipPort); + } + + String ipAddr = ipPort.substring(0, colonIndex); + int port = Integer.parseInt(ipPort.substring(colonIndex + 1)); + + if (ipAddr.charAt(0) == '[') { //IPv6 + ipAddr = ipAddr.substring(1, ipAddr.length() - 1); + } + String[] cols = ipAddr.split(","); + g_storages_address_map.puts(cols[0], cols[1], port); + } + } + } + /** * load global variables * @@ -114,11 +190,11 @@ public class ClientGlobal { InetSocketAddress[] tracker_servers = new InetSocketAddress[szTrackerServers.length]; for (int i = 0; i < szTrackerServers.length; i++) { - if(szTrackerServers[i].contains("[")){ + if (szTrackerServers[i].contains("[")) { parts = new String[2]; parts[0] = szTrackerServers[i].substring(1, szTrackerServers[i].indexOf("]")); parts[1] = szTrackerServers[i].substring(szTrackerServers[i].lastIndexOf(":") + 1); - }else { + } else { parts = szTrackerServers[i].split("\\:", 2); } @@ -130,6 +206,11 @@ public class ClientGlobal { } g_tracker_group = new TrackerGroup(tracker_servers); + String connect_first_by = iniReader.getStrValue("connect_first_by"); + if (connect_first_by != null && connect_first_by.equalsIgnoreCase("last-connected")) { + g_connect_first_by = CONNECT_FIRST_BY_LAST_CONNECTED; + } + g_tracker_http_port = iniReader.getIntValue("http.tracker_http_port", 80); g_anti_steal_token = iniReader.getBoolValue("http.anti_steal_token", false); if (g_anti_steal_token) { @@ -146,6 +227,8 @@ public class ClientGlobal { if (g_connection_pool_max_wait_time_in_ms < 0) { g_connection_pool_max_wait_time_in_ms = DEFAULT_CONNECTION_POOL_MAX_WAIT_TIME_IN_MS; } + + loadStoragsFromTracker(); } /** @@ -183,6 +266,8 @@ public class ClientGlobal { String httpAntiStealTokenConf = props.getProperty(PROP_KEY_HTTP_ANTI_STEAL_TOKEN); String httpSecretKeyConf = props.getProperty(PROP_KEY_HTTP_SECRET_KEY); String httpTrackerHttpPortConf = props.getProperty(PROP_KEY_HTTP_TRACKER_HTTP_PORT); + + String connectFirstBy = props.getProperty(PROP_KEY_CONNECT_FIRST_BY); String poolEnabled = props.getProperty(PROP_KEY_CONNECTION_POOL_ENABLED); String poolMaxCountPerEntry = props.getProperty(PROP_KEY_CONNECTION_POOL_MAX_COUNT_PER_ENTRY); String poolMaxIdleTime = props.getProperty(PROP_KEY_CONNECTION_POOL_MAX_IDLE_TIME); @@ -205,6 +290,10 @@ public class ClientGlobal { if (httpTrackerHttpPortConf != null && httpTrackerHttpPortConf.trim().length() != 0) { g_tracker_http_port = Integer.parseInt(httpTrackerHttpPortConf); } + + if (connectFirstBy != null && connectFirstBy.equalsIgnoreCase("last-connected")) { + g_connect_first_by = CONNECT_FIRST_BY_LAST_CONNECTED; + } if (poolEnabled != null && poolEnabled.trim().length() != 0) { g_connection_pool_enabled = Boolean.parseBoolean(poolEnabled); } @@ -217,6 +306,8 @@ public class ClientGlobal { if (poolMaxWaitTimeInMS != null && poolMaxWaitTimeInMS.trim().length() != 0) { g_connection_pool_max_wait_time_in_ms = Integer.parseInt(poolMaxWaitTimeInMS); } + + loadStoragsFromTracker(); } /** @@ -360,6 +451,8 @@ public class ClientGlobal { + "\n g_anti_steal_token = " + g_anti_steal_token + "\n g_secret_key = " + g_secret_key + "\n g_tracker_http_port = " + g_tracker_http_port + + "\n g_multi_storage_ips = " + g_multi_storage_ips + + "\n g_connect_first_by = " + (g_connect_first_by == CONNECT_FIRST_BY_TRACKER ? "tracker" : "last-connected") + "\n g_connection_pool_enabled = " + g_connection_pool_enabled + "\n g_connection_pool_max_count_per_entry = " + g_connection_pool_max_count_per_entry + "\n g_connection_pool_max_idle_time(ms) = " + g_connection_pool_max_idle_time diff --git a/src/main/java/org/csource/fastdfs/ProtoCommon.java b/src/main/java/org/csource/fastdfs/ProtoCommon.java index a8a4de7..b13be40 100644 --- a/src/main/java/org/csource/fastdfs/ProtoCommon.java +++ b/src/main/java/org/csource/fastdfs/ProtoCommon.java @@ -26,6 +26,7 @@ import java.util.Arrays; */ public class ProtoCommon { public static final byte FDFS_PROTO_CMD_QUIT = 82; + public static final byte TRACKER_PROTO_CMD_FETCH_STORAGE_IDS = 69; public static final byte TRACKER_PROTO_CMD_SERVER_LIST_GROUP = 91; public static final byte TRACKER_PROTO_CMD_SERVER_LIST_STORAGE = 92; public static final byte TRACKER_PROTO_CMD_SERVER_DELETE_STORAGE = 93; @@ -310,6 +311,23 @@ public class ProtoCommon { return headerInfo.errno == 0 ? true : false; } + /** + * int convert to buff (big-endian) + * + * @param n int number + * @return 4 bytes buff + */ + public static byte[] int2buff(int n) { + byte[] bs; + + bs = new byte[4]; + bs[0] = (byte) ((n >> 24) & 0xFF); + bs[1] = (byte) ((n >> 16) & 0xFF); + bs[2] = (byte) ((n >> 8) & 0xFF); + bs[3] = (byte) (n & 0xFF); + return bs; + } + /** * long convert to buff (big-endian) * @@ -317,19 +335,18 @@ public class ProtoCommon { * @return 8 bytes buff */ public static byte[] long2buff(long n) { - byte[] bs; + byte[] bs; - bs = new byte[8]; - bs[0] = (byte) ((n >> 56) & 0xFF); - bs[1] = (byte) ((n >> 48) & 0xFF); - bs[2] = (byte) ((n >> 40) & 0xFF); - bs[3] = (byte) ((n >> 32) & 0xFF); - bs[4] = (byte) ((n >> 24) & 0xFF); - bs[5] = (byte) ((n >> 16) & 0xFF); - bs[6] = (byte) ((n >> 8) & 0xFF); - bs[7] = (byte) (n & 0xFF); - - return bs; + bs = new byte[8]; + bs[0] = (byte) ((n >> 56) & 0xFF); + bs[1] = (byte) ((n >> 48) & 0xFF); + bs[2] = (byte) ((n >> 40) & 0xFF); + bs[3] = (byte) ((n >> 32) & 0xFF); + bs[4] = (byte) ((n >> 24) & 0xFF); + bs[5] = (byte) ((n >> 16) & 0xFF); + bs[6] = (byte) ((n >> 8) & 0xFF); + bs[7] = (byte) (n & 0xFF); + return bs; } /** diff --git a/src/main/java/org/csource/fastdfs/StorageAddressMap.java b/src/main/java/org/csource/fastdfs/StorageAddressMap.java new file mode 100644 index 0000000..8c53370 --- /dev/null +++ b/src/main/java/org/csource/fastdfs/StorageAddressMap.java @@ -0,0 +1,62 @@ +/** + * Copyright (C) 2023 Happy Fish / YuQing + *
+ * FastDFS Java Client may be copied only under the terms of the GNU Lesser
+ * General Public License (LGPL).
+ * Please visit the FastDFS Home Page https://github.com/happyfish100/fastdfs for more detail.
+ */
+
+package org.csource.fastdfs;
+
+import java.util.HashMap;
+import java.net.InetSocketAddress;
+
+/**
+ * Storage Server Address Map
+ *
+ * @author Happy Fish / YuQing
+ * @version Version 1.31
+ */
+public class StorageAddressMap {
+
+ protected boolean without_port;
+ protected HashMap