summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorAndy McCurdy <andy@andymccurdy.com>2011-05-22 21:26:44 -0700
committerAndy McCurdy <andy@andymccurdy.com>2011-05-22 21:26:44 -0700
commit2810533c1b77b3f8028cc5611359c4dd10492b93 (patch)
treef38c5ee2933b31645623e8139073b3b6fc3380d1
parentd1cd365006fdcce4739adfa0d8a1bd764db0d142 (diff)
downloadredis-py-2810533c1b77b3f8028cc5611359c4dd10492b93.tar.gz
allow shard hints to be passed to pipeline and pubsub objects. a smart connection pool could use these hints to determine the correct shard to run on.
removed all previously deprecated parameters and commands
-rw-r--r--redis/client.py143
1 files changed, 53 insertions, 90 deletions
diff --git a/redis/client.py b/redis/client.py
index 41b80d9..98b8a0b 100644
--- a/redis/client.py
+++ b/redis/client.py
@@ -13,26 +13,17 @@ from redis.exceptions import (
WatchError,
)
-
-def list_or_args(command, keys, args):
+def list_or_args(keys, args):
# returns a single list combining keys and args
- # if keys is not a list or args has items, issue a
- # deprecation warning
- oldapi = bool(args)
try:
i = iter(keys)
# a string can be iterated, but indicates
# keys wasn't passed as a list
if isinstance(keys, basestring):
- oldapi = True
+ keys = [keys]
except TypeError:
- oldapi = True
keys = [keys]
- if oldapi:
- warnings.warn(DeprecationWarning(
- "Passing *args to Redis.%s has been deprecated. "
- "Pass an iterable to ``keys`` instead" % command
- ))
+ if args:
keys.extend(args)
return keys
@@ -184,15 +175,15 @@ class Redis(object):
encoding_errors=errors
)
- def pipeline(self, transaction=True):
+ def pipeline(self, transaction=True, shard_hint=None):
"""
Return a new pipeline object that can queue multiple commands for
later execution. ``transaction`` indicates whether all commands
- should be executed atomically. Apart from multiple atomic operations,
- pipelines are useful for batch loading of data as they reduce the
- number of back and forth network operations between client and server.
+ should be executed atomically. Apart from making a group of operations
+ atomic, pipelines are useful for reducing the back-and-forth overhead
+ between the client and server.
"""
- return Pipeline(self.connection_pool, transaction)
+ return Pipeline(self.connection_pool, transaction, shard_hint)
def lock(self, name, timeout=None, sleep=0.1):
"""
@@ -208,11 +199,17 @@ class Redis(object):
"""
return Lock(self, name, timeout=timeout, sleep=sleep)
- def pubsub(self):
- return PubSub(self.connection_pool)
+ def pubsub(self, shard_hint=None):
+ """
+ Return a Publish/Subscribe object. With this object, you can
+ subscribe to channels and listen for messages that get published to
+ them.
+ """
+ return PubSub(self.connection_pool, shard_hint)
#### COMMAND EXECUTION AND PROTOCOL PARSING ####
def execute_command(self, *args, **options):
+ "Execute a command and return a parsed response"
command_name = args[0]
connection = self.connection_pool.get_connection(command_name)
try:
@@ -261,14 +258,6 @@ class Redis(object):
return self.execute_command('DEL', *names)
__delitem__ = delete
- def flush(self, all_dbs=False):
- warnings.warn(DeprecationWarning(
- "'flush' has been deprecated. "
- "Use Redis.flushdb() or Redis.flushall() instead"))
- if all_dbs:
- return self.flushall()
- return self.flushdb()
-
def flushall(self):
"Delete all keys in all databases on the current host"
return self.execute_command('FLUSHALL')
@@ -300,7 +289,13 @@ class Redis(object):
return self.execute_command('SAVE')
def select(self, db):
- "Select a differnet Redis database"
+ """
+ Select a differnet Redis database.
+
+ WARNING: this could have severe consequences for pooled connections.
+ It's highly advised to use a separate connection pool and client
+ instance to work with multiple databases. Use this at your own risk.
+ """
return self.execute_command('SELECT', db)
def shutdown(self):
@@ -370,8 +365,7 @@ class Redis(object):
value = self.get(name)
if value:
return value
- else:
- raise KeyError(name)
+ raise KeyError(name)
def getbit(self, name, offset):
"Returns a boolean indicating the value of ``offset`` in ``name``"
@@ -398,10 +392,8 @@ class Redis(object):
def mget(self, keys, *args):
"""
Returns a list of values ordered identically to ``keys``
-
- * Passing *args to this method has been deprecated *
"""
- keys = list_or_args('mget', keys, args)
+ keys = list_or_args(keys, args)
return self.execute_command('MGET', *keys)
def mset(self, mapping):
@@ -433,29 +425,17 @@ class Redis(object):
"Returns the name of a random key"
return self.execute_command('RANDOMKEY')
- def rename(self, src, dst, **kwargs):
+ def rename(self, src, dst):
"""
Rename key ``src`` to ``dst``
-
- * The following flags have been deprecated *
- If ``preserve`` is True, rename the key only if the destination name
- doesn't already exist
- """
- if kwargs:
- if 'preserve' in kwargs:
- warnings.warn(DeprecationWarning(
- "preserve option to 'rename' is deprecated, "
- "use Redis.renamenx instead"))
- if kwargs['preserve']:
- return self.renamenx(src, dst)
+ """
return self.execute_command('RENAME', src, dst)
def renamenx(self, src, dst):
"Rename key ``src`` to ``dst`` if ``dst`` doesn't already exist"
return self.execute_command('RENAMENX', src, dst)
-
- def set(self, name, value, **kwargs):
+ def set(self, name, value):
"""
Set the value at key ``name`` to ``value``
@@ -465,19 +445,6 @@ class Redis(object):
If ``getset`` is True, set the value only if key doesn't already exist
and return the resulting value of key
"""
- if kwargs:
- if 'getset' in kwargs:
- warnings.warn(DeprecationWarning(
- "getset option to 'set' is deprecated, "
- "use Redis.getset() instead"))
- if kwargs['getset']:
- return self.getset(name, value)
- if 'preserve' in kwargs:
- warnings.warn(DeprecationWarning(
- "preserve option to 'set' is deprecated, "
- "use Redis.setnx() instead"))
- if kwargs['preserve']:
- return self.setnx(name, value)
return self.execute_command('SET', name, value)
__setitem__ = set
@@ -778,7 +745,7 @@ class Redis(object):
def sdiff(self, keys, *args):
"Return the difference of sets specified by ``keys``"
- keys = list_or_args('sdiff', keys, args)
+ keys = list_or_args(keys, args)
return self.execute_command('SDIFF', *keys)
def sdiffstore(self, dest, keys, *args):
@@ -786,12 +753,12 @@ class Redis(object):
Store the difference of sets specified by ``keys`` into a new
set named ``dest``. Returns the number of keys in the new set.
"""
- keys = list_or_args('sdiffstore', keys, args)
+ keys = list_or_args(keys, args)
return self.execute_command('SDIFFSTORE', dest, *keys)
def sinter(self, keys, *args):
"Return the intersection of sets specified by ``keys``"
- keys = list_or_args('sinter', keys, args)
+ keys = list_or_args(keys, args)
return self.execute_command('SINTER', *keys)
def sinterstore(self, dest, keys, *args):
@@ -799,7 +766,7 @@ class Redis(object):
Store the intersection of sets specified by ``keys`` into a new
set named ``dest``. Returns the number of keys in the new set.
"""
- keys = list_or_args('sinterstore', keys, args)
+ keys = list_or_args(keys, args)
return self.execute_command('SINTERSTORE', dest, *keys)
def sismember(self, name, value):
@@ -828,7 +795,7 @@ class Redis(object):
def sunion(self, keys, *args):
"Return the union of sets specifiued by ``keys``"
- keys = list_or_args('sunion', keys, args)
+ keys = list_or_args(keys, args)
return self.execute_command('SUNION', *keys)
def sunionstore(self, dest, keys, *args):
@@ -836,7 +803,7 @@ class Redis(object):
Store the union of sets specified by ``keys`` into a new
set named ``dest``. Returns the number of keys in the new set.
"""
- keys = list_or_args('sunionstore', keys, args)
+ keys = list_or_args(keys, args)
return self.execute_command('SUNIONSTORE', dest, *keys)
@@ -993,12 +960,6 @@ class Redis(object):
"Return the score of element ``value`` in sorted set ``name``"
return self.execute_command('ZSCORE', name, value)
- def zunion(self, dest, keys, aggregate=None):
- warnings.warn(DeprecationWarning(
- "Redis.zunion has been deprecated, use Redis.zunionstore instead"
- ))
- return self.zunionstore(dest, keys, aggregate)
-
def zunionstore(self, dest, keys, aggregate=None):
"""
Union multiple sorted sets specified by ``keys`` into
@@ -1094,8 +1055,9 @@ class Redis(object):
class PubSub(object):
- def __init__(self, connection_pool):
+ def __init__(self, connection_pool, shard_hint=None):
self.connection_pool = connection_pool
+ self.shard_hint = shard_hint
self.connection = None
self.channels = set()
self.patterns = set()
@@ -1107,7 +1069,10 @@ class PubSub(object):
def execute_command(self, *args, **kwargs):
"Execute a publish/subscribe command"
if self.connection is None:
- self.connection = self.connection_pool.get_connection('pubsub')
+ self.connection = self.connection_pool.get_connection(
+ 'pubsub',
+ self.shard_hint
+ )
connection = self.connection
try:
connection.send_command(*args)
@@ -1225,9 +1190,10 @@ class Pipeline(Redis):
ResponseError exceptions, such as those raised when issuing a command
on a key of a different datatype.
"""
- def __init__(self, connection_pool, transaction):
+ def __init__(self, connection_pool, transaction, shard_hint):
self.connection_pool = connection_pool
self.transaction = transaction
+ self.shard_hint = shard_hint
self.reset()
def reset(self):
@@ -1251,20 +1217,20 @@ class Pipeline(Redis):
return self
def _execute_transaction(self, commands):
- connection = self.connection_pool.get_connection('MULTI')
+ conn = self.connection_pool.get_connection('MULTI', self.shard_hint)
try:
- all_cmds = ''.join(starmap(connection.pack_command,
+ all_cmds = ''.join(starmap(conn.pack_command,
[args for args, options in commands]))
- connection.send_packed_command(all_cmds)
+ conn.send_packed_command(all_cmds)
# we don't care about the multi/exec any longer
commands = commands[1:-1]
# parse off the response for MULTI and all commands prior to EXEC.
# the only data we care about is the response the EXEC
# which is the last command
for i in range(len(commands)+1):
- _ = self.parse_response(connection, '_')
+ _ = self.parse_response(conn, '_')
# parse the EXEC.
- response = self.parse_response(connection, '_')
+ response = self.parse_response(conn, '_')
if response is None:
raise WatchError("Watched variable changed.")
@@ -1283,19 +1249,19 @@ class Pipeline(Redis):
data.append(r)
return data
finally:
- self.connection_pool.release(connection)
+ self.connection_pool.release(conn)
def _execute_pipeline(self, commands):
# build up all commands into a single request to increase network perf
- connection = self.connection_pool.get_connection('MULTI')
+ conn = self.connection_pool.get_connection('MULTI', self.shard_hint)
try:
- all_cmds = ''.join(starmap(connection.pack_command,
+ all_cmds = ''.join(starmap(conn.pack_command,
[args for args, options in commands]))
- connection.send_packed_command(all_cmds)
- return [self.parse_response(connection, args[0], **options)
+ conn.send_packed_command(all_cmds)
+ return [self.parse_response(conn, args[0], **options)
for args, options in commands]
finally:
- self.connection_pool.release(connection)
+ self.connection_pool.release(conn)
def execute(self):
"Execute all the commands in the current pipeline"
@@ -1312,9 +1278,6 @@ class Pipeline(Redis):
connection.disconnect()
return execute(stack)
- def select(self, *args, **kwargs):
- raise RedisError("Cannot select a different database from a pipeline")
-
class LockError(RedisError):
"Errors thrown from the Lock"
pass