最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Redis Lettuce連接redis集群實現(xiàn)過程詳細講解

 更新時間:2023年01月17日 16:26:30   作者:FlyLikeButterfly  
這篇文章主要介紹了Redis Lettuce連接redis集群實現(xiàn)過程,文中通過示例代碼介紹的非常詳細,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)吧

前言

Lettuce連接redis集群使用的都是集群專用類,像RedisClusterClient、StatefulRedisClusterConnection、RedisAdvancedClusterCommands、StatefulRedisClusterPubSubConnection等等;

Lettuce對rediscluster的支持:

  • 支持所有Cluster命令;
  • 基于鍵哈希槽的路由節(jié)點;
  • 對集群命令高級抽象;
  • 在多個集群節(jié)點上執(zhí)行命令;
  • 處理MOVED和ASK重定向;
  • 通過槽位和ip端口直接連接集群節(jié)點;
  • SSL和身份驗證;
  • 定期和自適應(yīng)集群拓撲更新;
  • 發(fā)布訂閱;

啟動時只需至少一個可以連接的集群節(jié)點就可以,能夠自動拓撲出集群全部節(jié)點;也可以使用ReadFrom設(shè)置讀取數(shù)據(jù)來源,跟主從模式一樣;

雖然redis本身的多鍵命令要求key必須都在同一個槽位,但Lettuce對一部分命令多了優(yōu)化,可以對多鍵命令進行跨槽位執(zhí)行,通過將對不同槽位鍵的操作命令分解為多條命令,單個命令以fork/join方式并發(fā)運行,最后將結(jié)果合并返回;

可以跨槽位的命令有

  • DEL:刪除鍵,返回刪除數(shù)量;
  • EXISTS:統(tǒng)計跨槽位的存在的鍵的數(shù)量;
  • MGET:獲取所有給定鍵的值,順序按照鍵的順序返回;
  • MSET:批量保存鍵值對,總是返回OK;
  • TOUCH:改變給定鍵的最后訪問時間,返回改變的鍵的數(shù)量;
  • UNLINK:刪除鍵并在另一個不同的線程中回收內(nèi)存,返回刪除數(shù)量;

提供跨槽位命令的api:RedisAdvancedClusterCommands、RedisAdvancedClusterAsyncCommands、RedisAdvancedClusterReactiveCommands;

可以在多個集群節(jié)點上執(zhí)行的命令有

  • CLIENT SETNAME:在所有已知的集群節(jié)點上設(shè)置客戶端的名稱,總是返回OK;
  • KEYS:返回所有master上存儲的key;
  • DBSIZE:返回所有master上存儲的key的數(shù)量;
  • FLUSHALL:清空master上的所有數(shù)據(jù),總是返回OK;
  • FLUSHDB:清空master上的所有數(shù)據(jù),總是返回OK;
  • RANDOMKEY:從隨機master上返回隨機的key;
  • SCAN:根據(jù)ReadFrom設(shè)置掃描整個集群的鍵空間;
  • SCRIPT FLUSH:從所有的集群節(jié)點腳本緩存中刪除所有腳本;
  • SCRIPT LOAD:在所有的集群節(jié)點上加載lua腳本;
  • SCRIPT KILL:在所有集群節(jié)點上殺死腳本;(即使腳本沒有運行調(diào)用也不會失?。?/li>
  • SHUTDOWN:將數(shù)據(jù)集同步保存到磁盤,然后關(guān)閉集群所有節(jié)點;

關(guān)于發(fā)布訂閱

普通用戶空間的發(fā)布訂閱,redis集群會發(fā)送到每個節(jié)點,發(fā)布者和訂閱者不需要在同一個節(jié)點,普通訂閱發(fā)布消息可以在集群拓撲改變時重新連接。對于鍵空間事件,只會發(fā)到自己的節(jié)點,不會擴散到其他節(jié)點,要訂閱鍵空間事件可以去適當(dāng)?shù)亩鄠€節(jié)點上訂閱,或者使用RedisClusterClient消息傳播和NodeSelectionAPI獲得一個托管連接集合;

注意:由于主從同步,鍵會被復(fù)制到多個從節(jié)點上,特別是鍵過期事件,會在主從節(jié)點上都產(chǎn)生過期事件,如果訂閱從節(jié)點,可能會收到多條相同的過期事件;訂閱是通過NodeSelectionAPI或者單個節(jié)點調(diào)用subscribe(…)發(fā)出的,訂閱對于新增的節(jié)點無效;

測試Demo:(redis版本7.0.2,Lettuce版本6.1.8)

集群節(jié)點:虛擬機 192.168.1.31,端口 9001-9006,集群節(jié)點已設(shè)置notify-keyspace-events AK;

package testlettuce;
import java.time.Duration;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import io.lettuce.core.ClientOptions.DisconnectedBehavior;
import io.lettuce.core.KeyScanCursor;
import io.lettuce.core.KeyValue;
import io.lettuce.core.ReadFrom;
import io.lettuce.core.RedisURI;
import io.lettuce.core.ScanCursor;
import io.lettuce.core.SocketOptions;
import io.lettuce.core.SslOptions;
import io.lettuce.core.TimeoutOptions;
import io.lettuce.core.cluster.ClusterClientOptions;
import io.lettuce.core.cluster.ClusterTopologyRefreshOptions;
import io.lettuce.core.cluster.RedisClusterClient;
import io.lettuce.core.cluster.api.StatefulRedisClusterConnection;
import io.lettuce.core.cluster.api.sync.Executions;
import io.lettuce.core.cluster.api.sync.NodeSelection;
import io.lettuce.core.cluster.api.sync.NodeSelectionCommands;
import io.lettuce.core.cluster.api.sync.RedisAdvancedClusterCommands;
import io.lettuce.core.cluster.pubsub.StatefulRedisClusterPubSubConnection;
import io.lettuce.core.cluster.pubsub.api.async.NodeSelectionPubSubAsyncCommands;
import io.lettuce.core.cluster.pubsub.api.async.PubSubAsyncNodeSelection;
import io.lettuce.core.cluster.pubsub.api.reactive.RedisClusterPubSubReactiveCommands;
import io.lettuce.core.protocol.DecodeBufferPolicies;
import io.lettuce.core.protocol.ProtocolVersion;
import io.lettuce.core.pubsub.RedisPubSubListener;
import io.lettuce.core.pubsub.api.async.RedisPubSubAsyncCommands;
public class TestLettuceCluster {
	/**
	 * @param args
	 */
	public static void main(String[] args) {
		List<RedisURI> nodeList = new ArrayList<>();
		nodeList.add(RedisURI.builder().withHost("192.168.1.31").withPort(9001).withAuthentication("default", "123456").build());
		nodeList.add(RedisURI.builder().withHost("192.168.1.31").withPort(9002).withAuthentication("default", "123456").build());
		nodeList.add(RedisURI.builder().withHost("192.168.1.31").withPort(9003).withAuthentication("default", "123456").build());
		nodeList.add(RedisURI.builder().withHost("192.168.1.31").withPort(9004).withAuthentication("default", "123456").build());
		nodeList.add(RedisURI.builder().withHost("192.168.1.31").withPort(9005).withAuthentication("default", "123456").build());
		nodeList.add(RedisURI.builder().withHost("192.168.1.31").withPort(9006).withAuthentication("default", "123456").build());
		RedisClusterClient clusterClient = RedisClusterClient.create(nodeList);
		ClusterTopologyRefreshOptions clusterTopologyRefreshOptions = ClusterTopologyRefreshOptions.builder()
	            .adaptiveRefreshTriggersTimeout(Duration.ofSeconds(5L))//設(shè)置自適應(yīng)拓撲刷新超時,每次超時刷新一次,默認30s;
	            .closeStaleConnections(false)//刷新拓撲時是否關(guān)閉失效連接,默認true,isPeriodicRefreshEnabled()為true時生效;
	            .dynamicRefreshSources(true)//從拓撲中發(fā)現(xiàn)新節(jié)點,并將新節(jié)點也作為拓撲的源節(jié)點,動態(tài)刷新可以發(fā)現(xiàn)全部節(jié)點并計算每個客戶端的數(shù)量,設(shè)置false則只有初始節(jié)點為源和計算客戶端數(shù)量;
	            .enableAllAdaptiveRefreshTriggers()//啟用全部觸發(fā)器自適應(yīng)刷新拓撲,默認關(guān)閉;
	            .enablePeriodicRefresh(Duration.ofSeconds(5L))//開啟定時拓撲刷新并設(shè)置周期;
	            .refreshTriggersReconnectAttempts(3)//長連接重新連接嘗試n次才拓撲刷新
	            .build();
		ClusterClientOptions clusterClientOptions = ClusterClientOptions.builder()
				.autoReconnect(true)//在連接丟失時開啟或關(guān)閉自動重連,默認true;
				.cancelCommandsOnReconnectFailure(true)//允許在重連失敗取消排隊命令,默認false;
				.decodeBufferPolicy(DecodeBufferPolicies.always())//設(shè)置丟棄解碼緩沖區(qū)的策略,以回收內(nèi)存;always:解碼后丟棄,最大內(nèi)存效率;alwaysSome:解碼后丟棄一部分;ratio(n)基于比率丟棄,n/(1+n),通常用1-10對應(yīng)50%-90%;
				.disconnectedBehavior(DisconnectedBehavior.DEFAULT)//設(shè)置連接斷開時命令的調(diào)用行為,默認啟用重連;DEFAULT:啟用時重連中接收命令,禁用時重連中拒絕命令;ACCEPT_COMMANDS:重連中接收命令;REJECT_COMMANDS:重連中拒絕命令;
//				.maxRedirects(5)//當(dāng)鍵從一個節(jié)點遷移到另一個節(jié)點,集群重定向次數(shù),默認5;
//				.nodeFilter(nodeFilter)//設(shè)置節(jié)點過濾器
//				.pingBeforeActivateConnection(true)//激活連接前設(shè)置PING,默認true;
//				.protocolVersion(ProtocolVersion.RESP3)//設(shè)置協(xié)議版本,默認RESP3;
//				.publishOnScheduler(false)//使用專用的調(diào)度器發(fā)出響應(yīng)信號,默認false,啟用時數(shù)據(jù)信號將使用服務(wù)的多線程發(fā)出;
//				.requestQueueSize(requestQueueSize)//設(shè)置每個連接請求隊列大?。?
//				.scriptCharset(scriptCharset)//設(shè)置Lua腳本編碼為byte[]的字符集,默認StandardCharsets.UTF_8;
//				.socketOptions(SocketOptions.builder().connectTimeout(Duration.ofSeconds(10)).keepAlive(true).tcpNoDelay(true).build())//設(shè)置低級套接字的屬性
//				.sslOptions(SslOptions.builder().build())//設(shè)置ssl屬性
//				.suspendReconnectOnProtocolFailure(false)//當(dāng)重新連接遇到協(xié)議失敗時暫停重新連接(SSL驗證,連接失敗前PING),默認值為false;
//				.timeoutOptions(TimeoutOptions.enabled(Duration.ofSeconds(10)))//設(shè)置超時來取消和終止命令;
				.topologyRefreshOptions(clusterTopologyRefreshOptions)//設(shè)置拓撲更新設(shè)置
				.validateClusterNodeMembership(true)//在允許連接到集群節(jié)點之前,驗證集群節(jié)點成員關(guān)系,默認值為true;
				.build();
		clusterClient.setDefaultTimeout(Duration.ofSeconds(5L));
		clusterClient.setOptions(clusterClientOptions);
		StatefulRedisClusterConnection<String, String> clusterConn = clusterClient.connect();
		clusterConn.setReadFrom(ReadFrom.ANY);//設(shè)置從哪些節(jié)點讀取數(shù)據(jù);
		RedisAdvancedClusterCommands<String, String> clusterCmd = clusterConn.sync();
		clusterCmd.set("a", "A");
		clusterCmd.set("b", "B");
		clusterCmd.set("c", "C");
		clusterCmd.set("d", "D"); 
		System.out.println("get a=" + clusterCmd.get("a"));
		System.out.println("get b=" + clusterCmd.get("b"));
		System.out.println("get c=" + clusterCmd.get("c"));
		System.out.println("get d=" + clusterCmd.get("d"));
		//跨槽位命令
		Map<String, String> kvmap = new HashMap<>();
		kvmap.put("a", "AA");
		kvmap.put("b", "BB");
		kvmap.put("c", "CC");
		kvmap.put("d", "DD");
		clusterCmd.mset(kvmap);//Lettuce做了優(yōu)化,支持一些命令的跨槽位命令;
		System.out.println("Lettuce mget:" + clusterCmd.mget("a", "b", "c", "d"));
		//選定部分節(jié)點操作
		NodeSelection<String, String> replicas = clusterCmd.replicas();
		NodeSelectionCommands<String, String> replicaseCmd = replicas.commands();
		Executions<KeyScanCursor<String>> executions = replicaseCmd.scan(ScanCursor.INITIAL);
		executions.forEach(s -> {System.out.println(s.getKeys());});
		//訂閱發(fā)布消息
		StatefulRedisClusterPubSubConnection<String, String> pubSubConn = clusterClient.connectPubSub();
		pubSubConn.addListener(new RedisPubSubListener<String, String>() {
			@Override
			public void message(String channel, String message) {
				System.out.println("[message]ch:" + channel + ",msg:" + message);
			}
			@Override
			public void message(String pattern, String channel, String message) {
			}
			@Override
			public void subscribed(String channel, long count) {
				System.out.println("[subscribed]ch:" + channel);
			}
			@Override
			public void psubscribed(String pattern, long count) {
			}
			@Override
			public void unsubscribed(String channel, long count) {
			}
			@Override
			public void punsubscribed(String pattern, long count) {
			}
		});
		pubSubConn.sync().subscribe("TEST_Ch");//(回調(diào)內(nèi)部使用阻塞調(diào)用或者lettuce同步api調(diào)用,需使用異步訂閱)
		clusterCmd.publish("TEST_Ch", "MSGMSGMSG");
		//響應(yīng)式訂閱,可以監(jiān)聽ChannelMessage和PatternMessage,使用鏈式過濾處理計算等操作
		RedisClusterPubSubReactiveCommands<String, String> pubsubReactive = pubSubConn.reactive();
		pubsubReactive.subscribe("TEST_Ch2").subscribe();
		pubsubReactive.observeChannels()
			.filter(chmsg -> {return chmsg.getMessage().contains("tom");})
			.doOnNext(chmsg -> {System.out.println("<tom>" + chmsg.getChannel() + ">>" + chmsg.getMessage());})
			.subscribe();
		clusterCmd.publish("TEST_Ch2", "send to jerry");
		clusterCmd.publish("TEST_Ch", "tom MSG");
		clusterCmd.publish("TEST_Ch2", "this is tom");
		//keySpaceEvent事件
		StatefulRedisClusterPubSubConnection<String, String> clusterPubSubConn = clusterClient.connectPubSub();
		clusterPubSubConn.setNodeMessagePropagation(true);//啟用禁用節(jié)點消息傳播到該listener,例如只能在本節(jié)點通知的鍵事件通知;
		RedisPubSubListener<String, String> listener  = new RedisPubSubListener<String, String>() {
			@Override
			public void unsubscribed(String channel, long count) {
				System.out.println("unsubscribed_ch:" + channel);
			}
			@Override
			public void subscribed(String channel, long count) {
				System.out.println("subscribed_ch:" + channel);
			}
			@Override
			public void punsubscribed(String pattern, long count) {
				System.out.println("punsubscribed_pattern:" + pattern);
			}
			@Override
			public void psubscribed(String pattern, long count) {
				System.out.println("psubscribed_pattern:" + pattern);
			}
			@Override
			public void message(String pattern, String channel, String message) {
				System.out.println("message_pattern:" + pattern + " ch:" + channel + " msg:" + message);
			}
			@Override
			public void message(String channel, String message) {
				System.out.println("message_ch:" + channel + " msg:" + message);
			}
		};
		clusterPubSubConn.addListener(listener);
		PubSubAsyncNodeSelection<String, String> allPubSubAsyncNodeSelection = clusterPubSubConn.async().all();
		NodeSelectionPubSubAsyncCommands<String, String> pubsubAsyncCmd = allPubSubAsyncNodeSelection.commands();
		clusterCmd.setex("a", 1, "A");
		pubsubAsyncCmd.psubscribe("__keyspace@0__:*");
		try {
			Thread.sleep(3000);
		} catch (InterruptedException e) {
			e.printStackTrace();
		}
		System.out.println("end");
	}
}

運行結(jié)果:

另外,還有一個cluster專用的Listener:RedisClusterPubSubListener,可以從listener里獲得發(fā)布消息的節(jié)點信息:

RedisClusterPubSubListener<String, String> clusterListener = new RedisClusterPubSubListener<String, String>() {
			@Override
			public void message(RedisClusterNode node, String channel, String message) {
			}
			@Override
			public void message(RedisClusterNode node, String pattern, String channel, String message) {
			}
			@Override
			public void subscribed(RedisClusterNode node, String channel, long count) {
			}
			@Override
			public void psubscribed(RedisClusterNode node, String pattern, long count) {
			}
			@Override
			public void unsubscribed(RedisClusterNode node, String channel, long count) {
			}
			@Override
			public void punsubscribed(RedisClusterNode node, String pattern, long count) {
			}
		};

到此這篇關(guān)于Redis Lettuce連接redis集群實現(xiàn)過程詳細講解的文章就介紹到這了,更多相關(guān)Redis Lettuce連接redis集群內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • 5種Java中數(shù)組的拷貝方法總結(jié)分享

    5種Java中數(shù)組的拷貝方法總結(jié)分享

    這篇文章主要介紹了5種Java中數(shù)組的拷貝方法總結(jié)分享,文章圍繞主題展開詳細的內(nèi)容介紹,具有一定的參考價值,需要的朋友可以參考一下
    2022-07-07
  • SpringBoot實現(xiàn)其他普通類調(diào)用Spring管理的Service,dao等bean

    SpringBoot實現(xiàn)其他普通類調(diào)用Spring管理的Service,dao等bean

    這篇文章主要介紹了SpringBoot實現(xiàn)其他普通類調(diào)用Spring管理的Service,dao等bean,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-11-11
  • Java中l(wèi)ist集合為空或為null的區(qū)別說明

    Java中l(wèi)ist集合為空或為null的區(qū)別說明

    這篇文章主要介紹了Java中l(wèi)ist集合為空或為null的區(qū)別說明,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-11-11
  • Java?SSM實現(xiàn)前后端協(xié)議聯(lián)調(diào)詳解下篇

    Java?SSM實現(xiàn)前后端協(xié)議聯(lián)調(diào)詳解下篇

    首先我們已經(jīng)知道,在現(xiàn)在流行的“前后端完全分離”架構(gòu)中,前后端聯(lián)調(diào)是一個不可能避免的問題,這篇文章主要介紹了Java?SSM實現(xiàn)前后端協(xié)議聯(lián)調(diào)過程
    2022-08-08
  • 關(guān)于Java?中?Future?的?get?方法超時問題

    關(guān)于Java?中?Future?的?get?方法超時問題

    這篇文章主要介紹了Java?中?Future?的?get?方法超時,最常見的理解就是,“超時以后,當(dāng)前線程繼續(xù)執(zhí)行,線程池里的對應(yīng)線程中斷”,真的是這樣嗎?本文給大家詳細介紹,需要的朋友參考下吧
    2022-06-06
  • 淺談Java當(dāng)作數(shù)組的幾個應(yīng)用場景

    淺談Java當(dāng)作數(shù)組的幾個應(yīng)用場景

    數(shù)組可以存放多個同一類型的數(shù)據(jù),可以存儲基本數(shù)據(jù)類型,引用數(shù)據(jù)類型(對象),下面這篇文章主要給大家介紹了關(guān)于Java當(dāng)作數(shù)組的幾個應(yīng)用場景,需要的朋友可以參考下
    2022-11-11
  • Spring使用Configuration注解管理bean的方式詳解

    Spring使用Configuration注解管理bean的方式詳解

    在Spring的世界里,Configuration注解就像是一位細心的園丁,它的主要職責(zé)是在這個繁花似錦的園子里,幫助我們聲明和管理各種各樣的bean,本文給大家介紹了在Spring中如何優(yōu)雅地管理你的bean,需要的朋友可以參考下
    2024-05-05
  • 最全LocalDateTime、LocalDate、Date、String相互轉(zhuǎn)化的方法

    最全LocalDateTime、LocalDate、Date、String相互轉(zhuǎn)化的方法

    大家在開發(fā)過程中必不可少的和日期打交道,對接別的系統(tǒng)時,時間日期格式不一致,每次都要轉(zhuǎn)化,本文為大家準備了最全的LocalDateTime、LocalDate、Date、String相互轉(zhuǎn)化方法,需要的可以參考一下
    2023-06-06
  • Spring@Value使用獲取配置信息為null的操作

    Spring@Value使用獲取配置信息為null的操作

    這篇文章主要介紹了Spring@Value使用獲取配置信息為null的操作,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-07-07
  • Java面試之如何實現(xiàn)10億數(shù)據(jù)判重

    Java面試之如何實現(xiàn)10億數(shù)據(jù)判重

    當(dāng)數(shù)據(jù)量比較大時,使用常規(guī)的方式來判重就不行了,所以這篇文章小編主要來和大家介紹一下Java實現(xiàn)10億數(shù)據(jù)判重的相關(guān)方法,希望對大家有所幫助
    2024-02-02

最新評論

开远市| 漳州市| 时尚| 兴文县| 安溪县| 临颍县| 肥东县| 通山县| 通江县| 长顺县| 弥渡县| 东山县| 建水县| 普陀区| 南和县| 拜城县| 景宁| 水富县| 绵阳市| 无为县| 大同市| 定日县| 定西市| 莱芜市| 固原市| 绩溪县| 城口县| 濉溪县| 同江市| 长葛市| 青海省| 江源县| 焦作市| 兴义市| 永昌县| 喜德县| 聂拉木县| 新昌县| 衡山县| 梧州市| 鹿邑县|