1from __future__ import annotations
2
3import asyncio
4import warnings
5from typing import (
6 TYPE_CHECKING,
7 Any,
8 AsyncIterator,
9 Awaitable,
10 Dict,
11 Iterable,
12 Iterator,
13 List,
14 Literal,
15 Mapping,
16 NoReturn,
17 Sequence,
18 overload,
19)
20
21from redis.crc import key_slot
22from redis.exceptions import RedisClusterException, RedisError
23from redis.typing import (
24 AnyKeyT,
25 AsyncClientProtocol,
26 ClusterCommandsProtocol,
27 ClusterLinksResponse,
28 ClusterNodeDetail,
29 ClusterShardsResponse,
30 EncodableT,
31 KeysT,
32 KeyT,
33 PatternT,
34 ResponseT,
35 StralgoResponse,
36 SyncClientProtocol,
37)
38from redis.utils import deprecated_function
39
40from .core import (
41 ACLCommands,
42 AsyncACLCommands,
43 AsyncDataAccessCommands,
44 AsyncFunctionCommands,
45 AsyncManagementCommands,
46 AsyncModuleCommands,
47 AsyncScriptCommands,
48 DataAccessCommands,
49 FunctionCommands,
50 HotkeysMetricsTypes,
51 ManagementCommands,
52 ModuleCommands,
53 PubSubCommands,
54 ScriptCommands,
55)
56from .helpers import list_or_args
57from .redismodules import AsyncRedisModuleCommands, RedisModuleCommands
58
59if TYPE_CHECKING:
60 from redis.asyncio.cluster import TargetNodesT
61 from redis.cluster import LoadBalancingStrategy
62
63# DEPRECATED - no longer consulted by the default metadata routing, and it will be removed in a future release.
64#
65# Replica safety is now decided from command metadata by the metadata resolver,
66# so this set no longer describes what commands are safe to execute on replicas.
67# It is kept as a public attribute only so an external caller reading it keeps working;
68# editing it changes nothing.
69#
70# To change replica safety, edit ``redis.commands.metadata._STATIC_COMMAND_METADATA`` or pass
71# a ``metadata_resolver`` to the client. Nothing here.
72READ_COMMANDS = frozenset(
73 [
74 # Bit Operations
75 "BITCOUNT",
76 "BITFIELD_RO",
77 "BITPOS",
78 # Scripting
79 "EVAL_RO",
80 "EVALSHA_RO",
81 "FCALL_RO",
82 # Key Operations
83 "DBSIZE",
84 "DIGEST",
85 "DUMP",
86 "EXISTS",
87 "EXPIRETIME",
88 "PEXPIRETIME",
89 "KEYS",
90 "SCAN",
91 "PTTL",
92 "RANDOMKEY",
93 "TTL",
94 "TYPE",
95 # String Operations
96 "GET",
97 "GETBIT",
98 "GETRANGE",
99 "MGET",
100 "STRLEN",
101 "LCS",
102 # Geo Operations
103 "GEODIST",
104 "GEOHASH",
105 "GEOPOS",
106 "GEOSEARCH",
107 # Hash Operations
108 "HEXISTS",
109 "HGET",
110 "HGETALL",
111 "HKEYS",
112 "HLEN",
113 "HMGET",
114 "HSTRLEN",
115 "HVALS",
116 "HRANDFIELD",
117 "HEXPIRETIME",
118 "HPEXPIRETIME",
119 "HTTL",
120 "HPTTL",
121 "HSCAN",
122 # List Operations
123 "LINDEX",
124 "LPOS",
125 "LLEN",
126 "LRANGE",
127 # Set Operations
128 "SCARD",
129 "SDIFF",
130 "SDIFFCARD",
131 "SINTER",
132 "SINTERCARD",
133 "SISMEMBER",
134 "SMISMEMBER",
135 "SMEMBERS",
136 "SRANDMEMBER",
137 "SUNION",
138 "SUNIONCARD",
139 "SSCAN",
140 # Sorted Set Operations
141 "ZCARD",
142 "ZCOUNT",
143 "ZDIFF",
144 "ZINTER",
145 "ZINTERCARD",
146 "ZLEXCOUNT",
147 "ZMSCORE",
148 "ZRANDMEMBER",
149 "ZRANGE",
150 "ZRANGEBYLEX",
151 "ZRANGEBYSCORE",
152 "ZRANK",
153 "ZREVRANGE",
154 "ZREVRANGEBYLEX",
155 "ZREVRANGEBYSCORE",
156 "ZREVRANK",
157 "ZSCAN",
158 "ZSCORE",
159 "ZUNION",
160 # Stream Operations
161 "XLEN",
162 "XPENDING",
163 "XRANGE",
164 "XREAD",
165 "XREVRANGE",
166 # JSON Module
167 "JSON.ARRINDEX",
168 "JSON.ARRLEN",
169 "JSON.GET",
170 "JSON.MGET",
171 "JSON.OBJKEYS",
172 "JSON.OBJLEN",
173 "JSON.RESP",
174 "JSON.STRLEN",
175 "JSON.TYPE",
176 # RediSearch Module
177 "FT.ALIASLIST",
178 "FT.EXPLAIN",
179 "FT.INFO",
180 "FT.PROFILE",
181 "FT.SEARCH",
182 ]
183)
184
185
186class ClusterMultiKeyCommands(ClusterCommandsProtocol):
187 """
188 A class containing commands that handle more than one key
189 """
190
191 # Read-routing configuration, which the cluster clients all set on the instance. The
192 # defaults are here for a host class that mixes these commands in without one - see
193 # ``_is_replica_safe`` below, which is the same fallback for the same audience - so
194 # that ``_execute_pipeline_by_slot`` reads a value rather than raising
195 # ``AttributeError``. Typed under TYPE_CHECKING because ``redis.cluster`` imports this
196 # module.
197 read_from_replicas: bool = False
198 load_balancing_strategy: LoadBalancingStrategy | None = None
199
200 def _partition_keys_by_slot(self, keys: Iterable[KeyT]) -> Dict[int, List[KeyT]]:
201 """Split keys into a dictionary that maps a slot to a list of keys."""
202
203 slots_to_keys = {}
204 for key in keys:
205 slot = key_slot(self.encoder.encode(key))
206 slots_to_keys.setdefault(slot, []).append(key)
207
208 return slots_to_keys
209
210 def _partition_pairs_by_slot(
211 self, mapping: Mapping[AnyKeyT, EncodableT]
212 ) -> Dict[int, List[EncodableT]]:
213 """Split pairs into a dictionary that maps a slot to a list of pairs."""
214
215 slots_to_pairs = {}
216 for pair in mapping.items():
217 slot = key_slot(self.encoder.encode(pair[0]))
218 slots_to_pairs.setdefault(slot, []).extend(pair)
219
220 return slots_to_pairs
221
222 def _is_replica_safe(self, command_name: str) -> bool:
223 warnings.warn(
224 "Using READ_COMMANDS for replica-safe checks is deprecated and will be removed in a future release. "
225 "Use a MetadataResolver instead.",
226 DeprecationWarning,
227 stacklevel=2,
228 )
229 return isinstance(command_name, str) and command_name.upper() in READ_COMMANDS
230
231 def _execute_pipeline_by_slot(
232 self, command: str, slots_to_args: Mapping[int, Iterable[EncodableT]]
233 ) -> List[Any]:
234 replica_safe = (
235 self.read_from_replicas or self.load_balancing_strategy is not None
236 ) and self._is_replica_safe(command)
237 pipe = self.pipeline()
238 [
239 pipe.execute_command(
240 command,
241 *slot_args,
242 target_nodes=[
243 self.nodes_manager.get_node_from_slot(
244 slot,
245 replica_safe,
246 self.load_balancing_strategy if replica_safe else None,
247 )
248 ],
249 )
250 for slot, slot_args in slots_to_args.items()
251 ]
252 return pipe.execute()
253
254 def _reorder_keys_by_command(
255 self,
256 keys: Iterable[KeyT],
257 slots_to_args: Mapping[int, Iterable[EncodableT]],
258 responses: Iterable[Any],
259 ) -> List[Any]:
260 results = {
261 k: v
262 for slot_values, response in zip(slots_to_args.values(), responses)
263 for k, v in zip(slot_values, response)
264 }
265 return [results[key] for key in keys]
266
267 def mget_nonatomic(self, keys: KeysT, *args: KeyT) -> List[Any | None]:
268 """
269 Splits the keys into different slots and then calls MGET
270 for the keys of every slot. This operation will not be atomic
271 if keys belong to more than one slot.
272
273 Returns a list of values ordered identically to ``keys``
274
275 For more information see https://redis.io/commands/mget
276 """
277
278 # Concatenate all keys into a list
279 keys = list_or_args(keys, args)
280
281 # Split keys into slots
282 slots_to_keys = self._partition_keys_by_slot(keys)
283
284 # Execute commands using a pipeline
285 res = self._execute_pipeline_by_slot("MGET", slots_to_keys)
286
287 # Reorder keys in the order the user provided & return
288 return self._reorder_keys_by_command(keys, slots_to_keys, res)
289
290 def mset_nonatomic(self, mapping: Mapping[AnyKeyT, EncodableT]) -> List[bool]:
291 """
292 Sets key/values based on a mapping. Mapping is a dictionary of
293 key/value pairs. Both keys and values should be strings or types that
294 can be cast to a string via str().
295
296 Splits the keys into different slots and then calls MSET
297 for the keys of every slot. This operation will not be atomic
298 if keys belong to more than one slot.
299
300 For more information see https://redis.io/commands/mset
301 """
302
303 # Partition the keys by slot
304 slots_to_pairs = self._partition_pairs_by_slot(mapping)
305
306 # Execute commands using a pipeline & return list of replies
307 return self._execute_pipeline_by_slot("MSET", slots_to_pairs)
308
309 def _split_command_across_slots(self, command: str, *keys: KeyT) -> int:
310 """
311 Runs the given command once for the keys
312 of each slot. Returns the sum of the return values.
313 """
314
315 # Partition the keys by slot
316 slots_to_keys = self._partition_keys_by_slot(keys)
317
318 # Sum up the reply from each command
319 return sum(self._execute_pipeline_by_slot(command, slots_to_keys))
320
321 @overload
322 def exists(self: SyncClientProtocol, *keys: KeyT) -> int: ...
323
324 @overload
325 def exists(self: AsyncClientProtocol, *keys: KeyT) -> Awaitable[int]: ...
326
327 def exists(self, *keys: KeyT) -> int | Awaitable[int]:
328 """
329 Returns the number of ``names`` that exist in the
330 whole cluster. The keys are first split up into slots
331 and then an EXISTS command is sent for every slot
332
333 For more information see https://redis.io/commands/exists
334 """
335 return self._split_command_across_slots("EXISTS", *keys)
336
337 @overload
338 def delete(self: SyncClientProtocol, *keys: KeyT) -> int: ...
339
340 @overload
341 def delete(self: AsyncClientProtocol, *keys: KeyT) -> Awaitable[int]: ...
342
343 def delete(self, *keys: KeyT) -> int | Awaitable[int]:
344 """
345 Deletes the given keys in the cluster.
346 The keys are first split up into slots
347 and then an DEL command is sent for every slot
348
349 Non-existent keys are ignored.
350 Returns the number of keys that were deleted.
351
352 For more information see https://redis.io/commands/del
353 """
354 return self._split_command_across_slots("DEL", *keys)
355
356 @overload
357 def touch(self: SyncClientProtocol, *keys: KeyT) -> int: ...
358
359 @overload
360 def touch(self: AsyncClientProtocol, *keys: KeyT) -> Awaitable[int]: ...
361
362 def touch(self, *keys: KeyT) -> int | Awaitable[int]:
363 """
364 Updates the last access time of given keys across the
365 cluster.
366
367 The keys are first split up into slots
368 and then an TOUCH command is sent for every slot
369
370 Non-existent keys are ignored.
371 Returns the number of keys that were touched.
372
373 For more information see https://redis.io/commands/touch
374 """
375 return self._split_command_across_slots("TOUCH", *keys)
376
377 @overload
378 def unlink(self: SyncClientProtocol, *keys: KeyT) -> int: ...
379
380 @overload
381 def unlink(self: AsyncClientProtocol, *keys: KeyT) -> Awaitable[int]: ...
382
383 def unlink(self, *keys: KeyT) -> int | Awaitable[int]:
384 """
385 Remove the specified keys in a different thread.
386
387 The keys are first split up into slots
388 and then an TOUCH command is sent for every slot
389
390 Non-existent keys are ignored.
391 Returns the number of keys that were unlinked.
392
393 For more information see https://redis.io/commands/unlink
394 """
395 return self._split_command_across_slots("UNLINK", *keys)
396
397
398class AsyncClusterMultiKeyCommands(ClusterMultiKeyCommands):
399 """
400 A class containing commands that handle more than one key
401 """
402
403 async def mget_nonatomic(self, keys: KeysT, *args: KeyT) -> List[Any | None]:
404 """
405 Splits the keys into different slots and then calls MGET
406 for the keys of every slot. This operation will not be atomic
407 if keys belong to more than one slot.
408
409 Returns a list of values ordered identically to ``keys``
410
411 For more information see https://redis.io/commands/mget
412 """
413
414 # Concatenate all keys into a list
415 keys = list_or_args(keys, args)
416
417 # Split keys into slots
418 slots_to_keys = self._partition_keys_by_slot(keys)
419
420 # Execute commands using a pipeline
421 res = await self._execute_pipeline_by_slot("MGET", slots_to_keys)
422
423 # Reorder keys in the order the user provided & return
424 return self._reorder_keys_by_command(keys, slots_to_keys, res)
425
426 async def mset_nonatomic(self, mapping: Mapping[AnyKeyT, EncodableT]) -> List[bool]:
427 """
428 Sets key/values based on a mapping. Mapping is a dictionary of
429 key/value pairs. Both keys and values should be strings or types that
430 can be cast to a string via str().
431
432 Splits the keys into different slots and then calls MSET
433 for the keys of every slot. This operation will not be atomic
434 if keys belong to more than one slot.
435
436 For more information see https://redis.io/commands/mset
437 """
438
439 # Partition the keys by slot
440 slots_to_pairs = self._partition_pairs_by_slot(mapping)
441
442 # Execute commands using a pipeline & return list of replies
443 return await self._execute_pipeline_by_slot("MSET", slots_to_pairs)
444
445 async def _split_command_across_slots(self, command: str, *keys: KeyT) -> int:
446 """
447 Runs the given command once for the keys
448 of each slot. Returns the sum of the return values.
449 """
450
451 # Partition the keys by slot
452 slots_to_keys = self._partition_keys_by_slot(keys)
453
454 # Sum up the reply from each command
455 return sum(await self._execute_pipeline_by_slot(command, slots_to_keys))
456
457 async def _is_replica_safe(self, command_name: str) -> bool:
458 warnings.warn(
459 "Using READ_COMMANDS for replica-safe checks is deprecated and will be removed in a future release. "
460 "Use an AsyncMetadataResolver instead.",
461 DeprecationWarning,
462 stacklevel=2,
463 )
464 return isinstance(command_name, str) and command_name.upper() in READ_COMMANDS
465
466 async def _execute_pipeline_by_slot(
467 self, command: str, slots_to_args: Mapping[int, Iterable[EncodableT]]
468 ) -> List[Any]:
469 if self._initialize:
470 await self.initialize()
471 replica_safe = (
472 self.read_from_replicas or self.load_balancing_strategy is not None
473 ) and await self._is_replica_safe(command)
474 pipe = self.pipeline()
475 [
476 pipe.execute_command(
477 command,
478 *slot_args,
479 target_nodes=[
480 self.nodes_manager.get_node_from_slot(
481 slot,
482 replica_safe,
483 self.load_balancing_strategy if replica_safe else None,
484 )
485 ],
486 )
487 for slot, slot_args in slots_to_args.items()
488 ]
489 return await pipe.execute()
490
491
492class ClusterManagementCommands(ManagementCommands):
493 """
494 A class for Redis Cluster management commands
495
496 The class inherits from Redis's core ManagementCommands class and do the
497 required adjustments to work with cluster mode
498 """
499
500 def slaveof(self, *args, **kwargs) -> NoReturn:
501 """
502 Make the server a replica of another instance, or promote it as master.
503
504 For more information see https://redis.io/commands/slaveof
505 """
506 raise RedisClusterException("SLAVEOF is not supported in cluster mode")
507
508 def replicaof(self, *args, **kwargs) -> NoReturn:
509 """
510 Make the server a replica of another instance, or promote it as master.
511
512 For more information see https://redis.io/commands/replicaof
513 """
514 raise RedisClusterException("REPLICAOF is not supported in cluster mode")
515
516 def swapdb(self, *args, **kwargs) -> NoReturn:
517 """
518 Swaps two Redis databases.
519
520 For more information see https://redis.io/commands/swapdb
521 """
522 raise RedisClusterException("SWAPDB is not supported in cluster mode")
523
524 @overload
525 def cluster_myid(
526 self: SyncClientProtocol, target_node: "TargetNodesT"
527 ) -> bytes | str: ...
528
529 @overload
530 def cluster_myid(
531 self: AsyncClientProtocol, target_node: "TargetNodesT"
532 ) -> Awaitable[bytes | str]: ...
533
534 def cluster_myid(self, target_node: "TargetNodesT") -> (bytes | str) | Awaitable[
535 bytes | str
536 ]:
537 """
538 Returns the node's id.
539
540 :target_node: 'ClusterNode'
541 The node to execute the command on
542
543 For more information check https://redis.io/commands/cluster-myid/
544 """
545 return self.execute_command("CLUSTER MYID", target_nodes=target_node)
546
547 @overload
548 def cluster_addslots(
549 self: SyncClientProtocol, target_node: "TargetNodesT", *slots: EncodableT
550 ) -> bool: ...
551
552 @overload
553 def cluster_addslots(
554 self: AsyncClientProtocol, target_node: "TargetNodesT", *slots: EncodableT
555 ) -> Awaitable[bool]: ...
556
557 def cluster_addslots(
558 self, target_node: "TargetNodesT", *slots: EncodableT
559 ) -> bool | Awaitable[bool]:
560 """
561 Assign new hash slots to receiving node. Sends to specified node.
562
563 :target_node: 'ClusterNode'
564 The node to execute the command on
565
566 For more information see https://redis.io/commands/cluster-addslots
567 """
568 return self.execute_command(
569 "CLUSTER ADDSLOTS", *slots, target_nodes=target_node
570 )
571
572 @overload
573 def cluster_addslotsrange(
574 self: SyncClientProtocol, target_node: "TargetNodesT", *slots: EncodableT
575 ) -> bool: ...
576
577 @overload
578 def cluster_addslotsrange(
579 self: AsyncClientProtocol, target_node: "TargetNodesT", *slots: EncodableT
580 ) -> Awaitable[bool]: ...
581
582 def cluster_addslotsrange(
583 self, target_node: "TargetNodesT", *slots: EncodableT
584 ) -> bool | Awaitable[bool]:
585 """
586 Similar to the CLUSTER ADDSLOTS command.
587 The difference between the two commands is that ADDSLOTS takes a list of slots
588 to assign to the node, while ADDSLOTSRANGE takes a list of slot ranges
589 (specified by start and end slots) to assign to the node.
590
591 :target_node: 'ClusterNode'
592 The node to execute the command on
593
594 For more information see https://redis.io/commands/cluster-addslotsrange
595 """
596 return self.execute_command(
597 "CLUSTER ADDSLOTSRANGE", *slots, target_nodes=target_node
598 )
599
600 @overload
601 def cluster_countkeysinslot(self: SyncClientProtocol, slot_id: int) -> int: ...
602
603 @overload
604 def cluster_countkeysinslot(
605 self: AsyncClientProtocol, slot_id: int
606 ) -> Awaitable[int]: ...
607
608 def cluster_countkeysinslot(self, slot_id: int) -> int | Awaitable[int]:
609 """
610 Return the number of local keys in the specified hash slot
611 Send to node based on specified slot_id
612
613 For more information see https://redis.io/commands/cluster-countkeysinslot
614 """
615 return self.execute_command("CLUSTER COUNTKEYSINSLOT", slot_id)
616
617 @overload
618 def cluster_count_failure_report(self: SyncClientProtocol, node_id: str) -> int: ...
619
620 @overload
621 def cluster_count_failure_report(
622 self: AsyncClientProtocol, node_id: str
623 ) -> Awaitable[int]: ...
624
625 def cluster_count_failure_report(self, node_id: str) -> int | Awaitable[int]:
626 """
627 Return the number of failure reports active for a given node
628 Sends to a random node
629
630 For more information see https://redis.io/commands/cluster-count-failure-reports
631 """
632 return self.execute_command("CLUSTER COUNT-FAILURE-REPORTS", node_id)
633
634 def cluster_delslots(self, *slots: EncodableT) -> List[bool]:
635 """
636 Set hash slots as unbound in the cluster.
637 It determines by it self what node the slot is in and sends it there
638
639 Returns a list of the results for each processed slot.
640
641 For more information see https://redis.io/commands/cluster-delslots
642 """
643 return [self.execute_command("CLUSTER DELSLOTS", slot) for slot in slots]
644
645 @overload
646 def cluster_delslotsrange(self: SyncClientProtocol, *slots: EncodableT) -> bool: ...
647
648 @overload
649 def cluster_delslotsrange(
650 self: AsyncClientProtocol, *slots: EncodableT
651 ) -> Awaitable[bool]: ...
652
653 def cluster_delslotsrange(self, *slots: EncodableT) -> bool | Awaitable[bool]:
654 """
655 Similar to the CLUSTER DELSLOTS command.
656 The difference is that CLUSTER DELSLOTS takes a list of hash slots to remove
657 from the node, while CLUSTER DELSLOTSRANGE takes a list of slot ranges to remove
658 from the node.
659
660 For more information see https://redis.io/commands/cluster-delslotsrange
661 """
662 return self.execute_command("CLUSTER DELSLOTSRANGE", *slots)
663
664 @overload
665 def cluster_failover(
666 self: SyncClientProtocol,
667 target_node: "TargetNodesT",
668 option: str | None = None,
669 ) -> bool: ...
670
671 @overload
672 def cluster_failover(
673 self: AsyncClientProtocol,
674 target_node: "TargetNodesT",
675 option: str | None = None,
676 ) -> Awaitable[bool]: ...
677
678 def cluster_failover(
679 self, target_node: "TargetNodesT", option: str | None = None
680 ) -> bool | Awaitable[bool]:
681 """
682 Forces a slave to perform a manual failover of its master
683 Sends to specified node
684
685 :target_node: 'ClusterNode'
686 The node to execute the command on
687
688 For more information see https://redis.io/commands/cluster-failover
689 """
690 if option:
691 if option.upper() not in ["FORCE", "TAKEOVER"]:
692 raise RedisError(
693 f"Invalid option for CLUSTER FAILOVER command: {option}"
694 )
695 else:
696 return self.execute_command(
697 "CLUSTER FAILOVER", option, target_nodes=target_node
698 )
699 else:
700 return self.execute_command("CLUSTER FAILOVER", target_nodes=target_node)
701
702 @overload
703 def cluster_info(
704 self: SyncClientProtocol, target_nodes: "TargetNodesT" | None = None
705 ) -> dict[str, str]: ...
706
707 @overload
708 def cluster_info(
709 self: AsyncClientProtocol, target_nodes: "TargetNodesT" | None = None
710 ) -> Awaitable[dict[str, str]]: ...
711
712 def cluster_info(
713 self, target_nodes: "TargetNodesT" | None = None
714 ) -> dict[str, str] | Awaitable[dict[str, str]]:
715 """
716 Provides info about Redis Cluster node state.
717 The command will be sent to a random node in the cluster if no target
718 node is specified.
719
720 For more information see https://redis.io/commands/cluster-info
721 """
722 return self.execute_command("CLUSTER INFO", target_nodes=target_nodes)
723
724 @overload
725 def cluster_keyslot(self: SyncClientProtocol, key: str) -> int: ...
726
727 @overload
728 def cluster_keyslot(self: AsyncClientProtocol, key: str) -> Awaitable[int]: ...
729
730 def cluster_keyslot(self, key: str) -> int | Awaitable[int]:
731 """
732 Returns the hash slot of the specified key
733 Sends to random node in the cluster
734
735 For more information see https://redis.io/commands/cluster-keyslot
736 """
737 return self.execute_command("CLUSTER KEYSLOT", key)
738
739 @overload
740 def cluster_meet(
741 self: SyncClientProtocol,
742 host: str,
743 port: int,
744 target_nodes: "TargetNodesT" | None = None,
745 ) -> bool: ...
746
747 @overload
748 def cluster_meet(
749 self: AsyncClientProtocol,
750 host: str,
751 port: int,
752 target_nodes: "TargetNodesT" | None = None,
753 ) -> Awaitable[bool]: ...
754
755 def cluster_meet(
756 self, host: str, port: int, target_nodes: "TargetNodesT" | None = None
757 ) -> bool | Awaitable[bool]:
758 """
759 Force a node cluster to handshake with another node.
760 Sends to specified node.
761
762 For more information see https://redis.io/commands/cluster-meet
763 """
764 return self.execute_command(
765 "CLUSTER MEET", host, port, target_nodes=target_nodes
766 )
767
768 @overload
769 def cluster_nodes(self: SyncClientProtocol) -> dict[str, ClusterNodeDetail]: ...
770
771 @overload
772 def cluster_nodes(
773 self: AsyncClientProtocol,
774 ) -> Awaitable[dict[str, ClusterNodeDetail]]: ...
775
776 def cluster_nodes(
777 self,
778 ) -> dict[str, ClusterNodeDetail] | Awaitable[dict[str, ClusterNodeDetail]]:
779 """
780 Get Cluster config for the node.
781 Sends to random node in the cluster
782
783 For more information see https://redis.io/commands/cluster-nodes
784 """
785 return self.execute_command("CLUSTER NODES")
786
787 @overload
788 def cluster_replicate(
789 self: SyncClientProtocol, target_nodes: "TargetNodesT", node_id: str
790 ) -> bool: ...
791
792 @overload
793 def cluster_replicate(
794 self: AsyncClientProtocol, target_nodes: "TargetNodesT", node_id: str
795 ) -> Awaitable[bool]: ...
796
797 def cluster_replicate(
798 self, target_nodes: "TargetNodesT", node_id: str
799 ) -> bool | Awaitable[bool]:
800 """
801 Reconfigure a node as a slave of the specified master node
802
803 For more information see https://redis.io/commands/cluster-replicate
804 """
805 return self.execute_command(
806 "CLUSTER REPLICATE", node_id, target_nodes=target_nodes
807 )
808
809 @overload
810 def cluster_reset(
811 self: SyncClientProtocol,
812 soft: bool = True,
813 target_nodes: "TargetNodesT" | None = None,
814 ) -> bool: ...
815
816 @overload
817 def cluster_reset(
818 self: AsyncClientProtocol,
819 soft: bool = True,
820 target_nodes: "TargetNodesT" | None = None,
821 ) -> Awaitable[bool]: ...
822
823 def cluster_reset(
824 self, soft: bool = True, target_nodes: "TargetNodesT" | None = None
825 ) -> bool | Awaitable[bool]:
826 """
827 Reset a Redis Cluster node
828
829 If 'soft' is True then it will send 'SOFT' argument
830 If 'soft' is False then it will send 'HARD' argument
831
832 For more information see https://redis.io/commands/cluster-reset
833 """
834 return self.execute_command(
835 "CLUSTER RESET", b"SOFT" if soft else b"HARD", target_nodes=target_nodes
836 )
837
838 @overload
839 def cluster_save_config(
840 self: SyncClientProtocol, target_nodes: "TargetNodesT" | None = None
841 ) -> bool: ...
842
843 @overload
844 def cluster_save_config(
845 self: AsyncClientProtocol, target_nodes: "TargetNodesT" | None = None
846 ) -> Awaitable[bool]: ...
847
848 def cluster_save_config(
849 self, target_nodes: "TargetNodesT" | None = None
850 ) -> bool | Awaitable[bool]:
851 """
852 Forces the node to save cluster state on disk
853
854 For more information see https://redis.io/commands/cluster-saveconfig
855 """
856 return self.execute_command("CLUSTER SAVECONFIG", target_nodes=target_nodes)
857
858 @overload
859 def cluster_get_keys_in_slot(
860 self: SyncClientProtocol, slot: int, num_keys: int
861 ) -> list[bytes | str]: ...
862
863 @overload
864 def cluster_get_keys_in_slot(
865 self: AsyncClientProtocol, slot: int, num_keys: int
866 ) -> Awaitable[list[bytes | str]]: ...
867
868 def cluster_get_keys_in_slot(
869 self, slot: int, num_keys: int
870 ) -> list[bytes | str] | Awaitable[list[bytes | str]]:
871 """
872 Returns the number of keys in the specified cluster slot
873
874 For more information see https://redis.io/commands/cluster-getkeysinslot
875 """
876 return self.execute_command("CLUSTER GETKEYSINSLOT", slot, num_keys)
877
878 @overload
879 def cluster_set_config_epoch(
880 self: SyncClientProtocol, epoch: int, target_nodes: "TargetNodesT" | None = None
881 ) -> bool: ...
882
883 @overload
884 def cluster_set_config_epoch(
885 self: AsyncClientProtocol,
886 epoch: int,
887 target_nodes: "TargetNodesT" | None = None,
888 ) -> Awaitable[bool]: ...
889
890 def cluster_set_config_epoch(
891 self, epoch: int, target_nodes: "TargetNodesT" | None = None
892 ) -> bool | Awaitable[bool]:
893 """
894 Set the configuration epoch in a new node
895
896 For more information see https://redis.io/commands/cluster-set-config-epoch
897 """
898 return self.execute_command(
899 "CLUSTER SET-CONFIG-EPOCH", epoch, target_nodes=target_nodes
900 )
901
902 @overload
903 def cluster_setslot(
904 self: SyncClientProtocol,
905 target_node: "TargetNodesT",
906 node_id: str,
907 slot_id: int,
908 state: str,
909 ) -> bool: ...
910
911 @overload
912 def cluster_setslot(
913 self: AsyncClientProtocol,
914 target_node: "TargetNodesT",
915 node_id: str,
916 slot_id: int,
917 state: str,
918 ) -> Awaitable[bool]: ...
919
920 def cluster_setslot(
921 self, target_node: "TargetNodesT", node_id: str, slot_id: int, state: str
922 ) -> bool | Awaitable[bool]:
923 """
924 Bind an hash slot to a specific node
925
926 :target_node: 'ClusterNode'
927 The node to execute the command on
928
929 For more information see https://redis.io/commands/cluster-setslot
930 """
931 if state.upper() in ("IMPORTING", "NODE", "MIGRATING"):
932 return self.execute_command(
933 "CLUSTER SETSLOT", slot_id, state, node_id, target_nodes=target_node
934 )
935 elif state.upper() == "STABLE":
936 raise RedisError('For "stable" state please use cluster_setslot_stable')
937 else:
938 raise RedisError(f"Invalid slot state: {state}")
939
940 @overload
941 def cluster_setslot_stable(self: SyncClientProtocol, slot_id: int) -> bool: ...
942
943 @overload
944 def cluster_setslot_stable(
945 self: AsyncClientProtocol, slot_id: int
946 ) -> Awaitable[bool]: ...
947
948 def cluster_setslot_stable(self, slot_id: int) -> bool | Awaitable[bool]:
949 """
950 Clears migrating / importing state from the slot.
951 It determines by it self what node the slot is in and sends it there.
952
953 For more information see https://redis.io/commands/cluster-setslot
954 """
955 return self.execute_command("CLUSTER SETSLOT", slot_id, "STABLE")
956
957 @overload
958 def cluster_replicas(
959 self: SyncClientProtocol,
960 node_id: str,
961 target_nodes: "TargetNodesT" | None = None,
962 ) -> dict[str, ClusterNodeDetail]: ...
963
964 @overload
965 def cluster_replicas(
966 self: AsyncClientProtocol,
967 node_id: str,
968 target_nodes: "TargetNodesT" | None = None,
969 ) -> Awaitable[dict[str, ClusterNodeDetail]]: ...
970
971 def cluster_replicas(
972 self, node_id: str, target_nodes: "TargetNodesT" | None = None
973 ) -> dict[str, ClusterNodeDetail] | Awaitable[dict[str, ClusterNodeDetail]]:
974 """
975 Provides a list of replica nodes replicating from the specified primary
976 target node.
977
978 For more information see https://redis.io/commands/cluster-replicas
979 """
980 return self.execute_command(
981 "CLUSTER REPLICAS", node_id, target_nodes=target_nodes
982 )
983
984 @overload
985 def cluster_slots(
986 self: SyncClientProtocol, target_nodes: "TargetNodesT" | None = None
987 ) -> list[Any]: ...
988
989 @overload
990 def cluster_slots(
991 self: AsyncClientProtocol, target_nodes: "TargetNodesT" | None = None
992 ) -> Awaitable[list[Any]]: ...
993
994 def cluster_slots(
995 self, target_nodes: "TargetNodesT" | None = None
996 ) -> list[Any] | Awaitable[list[Any]]:
997 """
998 Get array of Cluster slot to node mappings
999
1000 For more information see https://redis.io/commands/cluster-slots
1001 """
1002 return self.execute_command("CLUSTER SLOTS", target_nodes=target_nodes)
1003
1004 @overload
1005 def cluster_shards(
1006 self: SyncClientProtocol, target_nodes: "TargetNodesT" | None = None
1007 ) -> ClusterShardsResponse: ...
1008
1009 @overload
1010 def cluster_shards(
1011 self: AsyncClientProtocol, target_nodes: "TargetNodesT" | None = None
1012 ) -> Awaitable[ClusterShardsResponse]: ...
1013
1014 def cluster_shards(
1015 self, target_nodes: "TargetNodesT" | None = None
1016 ) -> ClusterShardsResponse | Awaitable[ClusterShardsResponse]:
1017 """
1018 Returns details about the shards of the cluster.
1019
1020 For more information see https://redis.io/commands/cluster-shards
1021 """
1022 return self.execute_command("CLUSTER SHARDS", target_nodes=target_nodes)
1023
1024 @overload
1025 def cluster_myshardid(
1026 self: SyncClientProtocol, target_nodes: "TargetNodesT" | None = None
1027 ) -> bytes | str: ...
1028
1029 @overload
1030 def cluster_myshardid(
1031 self: AsyncClientProtocol, target_nodes: "TargetNodesT" | None = None
1032 ) -> Awaitable[bytes | str]: ...
1033
1034 def cluster_myshardid(self, target_nodes: "TargetNodesT" | None = None) -> (
1035 bytes | str
1036 ) | Awaitable[bytes | str]:
1037 """
1038 Returns the shard ID of the node.
1039
1040 For more information see https://redis.io/commands/cluster-myshardid/
1041 """
1042 return self.execute_command("CLUSTER MYSHARDID", target_nodes=target_nodes)
1043
1044 @overload
1045 def cluster_links(
1046 self: SyncClientProtocol, target_node: "TargetNodesT"
1047 ) -> ClusterLinksResponse: ...
1048
1049 @overload
1050 def cluster_links(
1051 self: AsyncClientProtocol, target_node: "TargetNodesT"
1052 ) -> Awaitable[ClusterLinksResponse]: ...
1053
1054 def cluster_links(
1055 self, target_node: "TargetNodesT"
1056 ) -> ClusterLinksResponse | Awaitable[ClusterLinksResponse]:
1057 """
1058 Each node in a Redis Cluster maintains a pair of long-lived TCP link with each
1059 peer in the cluster: One for sending outbound messages towards the peer and one
1060 for receiving inbound messages from the peer.
1061
1062 This command outputs information of all such peer links as an array.
1063
1064 For more information see https://redis.io/commands/cluster-links
1065 """
1066 return self.execute_command("CLUSTER LINKS", target_nodes=target_node)
1067
1068 def cluster_flushslots(self, target_nodes: "TargetNodesT" | None = None) -> None:
1069 raise NotImplementedError(
1070 "CLUSTER FLUSHSLOTS is intentionally not implemented in the client."
1071 )
1072
1073 def cluster_bumpepoch(self, target_nodes: "TargetNodesT" | None = None) -> None:
1074 raise NotImplementedError(
1075 "CLUSTER BUMPEPOCH is intentionally not implemented in the client."
1076 )
1077
1078 def readonly(self, target_nodes: "TargetNodesT" | None = None) -> ResponseT:
1079 """
1080 Enables read queries.
1081 The command will be sent to the default cluster node if target_nodes is
1082 not specified.
1083
1084 For more information see https://redis.io/commands/readonly
1085 """
1086 if target_nodes == "replicas" or target_nodes == "all":
1087 # read_from_replicas will only be enabled if the READONLY command
1088 # is sent to all replicas
1089 self.read_from_replicas = True
1090 return self.execute_command("READONLY", target_nodes=target_nodes)
1091
1092 def readwrite(self, target_nodes: "TargetNodesT" | None = None) -> ResponseT:
1093 """
1094 Disables read queries.
1095 The command will be sent to the default cluster node if target_nodes is
1096 not specified.
1097
1098 For more information see https://redis.io/commands/readwrite
1099 """
1100 # Reset read from replicas flag
1101 self.read_from_replicas = False
1102 return self.execute_command("READWRITE", target_nodes=target_nodes)
1103
1104 @deprecated_function(
1105 version="7.2.0",
1106 reason="Use client-side caching feature instead.",
1107 )
1108 def client_tracking_on(
1109 self,
1110 clientid: int | None = None,
1111 prefix: Sequence[KeyT] = [],
1112 bcast: bool = False,
1113 optin: bool = False,
1114 optout: bool = False,
1115 noloop: bool = False,
1116 target_nodes: "TargetNodesT" | None = "all",
1117 ) -> ResponseT:
1118 """
1119 Enables the tracking feature of the Redis server, that is used
1120 for server assisted client side caching.
1121
1122 When clientid is provided - in target_nodes only the node that owns the
1123 connection with this id should be provided.
1124 When clientid is not provided - target_nodes can be any node.
1125
1126 For more information see https://redis.io/commands/client-tracking
1127 """
1128 return self.client_tracking(
1129 True,
1130 clientid,
1131 prefix,
1132 bcast,
1133 optin,
1134 optout,
1135 noloop,
1136 target_nodes=target_nodes,
1137 )
1138
1139 @deprecated_function(
1140 version="7.2.0",
1141 reason="Use client-side caching feature instead.",
1142 )
1143 def client_tracking_off(
1144 self,
1145 clientid: int | None = None,
1146 prefix: Sequence[KeyT] = [],
1147 bcast: bool = False,
1148 optin: bool = False,
1149 optout: bool = False,
1150 noloop: bool = False,
1151 target_nodes: "TargetNodesT" | None = "all",
1152 ) -> ResponseT:
1153 """
1154 Disables the tracking feature of the Redis server, that is used
1155 for server assisted client side caching.
1156
1157 When clientid is provided - in target_nodes only the node that owns the
1158 connection with this id should be provided.
1159 When clientid is not provided - target_nodes can be any node.
1160
1161 For more information see https://redis.io/commands/client-tracking
1162 """
1163 return self.client_tracking(
1164 False,
1165 clientid,
1166 prefix,
1167 bcast,
1168 optin,
1169 optout,
1170 noloop,
1171 target_nodes=target_nodes,
1172 )
1173
1174 def hotkeys_start(
1175 self,
1176 metrics: List[HotkeysMetricsTypes],
1177 count: int | None = None,
1178 duration: int | None = None,
1179 sample_ratio: int | None = None,
1180 slots: List[int] | None = None,
1181 **kwargs,
1182 ) -> str | bytes:
1183 """
1184 Cluster client does not support hotkeys command. Please use the non-cluster client.
1185
1186 For more information see https://redis.io/commands/hotkeys-start
1187 """
1188 raise NotImplementedError(
1189 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client."
1190 )
1191
1192 def hotkeys_stop(self, **kwargs) -> str | bytes:
1193 """
1194 Cluster client does not support hotkeys command. Please use the non-cluster client.
1195
1196 For more information see https://redis.io/commands/hotkeys-stop
1197 """
1198 raise NotImplementedError(
1199 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client."
1200 )
1201
1202 def hotkeys_reset(self, **kwargs) -> str | bytes:
1203 """
1204 Cluster client does not support hotkeys command. Please use the non-cluster client.
1205
1206 For more information see https://redis.io/commands/hotkeys-reset
1207 """
1208 raise NotImplementedError(
1209 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client."
1210 )
1211
1212 def hotkeys_get(self, **kwargs) -> list[dict[str | bytes, Any]]:
1213 """
1214 Cluster client does not support hotkeys command. Please use the non-cluster client.
1215
1216 For more information see https://redis.io/commands/hotkeys-get
1217 """
1218 raise NotImplementedError(
1219 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client."
1220 )
1221
1222
1223class AsyncClusterManagementCommands(
1224 ClusterManagementCommands, AsyncManagementCommands
1225):
1226 """
1227 A class for Redis Cluster management commands
1228
1229 The class inherits from Redis's core ManagementCommands class and do the
1230 required adjustments to work with cluster mode
1231 """
1232
1233 async def cluster_delslots(self, *slots: EncodableT) -> List[bool]:
1234 """
1235 Set hash slots as unbound in the cluster.
1236 It determines by it self what node the slot is in and sends it there
1237
1238 Returns a list of the results for each processed slot.
1239
1240 For more information see https://redis.io/commands/cluster-delslots
1241 """
1242 return await asyncio.gather(
1243 *(
1244 asyncio.create_task(self.execute_command("CLUSTER DELSLOTS", slot))
1245 for slot in slots
1246 )
1247 )
1248
1249 @deprecated_function(
1250 version="7.2.0",
1251 reason="Use client-side caching feature instead.",
1252 )
1253 async def client_tracking_on(
1254 self,
1255 clientid: int | None = None,
1256 prefix: Sequence[KeyT] = [],
1257 bcast: bool = False,
1258 optin: bool = False,
1259 optout: bool = False,
1260 noloop: bool = False,
1261 target_nodes: "TargetNodesT" | None = "all",
1262 ) -> ResponseT:
1263 """
1264 Enables the tracking feature of the Redis server, that is used
1265 for server assisted client side caching.
1266
1267 When clientid is provided - in target_nodes only the node that owns the
1268 connection with this id should be provided.
1269 When clientid is not provided - target_nodes can be any node.
1270
1271 For more information see https://redis.io/commands/client-tracking
1272 """
1273 return await self.client_tracking(
1274 True,
1275 clientid,
1276 prefix,
1277 bcast,
1278 optin,
1279 optout,
1280 noloop,
1281 target_nodes=target_nodes,
1282 )
1283
1284 @deprecated_function(
1285 version="7.2.0",
1286 reason="Use client-side caching feature instead.",
1287 )
1288 async def client_tracking_off(
1289 self,
1290 clientid: int | None = None,
1291 prefix: Sequence[KeyT] = [],
1292 bcast: bool = False,
1293 optin: bool = False,
1294 optout: bool = False,
1295 noloop: bool = False,
1296 target_nodes: "TargetNodesT" | None = "all",
1297 ) -> ResponseT:
1298 """
1299 Disables the tracking feature of the Redis server, that is used
1300 for server assisted client side caching.
1301
1302 When clientid is provided - in target_nodes only the node that owns the
1303 connection with this id should be provided.
1304 When clientid is not provided - target_nodes can be any node.
1305
1306 For more information see https://redis.io/commands/client-tracking
1307 """
1308 return await self.client_tracking(
1309 False,
1310 clientid,
1311 prefix,
1312 bcast,
1313 optin,
1314 optout,
1315 noloop,
1316 target_nodes=target_nodes,
1317 )
1318
1319 async def hotkeys_start(
1320 self,
1321 metrics: List[HotkeysMetricsTypes],
1322 count: int | None = None,
1323 duration: int | None = None,
1324 sample_ratio: int | None = None,
1325 slots: List[int] | None = None,
1326 **kwargs,
1327 ) -> Awaitable[str | bytes]:
1328 """
1329 Cluster client does not support hotkeys command. Please use the non-cluster client.
1330
1331 For more information see https://redis.io/commands/hotkeys-start
1332 """
1333 raise NotImplementedError(
1334 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client."
1335 )
1336
1337 async def hotkeys_stop(self, **kwargs) -> Awaitable[str | bytes]:
1338 """
1339 Cluster client does not support hotkeys command. Please use the non-cluster client.
1340
1341 For more information see https://redis.io/commands/hotkeys-stop
1342 """
1343 raise NotImplementedError(
1344 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client."
1345 )
1346
1347 async def hotkeys_reset(self, **kwargs) -> Awaitable[str | bytes]:
1348 """
1349 Cluster client does not support hotkeys command. Please use the non-cluster client.
1350
1351 For more information see https://redis.io/commands/hotkeys-reset
1352 """
1353 raise NotImplementedError(
1354 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client."
1355 )
1356
1357 async def hotkeys_get(self, **kwargs) -> Awaitable[list[dict[str | bytes, Any]]]:
1358 """
1359 Cluster client does not support hotkeys command. Please use the non-cluster client.
1360
1361 For more information see https://redis.io/commands/hotkeys-get
1362 """
1363 raise NotImplementedError(
1364 "HOTKEYS commands are not supported in cluster mode. Please use the non-cluster client."
1365 )
1366
1367
1368class ClusterDataAccessCommands(DataAccessCommands):
1369 """
1370 A class for Redis Cluster Data Access Commands
1371
1372 The class inherits from Redis's core DataAccessCommand class and do the
1373 required adjustments to work with cluster mode
1374 """
1375
1376 @overload
1377 def stralgo(
1378 self: SyncClientProtocol,
1379 algo: Literal["LCS"],
1380 value1: KeyT,
1381 value2: KeyT,
1382 specific_argument: Literal["strings"] | Literal["keys"] = "strings",
1383 len: bool = False,
1384 idx: bool = False,
1385 minmatchlen: int | None = None,
1386 withmatchlen: bool = False,
1387 **kwargs,
1388 ) -> StralgoResponse: ...
1389
1390 @overload
1391 def stralgo(
1392 self: AsyncClientProtocol,
1393 algo: Literal["LCS"],
1394 value1: KeyT,
1395 value2: KeyT,
1396 specific_argument: Literal["strings"] | Literal["keys"] = "strings",
1397 len: bool = False,
1398 idx: bool = False,
1399 minmatchlen: int | None = None,
1400 withmatchlen: bool = False,
1401 **kwargs,
1402 ) -> Awaitable[StralgoResponse]: ...
1403
1404 def stralgo(
1405 self,
1406 algo: Literal["LCS"],
1407 value1: KeyT,
1408 value2: KeyT,
1409 specific_argument: Literal["strings"] | Literal["keys"] = "strings",
1410 len: bool = False,
1411 idx: bool = False,
1412 minmatchlen: int | None = None,
1413 withmatchlen: bool = False,
1414 **kwargs,
1415 ) -> StralgoResponse | Awaitable[StralgoResponse]:
1416 """
1417 Implements complex algorithms that operate on strings.
1418 Right now the only algorithm implemented is the LCS algorithm
1419 (longest common substring). However new algorithms could be
1420 implemented in the future.
1421
1422 ``algo`` Right now must be LCS
1423 ``value1`` and ``value2`` Can be two strings or two keys
1424 ``specific_argument`` Specifying if the arguments to the algorithm
1425 will be keys or strings. strings is the default.
1426 ``len`` Returns just the len of the match.
1427 ``idx`` Returns the match positions in each string.
1428 ``minmatchlen`` Restrict the list of matches to the ones of a given
1429 minimal length. Can be provided only when ``idx`` set to True.
1430 ``withmatchlen`` Returns the matches with the len of the match.
1431 Can be provided only when ``idx`` set to True.
1432
1433 For more information see https://redis.io/commands/stralgo
1434 """
1435 target_nodes = kwargs.pop("target_nodes", None)
1436 if specific_argument == "strings" and target_nodes is None:
1437 target_nodes = "default-node"
1438 kwargs.update({"target_nodes": target_nodes})
1439 return super().stralgo(
1440 algo,
1441 value1,
1442 value2,
1443 specific_argument,
1444 len,
1445 idx,
1446 minmatchlen,
1447 withmatchlen,
1448 **kwargs,
1449 )
1450
1451 def scan_iter(
1452 self,
1453 match: PatternT | None = None,
1454 count: int | None = None,
1455 _type: str | None = None,
1456 **kwargs,
1457 ) -> Iterator:
1458 # Do the first query with cursor=0 for all nodes
1459 cursors, data = self.scan(match=match, count=count, _type=_type, **kwargs)
1460 yield from data
1461
1462 cursors = {name: cursor for name, cursor in cursors.items() if cursor != 0}
1463 if cursors:
1464 # Get nodes by name
1465 nodes = {name: self.get_node(node_name=name) for name in cursors.keys()}
1466
1467 # Iterate over each node till its cursor is 0
1468 kwargs.pop("target_nodes", None)
1469 while cursors:
1470 for name, cursor in cursors.items():
1471 cur, data = self.scan(
1472 cursor=cursor,
1473 match=match,
1474 count=count,
1475 _type=_type,
1476 target_nodes=nodes[name],
1477 **kwargs,
1478 )
1479 yield from data
1480 cursors[name] = cur[name]
1481
1482 cursors = {
1483 name: cursor for name, cursor in cursors.items() if cursor != 0
1484 }
1485
1486
1487class AsyncClusterDataAccessCommands(
1488 ClusterDataAccessCommands, AsyncDataAccessCommands
1489):
1490 """
1491 A class for Redis Cluster Data Access Commands
1492
1493 The class inherits from Redis's core DataAccessCommand class and do the
1494 required adjustments to work with cluster mode
1495 """
1496
1497 async def scan_iter(
1498 self,
1499 match: PatternT | None = None,
1500 count: int | None = None,
1501 _type: str | None = None,
1502 **kwargs,
1503 ) -> AsyncIterator:
1504 # Do the first query with cursor=0 for all nodes
1505 cursors, data = await self.scan(match=match, count=count, _type=_type, **kwargs)
1506 for value in data:
1507 yield value
1508
1509 cursors = {name: cursor for name, cursor in cursors.items() if cursor != 0}
1510 if cursors:
1511 # Get nodes by name
1512 nodes = {name: self.get_node(node_name=name) for name in cursors.keys()}
1513
1514 # Iterate over each node till its cursor is 0
1515 kwargs.pop("target_nodes", None)
1516 while cursors:
1517 for name, cursor in cursors.items():
1518 cur, data = await self.scan(
1519 cursor=cursor,
1520 match=match,
1521 count=count,
1522 _type=_type,
1523 target_nodes=nodes[name],
1524 **kwargs,
1525 )
1526 for value in data:
1527 yield value
1528 cursors[name] = cur[name]
1529
1530 cursors = {
1531 name: cursor for name, cursor in cursors.items() if cursor != 0
1532 }
1533
1534
1535class RedisClusterCommands(
1536 ClusterMultiKeyCommands,
1537 ClusterManagementCommands,
1538 ACLCommands,
1539 PubSubCommands,
1540 ClusterDataAccessCommands,
1541 ScriptCommands,
1542 FunctionCommands,
1543 ModuleCommands,
1544 RedisModuleCommands,
1545):
1546 """
1547 A class for all Redis Cluster commands
1548
1549 For key-based commands, the target node(s) will be internally determined
1550 by the keys' hash slot.
1551 Non-key-based commands can be executed with the 'target_nodes' argument to
1552 target specific nodes. By default, if target_nodes is not specified, the
1553 command will be executed on the default cluster node.
1554
1555 :param target_nodes: type can be one of the following:
1556 - nodes flag: ALL_NODES, PRIMARIES, REPLICAS, RANDOM
1557 - 'ClusterNode'
1558 - 'list(ClusterNodes)'
1559 - 'dict(any:clusterNodes)'
1560
1561 for example:
1562 r.cluster_info(target_nodes=RedisCluster.ALL_NODES)
1563 """
1564
1565
1566class AsyncRedisClusterCommands(
1567 AsyncClusterMultiKeyCommands,
1568 AsyncClusterManagementCommands,
1569 AsyncACLCommands,
1570 PubSubCommands,
1571 AsyncClusterDataAccessCommands,
1572 AsyncScriptCommands,
1573 AsyncFunctionCommands,
1574 AsyncModuleCommands,
1575 AsyncRedisModuleCommands,
1576):
1577 """
1578 A class for all Redis Cluster commands
1579
1580 For key-based commands, the target node(s) will be internally determined
1581 by the keys' hash slot.
1582 Non-key-based commands can be executed with the 'target_nodes' argument to
1583 target specific nodes. By default, if target_nodes is not specified, the
1584 command will be executed on the default cluster node.
1585
1586 :param target_nodes: type can be one of the following:
1587 - nodes flag: ALL_NODES, PRIMARIES, REPLICAS, RANDOM
1588 - 'ClusterNode'
1589 - 'list(ClusterNodes)'
1590 - 'dict(any:clusterNodes)'
1591
1592 for example:
1593 r.cluster_info(target_nodes=RedisCluster.ALL_NODES)
1594 """