From 7dd8166643e9868fe930f6f703b9a058c6441545 Mon Sep 17 00:00:00 2001 From: Zsolt Dollenstein Date: Fri, 8 Aug 2014 15:12:47 +0200 Subject: [PATCH 1/3] support redis over tcp --- src/analyzer/analyzer.py | 6 +++--- src/horizon/roomba.py | 7 ++++--- src/horizon/worker.py | 7 ++++--- src/utils.py | 10 ++++++++++ src/webapp/webapp.py | 5 +++-- utils/continuity.py | 4 ++-- utils/seed_data.py | 4 ++-- 7 files changed, 28 insertions(+), 15 deletions(-) create mode 100644 src/utils.py diff --git a/src/analyzer/analyzer.py b/src/analyzer/analyzer.py index 51bf4d95..bf7b74e3 100644 --- a/src/analyzer/analyzer.py +++ b/src/analyzer/analyzer.py @@ -1,6 +1,5 @@ import logging from Queue import Empty -from redis import StrictRedis from time import time, sleep from threading import Thread from collections import defaultdict @@ -16,6 +15,7 @@ from alerters import trigger_alert from algorithms import run_selected_algorithm from algorithm_exceptions import * +from utils import redis_conn logger = logging.getLogger("AnalyzerLog") @@ -26,7 +26,7 @@ def __init__(self, parent_pid): Initialize the Analyzer """ super(Analyzer, self).__init__() - self.redis_conn = StrictRedis(unix_socket_path = settings.REDIS_SOCKET_PATH) + self.redis_conn = redis_conn() self.daemon = True self.parent_pid = parent_pid self.current_pid = getpid() @@ -138,7 +138,7 @@ def run(self): except: logger.error('skyline can\'t connect to redis at socket path %s' % settings.REDIS_SOCKET_PATH) sleep(10) - self.redis_conn = StrictRedis(unix_socket_path = settings.REDIS_SOCKET_PATH) + self.redis_conn = redis_conn() continue # Discover unique metrics diff --git a/src/horizon/roomba.py b/src/horizon/roomba.py index 3c5e466e..c5c9840d 100644 --- a/src/horizon/roomba.py +++ b/src/horizon/roomba.py @@ -1,10 +1,11 @@ from os import kill -from redis import StrictRedis, WatchError +from redis import WatchError from multiprocessing import Process from threading import Thread from msgpack import Unpacker, packb from types import TupleType from time import time, sleep +from utils import redis_conn import logging import settings @@ -18,7 +19,7 @@ class Roomba(Thread): """ def __init__(self, parent_pid, skip_mini): super(Roomba, self).__init__() - self.redis_conn = StrictRedis(unix_socket_path = settings.REDIS_SOCKET_PATH) + self.redis_conn = redis_conn() self.daemon = True self.parent_pid = parent_pid self.skip_mini = skip_mini @@ -160,7 +161,7 @@ def run(self): except: logger.error('roomba can\'t connect to redis at socket path %s' % settings.REDIS_SOCKET_PATH) sleep(10) - self.redis_conn = StrictRedis(unix_socket_path = settings.REDIS_SOCKET_PATH) + self.redis_conn = redis_conn() continue # Spawn processes diff --git a/src/horizon/worker.py b/src/horizon/worker.py index 6bdd44ff..8260fcec 100644 --- a/src/horizon/worker.py +++ b/src/horizon/worker.py @@ -1,9 +1,10 @@ from os import kill, system -from redis import StrictRedis, WatchError +from redis import WatchError from multiprocessing import Process from Queue import Empty from msgpack import packb from time import time, sleep +from utils import redis_conn import logging import socket @@ -19,7 +20,7 @@ class Worker(Process): """ def __init__(self, queue, parent_pid, skip_mini, canary=False): super(Worker, self).__init__() - self.redis_conn = StrictRedis(unix_socket_path = settings.REDIS_SOCKET_PATH) + self.redis_conn = redis_conn() self.q = queue self.parent_pid = parent_pid self.daemon = True @@ -76,7 +77,7 @@ def run(self): except: logger.error('worker can\'t connect to redis at socket path %s' % settings.REDIS_SOCKET_PATH) sleep(10) - self.redis_conn = StrictRedis(unix_socket_path = settings.REDIS_SOCKET_PATH) + self.redis_conn = redis_conn() pipe = self.redis_conn.pipeline() continue diff --git a/src/utils.py b/src/utils.py new file mode 100644 index 00000000..425f912d --- /dev/null +++ b/src/utils.py @@ -0,0 +1,10 @@ +from redis import StrictRedis + +import settings + +def redis_conn(): + if settings.REDIS_HOST_PORT is not None: + (host, port) = settings.REDIS_HOST_PORT + return StrictRedis(host=host, port=port) + else: + return StrictRedis(unix_socket_path=settings.REDIS_SOCKET_PATH) diff --git a/src/webapp/webapp.py b/src/webapp/webapp.py index 6751d761..537befc3 100644 --- a/src/webapp/webapp.py +++ b/src/webapp/webapp.py @@ -1,4 +1,3 @@ -import redis import logging import simplejson as json import sys @@ -7,11 +6,13 @@ from daemon import runner from os.path import dirname, abspath +from utils import redis_conn + # add the shared settings file to namespace sys.path.insert(0, dirname(dirname(abspath(__file__)))) import settings -REDIS_CONN = redis.StrictRedis(unix_socket_path=settings.REDIS_SOCKET_PATH) +REDIS_CONN = redis_conn() app = Flask(__name__) app.config['PROPAGATE_EXCEPTIONS'] = True diff --git a/utils/continuity.py b/utils/continuity.py index c25ef12c..defb62a1 100644 --- a/utils/continuity.py +++ b/utils/continuity.py @@ -1,9 +1,9 @@ -import redis import msgpack import sys import time from os.path import dirname, abspath +from utils import redis_conn # add the shared settings file to namespace sys.path.insert(0, ''.join((dirname(dirname(abspath(__file__))), "/src"))) import settings @@ -12,7 +12,7 @@ def check_continuity(metric, mini = False): - r = redis.StrictRedis(unix_socket_path=settings.REDIS_SOCKET_PATH) + r = redis_conn() if mini: raw_series = r.get(settings.MINI_NAMESPACE + metric) else: diff --git a/utils/seed_data.py b/utils/seed_data.py index 087e39a1..e83057eb 100755 --- a/utils/seed_data.py +++ b/utils/seed_data.py @@ -10,8 +10,8 @@ from multiprocessing import Manager, Process, log_to_stderr from struct import Struct, pack -import redis import msgpack +from utils import redis_conn # Get the current working directory of this file. # http://stackoverflow.com/a/4060259/120999 @@ -44,7 +44,7 @@ def seed(): sock.sendto(packet, (socket.gethostname(), settings.UDP_PORT)) print "Connecting to Redis..." - r = redis.StrictRedis(unix_socket_path=settings.REDIS_SOCKET_PATH) + r = redis_conn() time.sleep(5) try: From e3c52c892e15fe6e2247b2f2d4c7e41aa95b9ff6 Mon Sep 17 00:00:00 2001 From: Zsolt Dollenstein Date: Fri, 8 Aug 2014 15:30:22 +0200 Subject: [PATCH 2/3] change error messages to include proper redis connection string --- src/analyzer/algorithms.py | 5 ++--- src/analyzer/analyzer.py | 4 ++-- src/horizon/roomba.py | 4 ++-- src/horizon/worker.py | 4 ++-- src/utils.py | 8 ++++++++ 5 files changed, 16 insertions(+), 9 deletions(-) diff --git a/src/analyzer/algorithms.py b/src/analyzer/algorithms.py index 8ee5c80f..43b3b398 100644 --- a/src/analyzer/algorithms.py +++ b/src/analyzer/algorithms.py @@ -6,7 +6,7 @@ import logging from time import time from msgpack import unpackb, packb -from redis import StrictRedis +import utils from settings import ( ALGORITHMS, @@ -15,7 +15,6 @@ MAX_TOLERABLE_BOREDOM, MIN_TOLERABLE_LENGTH, STALE_PERIOD, - REDIS_SOCKET_PATH, ENABLE_SECOND_ORDER, BOREDOM_SET_SIZE, ) @@ -23,7 +22,7 @@ from algorithm_exceptions import * logger = logging.getLogger("AnalyzerLog") -redis_conn = StrictRedis(unix_socket_path=REDIS_SOCKET_PATH) +redis_conn = utils.redis_conn() """ This is no man's land. Do anything you want in here, diff --git a/src/analyzer/analyzer.py b/src/analyzer/analyzer.py index bf7b74e3..b4d4aa72 100644 --- a/src/analyzer/analyzer.py +++ b/src/analyzer/analyzer.py @@ -15,7 +15,7 @@ from alerters import trigger_alert from algorithms import run_selected_algorithm from algorithm_exceptions import * -from utils import redis_conn +from utils import redis_conn, redis_conn_string logger = logging.getLogger("AnalyzerLog") @@ -136,7 +136,7 @@ def run(self): try: self.redis_conn.ping() except: - logger.error('skyline can\'t connect to redis at socket path %s' % settings.REDIS_SOCKET_PATH) + logger.error('skyline can\'t connect to redis at socket path %s' % redis_conn_string()) sleep(10) self.redis_conn = redis_conn() continue diff --git a/src/horizon/roomba.py b/src/horizon/roomba.py index c5c9840d..65101afa 100644 --- a/src/horizon/roomba.py +++ b/src/horizon/roomba.py @@ -5,7 +5,7 @@ from msgpack import Unpacker, packb from types import TupleType from time import time, sleep -from utils import redis_conn +from utils import redis_conn, redis_conn_string import logging import settings @@ -159,7 +159,7 @@ def run(self): try: self.redis_conn.ping() except: - logger.error('roomba can\'t connect to redis at socket path %s' % settings.REDIS_SOCKET_PATH) + logger.error('roomba can\'t connect to redis at socket path %s' % redis_conn_string()) sleep(10) self.redis_conn = redis_conn() continue diff --git a/src/horizon/worker.py b/src/horizon/worker.py index 8260fcec..6d257de8 100644 --- a/src/horizon/worker.py +++ b/src/horizon/worker.py @@ -4,7 +4,7 @@ from Queue import Empty from msgpack import packb from time import time, sleep -from utils import redis_conn +from utils import redis_conn, redis_conn_string import logging import socket @@ -75,7 +75,7 @@ def run(self): try: self.redis_conn.ping() except: - logger.error('worker can\'t connect to redis at socket path %s' % settings.REDIS_SOCKET_PATH) + logger.error('worker can\'t connect to redis at socket path %s' % redis_conn_string()) sleep(10) self.redis_conn = redis_conn() pipe = self.redis_conn.pipeline() diff --git a/src/utils.py b/src/utils.py index 425f912d..e2f839d4 100644 --- a/src/utils.py +++ b/src/utils.py @@ -2,6 +2,14 @@ import settings + +def redis_conn_string(): + if settings.REDIS_HOST_PORT is None: + return settings.REDIS_SOCKET_PATH + else: + return settings.REDIS_HOST_PORT + + def redis_conn(): if settings.REDIS_HOST_PORT is not None: (host, port) = settings.REDIS_HOST_PORT From 1881102c3b71fbf8a3549e5de13f3c473db0372b Mon Sep 17 00:00:00 2001 From: Zsolt Dollenstein Date: Fri, 8 Aug 2014 15:32:10 +0200 Subject: [PATCH 3/3] update example config --- src/settings.py.example | 2 ++ 1 file changed, 2 insertions(+) diff --git a/src/settings.py.example b/src/settings.py.example index 5a65a1f0..c5cb11e5 100644 --- a/src/settings.py.example +++ b/src/settings.py.example @@ -4,6 +4,8 @@ Shared settings # The path for the Redis unix socket REDIS_SOCKET_PATH = '/tmp/redis.sock' +# Or you can use a tcp socket +# REDIS_HOST_PORT = ('localhost', 6379) # The Skyline logs directory. Do not include a trailing slash. LOG_PATH = '/var/log/skyline'