ZhiYeJianKang_PeiXun/cyqdata-master/DistributedCache/CacheImplement/Redis/RedisClient.cs

727 lines
26 KiB
C#
Raw Permalink Normal View History

2025-02-20 15:41:53 +08:00
using System;
using System.Collections.Generic;
using System.Collections.Specialized;
using System.Globalization;
using System.Text;
using CYQ.Data.Tool;
namespace CYQ.Data.Cache
{
/// <summary>
/// Redis client main class.
/// Use the static methods Setup and GetInstance to setup and get an instance of the client for use.
/// </summary>
internal class RedisClient : ClientBase
{
#region Static fields and methods.
public static RedisClient Create(string configValue)
{
return new RedisClient(configValue);
}
private RedisClient(string configValue)
{
hostServer = new HostServer(CacheType.Redis, configValue);
hostServer.OnAuthEvent += new HostServer.AuthDelegate(hostServer_OnAuthEvent);
}
bool hostServer_OnAuthEvent(MSocket socket)
{
if (!Auth(socket.HostNode.Password, socket))
{
string err = "Auth password fail : " + socket.HostNode.Password;
socket.HostNode.Error = err;
Error.Throw(err);
}
return true;
}
#endregion
#region Add<EFBFBD><EFBFBD>SetNX
public bool Add(string key, object value, int seconds) { return Add("setnx", key, true, value, hash(key), seconds); }
private bool Add(string command, string key, bool keyIsChecked, object value, uint hash, int expirySeconds)
{
if (!keyIsChecked)
{
checkKey(key);
}
string result = hostServer.Execute<string>(hash, "", delegate (MSocket socket, out bool isNoResponse)
{
SerializedType type;
byte[] bytes;
byte[] typeBit = new byte[1];
bytes = Serializer.Serialize(value, out type, compressionThreshold);
typeBit[0] = (byte)type;
// CheckDB(socket, hash);
int db = GetDBIndex(socket, hash);
// Console.WriteLine("Set :" + key + ":" + hash + " db." + db);
int skipCmd = 0;
using (RedisCommand cmd = new RedisCommand(socket))
{
if (db > -1)
{
cmd.Reset(2, "Select");
cmd.AddKey(db.ToString());
skipCmd++;
}
cmd.Reset(3, command);
cmd.AddKey(key);
cmd.AddValue(typeBit, bytes);
cmd.Reset(2, "ttl");
cmd.AddKey(key);//<2F><><EFBFBD><EFBFBD>ʧЧʱ<D0A7><EFBFBD>Ƿ񷵻<C7B7>-1<><31><EFBFBD><EFBFBD><EFBFBD><EFBFBD>δ<EFBFBD><CEB4><EFBFBD>ù<EFBFBD><C3B9><EFBFBD>ʱ<EFBFBD><EFBFBD><E4A3AC><EFBFBD><EFBFBD>setNx<4E><78>expireԭ<65><D4AD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD>
cmd.Send();
socket.SkipToEndOfLine(skipCmd);
result = socket.ReadResponse();
string ttl = socket.ReadResponse();
if (result == ":1" || ttl == ":-1")
{
if (expirySeconds > 0)
{
cmd.Reset(3, "EXPIRE");
cmd.AddKey(key);
cmd.AddKey(expirySeconds.ToString());
cmd.Send();
//result = socket.ReadResponse();
//if (result != ":1")
//{
// cmd.Reset(2, "DEL");
// cmd.AddKey(key);
// cmd.Send();
socket.SkipToEndOfLine(1);//<2F><><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD>
//}
}
}
isNoResponse = string.IsNullOrEmpty(result);
return result;
}
});
return result.StartsWith("+OK") || result.StartsWith(":1");
}
#endregion
#region Set<EFBFBD><EFBFBD>Append
public bool Append(string key, object value, int seconds) { return Set("append", key, true, value, hash(key), seconds); }
public bool Set(string key, object value, int seconds) { return Set("set", key, true, value, hash(key), seconds); }
private bool Set(string command, string key, bool keyIsChecked, object value, uint hash, int expirySeconds)
{
if (!keyIsChecked)
{
checkKey(key);
}
string result = hostServer.Execute<string>(hash, "", delegate (MSocket socket, out bool isNoResponse)
{
SerializedType type;
byte[] bytes;
byte[] typeBit = new byte[1];
bytes = Serializer.Serialize(value, out type, compressionThreshold);
typeBit[0] = (byte)type;
// CheckDB(socket, hash);
int db = GetDBIndex(socket, hash);
// Console.WriteLine("Set :" + key + ":" + hash + " db." + db);
int skipCmd = 0;
using (RedisCommand cmd = new RedisCommand(socket))
{
if (db > -1)
{
cmd.Reset(2, "Select");
cmd.AddKey(db.ToString());
skipCmd++;
}
cmd.Reset(3, command);
cmd.AddKey(key);
cmd.AddValue(typeBit, bytes);
skipCmd++;
if (expirySeconds > 0)
{
cmd.Reset(3, "EXPIRE");
cmd.AddKey(key);
cmd.AddKey(expirySeconds.ToString());
skipCmd++;
}
}
socket.SkipToEndOfLine(skipCmd - 1);//ȡ<><C8A1><EFBFBD><EFBFBD>1<EFBFBD><31><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD>Ľ<EFBFBD><C4BD><EFBFBD>
result = socket.ReadResponse();
isNoResponse = string.IsNullOrEmpty(result);
return result;
});
bool ret = result == "+OK" || result == ":1";
//if (!ret)
//{
//}
return ret;
}
#endregion
#region Get
public object Get(string key) { return Get("get", key, true, hash(key)); }
private object Get(string command, string key, bool keyIsChecked, uint hash)
{
if (!keyIsChecked)
{
checkKey(key);
}
object value = hostServer.Execute<object>(hash, null, delegate (MSocket socket, out bool isNoResponse)
{
int db = GetDBIndex(socket, hash);
using (RedisCommand cmd = new RedisCommand(socket))
{
if (db > -1)
{
cmd.Reset(2, "Select");
cmd.AddKey(db.ToString());
}
cmd.Reset(2, command);
cmd.AddKey(key);
}
if (db > -1) { socket.SkipToEndOfLine(); }
string result = socket.ReadResponse();
isNoResponse = string.IsNullOrEmpty(result);
if (!string.IsNullOrEmpty(result) && result[0] == '$')
{
int len = 0;
if (int.TryParse(result.Substring(1), out len) && len > 0)
{
byte[] bytes = socket.ReadBytes(len);
if (bytes.Length > 0)
{
byte[] data = new byte[bytes.Length - 1];
SerializedType st = (SerializedType)bytes[0];
Array.Copy(bytes, 1, data, 0, data.Length);
bytes = null;
return Serializer.DeSerialize(data, st);
}
}
}
return null;
});
return value;
}
#endregion
#region Exists
public bool ContainsKey(string key) { return ContainsKey(key, true, hash(key)); }
private bool ContainsKey(string key, bool keyIsChecked, uint hash)
{
if (!keyIsChecked)
{
checkKey(key);
}
return hostServer.Execute<bool>(hash, false, delegate (MSocket socket, out bool isNoResponse)
{
int db = GetDBIndex(socket, hash);
//Console.WriteLine("ContainsKey :" + key + ":" + hash + " db." + db);
using (RedisCommand cmd = new RedisCommand(socket))
{
if (db > -1)
{
cmd.Reset(2, "Select");
cmd.AddKey(db.ToString());
}
cmd.Reset(2, "Exists");
cmd.AddKey(key);
}
if (db > -1) { socket.SkipToEndOfLine(); }
string result = socket.ReadResponse();
isNoResponse = string.IsNullOrEmpty(result);
return !result.StartsWith(":0") && !result.StartsWith("-");
});
}
#endregion
#region Select DB
internal int GetDBIndex(MSocket socket, uint hash)
{
if (AppConfig.Redis.UseDBCount > 1 || AppConfig.Redis.UseDBIndex > 0)
{
return AppConfig.Redis.UseDBIndex > 0 ? AppConfig.Redis.UseDBIndex : (int)(hash % AppConfig.Redis.UseDBCount);//Ĭ<>Ϸ<EFBFBD>ɢ<EFBFBD><C9A2>16<31><36>DB<44>С<EFBFBD>
}
return -1;
}
#endregion
#region Delete
public bool Delete(string key) { return Delete(key, true, hash(key), 0); }
private bool Delete(string key, bool keyIsChecked, uint hash, int time)
{
if (!keyIsChecked)
{
checkKey(key);
}
return hostServer.Execute<bool>(hash, false, delegate (MSocket socket, out bool isNoResponse)
{
int db = GetDBIndex(socket, hash);
using (RedisCommand cmd = new RedisCommand(socket))
{
if (db > -1)
{
cmd.Reset(2, "Select");
cmd.AddKey(db.ToString());
}
cmd.Reset(2, "DEL");
cmd.AddKey(key);
}
if (db > -1)
{
socket.SkipToEndOfLine();
}
string result = socket.ReadResponse();
isNoResponse = string.IsNullOrEmpty(result);
return result.StartsWith(":1");
});
}
#endregion
#region Auth
private bool Auth(string password, MSocket socket)
{
if (!string.IsNullOrEmpty(password))
{
using (RedisCommand cmd = new RedisCommand(socket, 2, "AUTH"))
{
cmd.AddKey(password);
}
string result = socket.ReadLine();
return result.StartsWith("+OK");
}
return true;
}
#endregion
#region Flush All
public bool FlushAll()
{
foreach (KeyValuePair<string, HostNode> item in hostServer.HostList)
{
HostNode pool = item.Value;
hostServer.Execute(pool, delegate (MSocket socket)
{
using (RedisCommand cmd = new RedisCommand(socket, 1, "flushall"))
{
cmd.Send();
socket.SkipToEndOfLine();
}
});
}
return true;
}
#endregion
#region Stats
/// <summary>
/// This method corresponds to the "stats" command in the memcached protocol.
/// It will send the stats command to all servers, and it will return a Dictionary for each server
/// containing the results of the command.
/// </summary>
public Dictionary<string, Dictionary<string, string>> Stats()
{
Dictionary<string, Dictionary<string, string>> results = new Dictionary<string, Dictionary<string, string>>();
foreach (KeyValuePair<string, HostNode> item in hostServer.HostList)
{
results.Add(item.Key, Stats(item.Value));
}
return results;
}
private Dictionary<string, string> Stats(HostNode pool)
{
Dictionary<string, string> dic = new Dictionary<string, string>();
hostServer.Execute(pool, delegate (MSocket socket)
{
using (RedisCommand cmd = new RedisCommand(socket, 1, "info"))
{
}
string result = socket.ReadResponse();
if (!string.IsNullOrEmpty(result) && (result[0] == '$' || result == "+OK"))
{
string line = null;
bool isEnd = false;
while (true)
{
try
{
line = socket.ReadLine();
}
catch (Exception err)
{
break;
}
if (line == null)
{
if (isEnd)
{
break;
}
continue;
}
else if (line == "# Keyspace")
{
isEnd = true;
}
string[] s = line.Split(':');
if (s.Length > 1)
{
dic.Add(s[0], s[1]);
}
else
{
dic.Add(line, "- - -");
}
}
}
});
return dic;
}
#endregion
#region Exe All
public int SetAll(string key, object value, int seconds) { return SetAll("set", key, true, value, hash(key), seconds); }
private int SetAll(string command, string key, bool keyIsChecked, object value, uint hash, int expirySeconds)
{
if (!keyIsChecked)
{
checkKey(key);
}
UseSocket<bool> useSocket = delegate (MSocket socket, out bool isNoResponse)
{
SerializedType type;
byte[] bytes;
byte[] typeBit = new byte[1];
bytes = Serializer.Serialize(value, out type, compressionThreshold);
typeBit[0] = (byte)type;
// CheckDB(socket, hash);
int db = GetDBIndex(socket, hash);
// Console.WriteLine("Set :" + key + ":" + hash + " db." + db);
int skipCmd = 0;
using (RedisCommand cmd = new RedisCommand(socket))
{
if (db > -1)
{
cmd.Reset(2, "Select");
cmd.AddKey(db.ToString());
skipCmd++;
}
cmd.Reset(3, command);
cmd.AddKey(key);
cmd.AddValue(typeBit, bytes);
skipCmd++;
if (expirySeconds > 0)
{
cmd.Reset(3, "EXPIRE");
cmd.AddKey(key);
cmd.AddKey(expirySeconds.ToString());
skipCmd++;
}
}
socket.SkipToEndOfLine(skipCmd - 1);//ȡ<><C8A1><EFBFBD><EFBFBD>1<EFBFBD><31><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD>Ľ<EFBFBD><C4BD><EFBFBD>
string result = socket.ReadResponse();
isNoResponse = string.IsNullOrEmpty(result);
bool ret = result == "+OK" || result == ":1";
return ret;
};
return hostServer.ExecuteAll(useSocket, hash);
}
public int DeleteAll(string key)
{
return DeleteAll(key, true, hash(key), 0);
}
private int DeleteAll(string key, bool keyIsChecked, uint hash, int time)
{
if (!keyIsChecked)
{
checkKey(key);
}
UseSocket<bool> useSocket = delegate (MSocket socket, out bool isNoResponse)
{
int db = GetDBIndex(socket, hash);
using (RedisCommand cmd = new RedisCommand(socket))
{
if (db > -1)
{
cmd.Reset(2, "Select");
cmd.AddKey(db.ToString());
}
cmd.Reset(2, "DEL");
cmd.AddKey(key);
}
if (db > -1)
{
socket.SkipToEndOfLine();
}
string result = socket.ReadResponse();
isNoResponse = string.IsNullOrEmpty(result);
return result.StartsWith(":1");
};
return hostServer.ExecuteAll(useSocket, hash);
}
public int AddAll(string key, object value, int seconds) { return AddAll("setnx", key, true, value, hash(key), seconds); }
private int AddAll(string command, string key, bool keyIsChecked, object value, uint hash, int expirySeconds)
{
if (!keyIsChecked)
{
checkKey(key);
}
UseSocket<bool> useSocket = delegate (MSocket socket, out bool isNoResponse)
{
string result = string.Empty;
SerializedType type;
byte[] bytes;
byte[] typeBit = new byte[1];
bytes = Serializer.Serialize(value, out type, compressionThreshold);
typeBit[0] = (byte)type;
// CheckDB(socket, hash);
int db = GetDBIndex(socket, hash);
// Console.WriteLine("Set :" + key + ":" + hash + " db." + db);
int skipCmd = 0;
using (RedisCommand cmd = new RedisCommand(socket))
{
if (db > -1)
{
cmd.Reset(2, "Select");
cmd.AddKey(db.ToString());
skipCmd++;
}
cmd.Reset(3, command);
cmd.AddKey(key);
cmd.AddValue(typeBit, bytes);
cmd.Reset(2, "ttl");
cmd.AddKey(key);//<2F><><EFBFBD><EFBFBD>ʧЧʱ<D0A7><EFBFBD>Ƿ񷵻<C7B7>-1<><31><EFBFBD><EFBFBD><EFBFBD><EFBFBD>δ<EFBFBD><CEB4><EFBFBD>ù<EFBFBD><C3B9><EFBFBD>ʱ<EFBFBD><EFBFBD><E4A3AC><EFBFBD><EFBFBD>setNx<4E><78>expireԭ<65><D4AD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD>
cmd.Send();
socket.SkipToEndOfLine(skipCmd);
result = socket.ReadResponse();
string ttl = socket.ReadResponse();
if (result == ":1" || ttl == ":-1")
{
if (expirySeconds > 0)
{
cmd.Reset(3, "EXPIRE");
cmd.AddKey(key);
cmd.AddKey(expirySeconds.ToString());
cmd.Send();
socket.SkipToEndOfLine(1);//<2F><><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD>
}
}
isNoResponse = string.IsNullOrEmpty(result);
return result.StartsWith("+OK") || result.StartsWith(":1");
}
};
return hostServer.ExecuteAll(useSocket, hash);
//List<HostNode> okNodeList = new List<HostNode>();
//List<string> hosts = hostServer.HostList.GetKeys();
//int exeCount = 0;
//int okCount = 0;
//foreach (string host in hosts)
//{
// HostNode hostNode = hostServer.HostList[host];
// if (!hostNode.IsEndPointDead)
// {
// exeCount++;
// }
// bool isOK = hostServer.Execute<bool>(hostNode, hash, false, , false);
// if (isOK)
// {
// okCount++;
// okNodeList.Add(hostNode);
// }
//}
//bool retResult = false;
//if (exeCount < 3)
//{
// retResult = okCount > 0 && okCount == exeCount;//2<><32><EFBFBD>ڵ<EFBFBD><DAB5><EFBFBD><EFBFBD>£<EFBFBD>Ҫ<EFBFBD><D2AA>ȫ<EFBFBD><C8AB><EFBFBD>ɹ<EFBFBD><C9B9><EFBFBD>
//}
//else
//{
// retResult = okCount > exeCount / 2 + 1;//<2F><><EFBFBD><EFBFBD>1<EFBFBD><31><EFBFBD>ijɹ<C4B3><C9B9><EFBFBD>
//}
//if (!retResult && okCount > 0)
//{
// foreach (var node in okNodeList)
// {
// hostServer.Execute(node, delegate (MSocket socket)
// {
// int db = GetDBIndex(socket, hash);
// using (RedisCommand cmd = new RedisCommand(socket))
// {
// if (db > -1)
// {
// cmd.Reset(2, "Select");
// cmd.AddKey(db.ToString());
// }
// cmd.Reset(2, "DEL");
// cmd.AddKey(key);
// }
// if (db > -1)
// {
// socket.SkipToEndOfLine();
// }
// socket.SkipToEndOfLine();
// });
// }
//}
//return retResult;
}
public bool SetNXAll(string key, object value, int seconds) { return SetNXAll("setnx", key, true, value, hash(key), seconds); }
private bool SetNXAll(string command, string key, bool keyIsChecked, object value, uint hash, int expirySeconds)
{
if (!keyIsChecked)
{
checkKey(key);
}
UseSocket<bool> useSocket = delegate (MSocket socket, out bool isNoResponse)
{
string result = string.Empty;
SerializedType type;
byte[] bytes;
byte[] typeBit = new byte[1];
bytes = Serializer.Serialize(value, out type, compressionThreshold);
typeBit[0] = (byte)type;
// CheckDB(socket, hash);
int db = GetDBIndex(socket, hash);
// Console.WriteLine("Set :" + key + ":" + hash + " db." + db);
int skipCmd = 0;
using (RedisCommand cmd = new RedisCommand(socket))
{
if (db > -1)
{
cmd.Reset(2, "Select");
cmd.AddKey(db.ToString());
skipCmd++;
}
cmd.Reset(3, command);
cmd.AddKey(key);
cmd.AddValue(typeBit, bytes);
cmd.Reset(2, "ttl");
cmd.AddKey(key);//<2F><><EFBFBD><EFBFBD>ʧЧʱ<D0A7><EFBFBD>Ƿ񷵻<C7B7>-1<><31><EFBFBD><EFBFBD><EFBFBD><EFBFBD>δ<EFBFBD><CEB4><EFBFBD>ù<EFBFBD><C3B9><EFBFBD>ʱ<EFBFBD><EFBFBD><E4A3AC><EFBFBD><EFBFBD>setNx<4E><78>expireԭ<65><D4AD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD>
cmd.Send();
socket.SkipToEndOfLine(skipCmd);
result = socket.ReadResponse();
string ttl = socket.ReadResponse();
if (result == ":1" || ttl == ":-1")
{
if (expirySeconds > 0)
{
cmd.Reset(3, "EXPIRE");
cmd.AddKey(key);
cmd.AddKey(expirySeconds.ToString());
cmd.Send();
socket.SkipToEndOfLine(1);//<2F><><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD><EFBFBD>
}
}
isNoResponse = string.IsNullOrEmpty(result);
return result.StartsWith("+OK") || result.StartsWith(":1");
}
};
List<HostNode> okNodeList = new List<HostNode>();
List<string> hosts = hostServer.HostList.GetKeys();
int exeCount = 0;
int okCount = 0;
foreach (string host in hosts)
{
HostNode hostNode = hostServer.HostList[host];
if (!hostNode.IsEndPointDead)
{
exeCount++;
}
bool isOK = hostServer.Execute<bool>(hostNode, hash, false, useSocket, false);
if (isOK)
{
okCount++;
okNodeList.Add(hostNode);
}
}
bool retResult = false;
if (exeCount < 3)
{
retResult = okCount > 0 && okCount == exeCount;//2<><32><EFBFBD>ڵ<EFBFBD><DAB5><EFBFBD><EFBFBD>£<EFBFBD>Ҫ<EFBFBD><D2AA>ȫ<EFBFBD><C8AB><EFBFBD>ɹ<EFBFBD><C9B9><EFBFBD>
}
else
{
retResult = okCount > exeCount / 2 + 1;//<2F><><EFBFBD><EFBFBD>1<EFBFBD><31><EFBFBD>ijɹ<C4B3><C9B9><EFBFBD>
}
if (!retResult && okCount > 0)
{
foreach (var node in okNodeList)
{
hostServer.Execute(node, delegate (MSocket socket)
{
int db = GetDBIndex(socket, hash);
using (RedisCommand cmd = new RedisCommand(socket))
{
if (db > -1)
{
cmd.Reset(2, "Select");
cmd.AddKey(db.ToString());
}
cmd.Reset(2, "DEL");
cmd.AddKey(key);
}
if (db > -1)
{
socket.SkipToEndOfLine();
}
socket.SkipToEndOfLine();
});
}
}
return retResult;
}
#endregion
}
}