diff --git a/.circleci/config.yml b/.circleci/config.yml new file mode 100644 index 0000000..674b101 --- /dev/null +++ b/.circleci/config.yml @@ -0,0 +1,41 @@ +# Javascript Node CircleCI 2.0 configuration file +# +# Check https://circleci.com/docs/2.0/language-javascript/ for more details +# +version: 2 +jobs: + build: + docker: + # specify the version you desire here + - image: circleci/node:7.10 + + # Specify service dependencies here if necessary + # CircleCI maintains a library of pre-built images + # documented at https://circleci.com/docs/2.0/circleci-images/ + # - image: circleci/mongo:3.4.4 + + working_directory: ~/repo + + steps: + - checkout + + # Download and cache dependencies + - restore_cache: + keys: + - v1-dependencies-{{ checksum "package.json" }} + # fallback to using the latest cache if no exact match is found + - v1-dependencies- + + - run: npm install + + - run: sudo apt-get update + + - run: sudo apt-get install redis-server + + - save_cache: + paths: + - node_modules + key: v1-dependencies-{{ checksum "package.json" }} + + # run tests! + - run: npm run tests \ No newline at end of file diff --git a/.eslintrc.json b/.eslintrc.json new file mode 100644 index 0000000..bd0042c --- /dev/null +++ b/.eslintrc.json @@ -0,0 +1,12 @@ +{ + "extends": "airbnb", + "rules": { + "comma-dangle": ["error", { + "arrays": "never", + "objects": "never", + "imports": "never", + "exports": "never", + "functions": "never" + }] + } +} diff --git a/README.md b/README.md index 6cfe1cc..00b6438 100644 --- a/README.md +++ b/README.md @@ -1,11 +1,15 @@ -redis-shard -=========== +redis-shard-optimized +===================== -A consistent sharding library for redis in node +[![CircleCI](https://circleci.com/gh/ppsreejith/node-redis-shard.svg?style=svg)](https://circleci.com/gh/ppsreejith/node-redis-shard) - $ npm install redis-shard +A consistent sharding library for redis in node. - var RedisShard = require('redis-shard'); +This project improves over the parent project by consistently distributing mget and mset over the sharded redis instances. (The parent implementation sends mget and mset to a single instance). + + $ npm install redis-shard-optimized + + var RedisShard = require('redis-shard-optimized'); var options = { servers: [ '127.0.0.1:6379', '127.0.0.1:6479' ], database : 1, password : 'redis4pulseLocker' }; var redis = new RedisShard(options); @@ -13,6 +17,10 @@ A consistent sharding library for redis in node redis.set('foo', 'bar', console.log); redis.get('foo', console.log); + // SINGLE (Multi key commands) + redis.mset(['key1', 'val1', 'key2', 'val2', 'key3', 'val3'], console.log); + redis.mget(['key1', 'key2', 'key3'], console.log); + // MULTI var multi = redis.multi(); multi.set('foo', 'bar').set('bah', 'baz').expire('foo', 3600).expire('bah', 3600); @@ -26,4 +34,13 @@ The constructor accepts an object containing the following options: - `servers` (required) - An array of Redis servers (e.g. `'127.0.0.1:6379'`) to connect to - `database` - Redis database to select - `password` - Password for authentication -- `clientOptions` - Options object to be passed to each Redis client \ No newline at end of file +- `clientOptions` - Options object to be passed to each Redis client + +Tests +----- + +To run tests, use + +``` +npm run tests +``` \ No newline at end of file diff --git a/index.js b/index.js index 03010f6..f2b7652 100644 --- a/index.js +++ b/index.js @@ -1,113 +1,191 @@ -var assert = require('assert'); -var HashRing = require('hashring'); -var redis = require('redis'); -var step = require('step'); +const assert = require('assert'); +const HashRing = require('hashring'); +const redis = require('redis'); +const step = require('step'); +const _ = require('lodash'); +const async = require('async'); module.exports = function RedisShard(options) { - assert(!!options, "options must be an object"); - assert(Array.isArray(options.servers), "servers must be an array"); - - var self = {}; - var clients = {}; - options.servers.forEach(function(server) { - var fields = server.split(/:/); - var clientOptions = options.clientOptions || {}; - var client = redis.createClient(parseInt(fields[1], 10), fields[0], clientOptions); - if ( options.database ) { - client.select(options.database, function(){}); + assert(!!options, 'options must be an object'); + assert(Array.isArray(options.servers), 'servers must be an array'); + + const self = {}; + const clients = {}; + options.servers.forEach((server) => { + const fields = server.split(/:/); + const clientOptions = options.clientOptions || {}; + const client = redis.createClient(parseInt(fields[1], 10), fields[0], clientOptions); + if (options.database) { + client.select(options.database, () => {}); } - if ( options.password ) { + if (options.password) { client.auth(options.password); } clients[server] = client; }); - var servers = {}; - for (var key in clients) { + const servers = {}; + for (const key in clients) { servers[key] = 1; // balanced ring for now } self.ring = new HashRing(servers); // All of these commands have 'key' as their first parameter - var SHARDABLE = [ - "append", "bitcount", "blpop", "brpop", "debug object", "decr", "decrby", "del", "dump", "exists", "expire", - "expireat", "get", "getbit", "getrange", "getset", "hdel", "hexists", "hget", "hgetall", "hincrby", - "hincrbyfloat", "hkeys", "hlen", "hmget", "hmset", "hset", "hsetnx", "hvals", "incr", "incrby", "incrbyfloat", - "lindex", "linsert", "llen", "lpop", "lpush", "lpushx", "lrange", "lrem", "lset", "ltrim", "mget", "move", - "persist", "pexpire", "pexpireat", "psetex", "pttl", "rename", "renamenx", "restore", "rpop", "rpush", "rpushx", - "sadd", "scard", "sdiff", "set", "setbit", "setex", "setnx", "setrange", "sinter", "sismember", "smembers", - "sort", "spop", "srandmember", "srem", "strlen", "sunion", "ttl", "type", "watch", "zadd", "zcard", "zcount", - "zincrby", "zrange", "zrangebyscore", "zrank", "zrem", "zremrangebyrank", "zremrangebyscore", "zrevrange", - "zrevrangebyscore", "zrevrank", "zscore" + const SHARDABLE = [ + 'append', 'bitcount', 'blpop', 'brpop', 'debug object', 'decr', 'decrby', 'del', 'dump', 'exists', 'expire', + 'expireat', 'get', 'getbit', 'getrange', 'getset', 'hdel', 'hexists', 'hget', 'hgetall', 'hincrby', + 'hincrbyfloat', 'hkeys', 'hlen', 'hmget', 'hmset', 'hset', 'hsetnx', 'hvals', 'incr', 'incrby', 'incrbyfloat', + 'lindex', 'linsert', 'llen', 'lpop', 'lpush', 'lpushx', 'lrange', 'lrem', 'lset', 'ltrim', 'move', + 'persist', 'pexpire', 'pexpireat', 'psetex', 'pttl', 'rename', 'renamenx', 'restore', 'rpop', 'rpush', 'rpushx', + 'sadd', 'scard', 'sdiff', 'set', 'setbit', 'setex', 'setnx', 'setrange', 'sinter', 'sismember', 'smembers', + 'sort', 'spop', 'srandmember', 'srem', 'strlen', 'sunion', 'ttl', 'type', 'watch', 'zadd', 'zcard', 'zcount', + 'zincrby', 'zrange', 'zrangebyscore', 'zrank', 'zrem', 'zremrangebyrank', 'zremrangebyscore', 'zrevrange', + 'zrevrangebyscore', 'zrevrank', 'zscore' ]; - SHARDABLE.forEach(function(command) { - self[command] = function() { - var node = self.ring.get(arguments[0]); - var client = clients[node]; - client[command].apply(client, arguments); + SHARDABLE.forEach((command) => { + self[command] = function () { + const node = self.ring.get(arguments[0]); + const client = clients[node]; + client[command](...arguments); }; }); + // mget + self.mget = function () { + const keys = _.first(arguments); + const callback = _.last(arguments); + const mapping = _.reduce(keys, (acc, key) => { + const node = self.ring.get(key); + if (_.has(acc, node)) { + acc[node].push(key); + } else { + acc[node] = [key]; + } + return acc; + }, {}); + const nodes = _.keys(mapping); + async.map(nodes, (node, next) => { + const keys = mapping[node]; + const client = clients[node]; + // console.log(node, keys.length); + client.mget(keys, next); + }, (err, results) => { + if (err) { + return callback(err); + } + const keyHash = _.reduce(nodes, (acc, node, index) => { + const keys = mapping[node]; + const values = _.get(results, [index], []); + return _.assign(acc, _.zipObject(keys, values)); + }, {}); + return callback(err, _.map(keys, key => keyHash[key])); + }); + }; + + // keys + self.keys = function () { + const key_regex = _.first(arguments); + const callback = _.last(arguments); + const nodes = _.keys(self.ring.vnodes); + async.map(nodes, (node, next) => { + const client = clients[node]; + client.keys(key_regex, next); + }, (err, results) => { + if (err) { + return callback(err); + } + let keys = []; + _.map(results, node_keys => keys = _.concat(keys, node_keys)); + return callback(err, keys); + }); + }; + + // mset + self.mset = function () { + const keys = _.first(arguments); + const callback = _.last(arguments); + const mapping = _.reduce(keys, (acc, key, index) => { + const keyCheck = index % 2 ? keys[index - 1] : key; + const node = self.ring.get(keyCheck); + if (_.has(acc, node)) { + acc[node].push(key); + } else { + acc[node] = [key]; + } + return acc; + }, {}); + const nodes = _.keys(mapping); + async.map(nodes, (node, next) => { + const keys = mapping[node]; + const client = clients[node]; + client.mset(keys, next); + }, (err, responses) => { + if (err) { + return callback(err); + } + return callback(null, 'OK'); + }); + }; + // No key parameter to shard on - throw Error - var UNSHARDABLE = [ - "auth", "bgrewriteaof", "bgsave", "bitop", "brpoplpush", "client kill", "client list", "client getname", - "client setname", "config get", "config set", "config resetstat", "dbsize", "debug segfault", "discard", - "echo", "eval", "evalsha", "exec", "flushall", "flushdb", "info", "keys", "lastsave", "migrate", "monitor", - "mset", "msetnx", "multi", "object", "ping", "psubscribe", "publish", "punsubscribe", "quit", "randomkey", - "rpoplpush", "save", "script exists", "script flush", "script kill", "script load", "sdiffstore", "select", - "shutdown", "sinterstore", "slaveof", "slowlog", "smove", "subscribe", "sunionstore", "sync", "time", - "unsubscribe", "unwatch", "zinterstore", "zunionstore" + const UNSHARDABLE = [ + 'auth', 'bgrewriteaof', 'bgsave', 'bitop', 'brpoplpush', 'client kill', 'client list', 'client getname', + 'client setname', 'config get', 'config set', 'config resetstat', 'dbsize', 'debug segfault', 'discard', + 'echo', 'eval', 'evalsha', 'exec', 'flushall', 'flushdb', 'info', 'lastsave', 'migrate', 'monitor', + 'msetnx', 'multi', 'object', 'ping', 'psubscribe', 'publish', 'punsubscribe', 'quit', 'randomkey', + 'rpoplpush', 'save', 'script exists', 'script flush', 'script kill', 'script load', 'sdiffstore', 'select', + 'shutdown', 'sinterstore', 'slaveof', 'slowlog', 'smove', 'subscribe', 'sunionstore', 'sync', 'time', + 'unsubscribe', 'unwatch', 'zinterstore', 'zunionstore' ]; - UNSHARDABLE.forEach(function(command) { - self[command] = function() { - throw new Error(command + ' is not shardable'); + UNSHARDABLE.forEach((command) => { + self[command] = function () { + throw new Error(`${command} is not shardable`); }; }); // This is the tricky part - pipeline commands to multiple servers self.multi = function Multi() { - - var self = {}; - var multis = {}; - var interlachen = []; + const self = {}; + const multis = {}; + const interlachen = []; // Setup chainable shardable commands - SHARDABLE.forEach(function(command) { - self[command] = function() { - var node = self.ring.get(arguments[0]); - var multi = multis[node]; + SHARDABLE.forEach((command) => { + self[command] = function () { + const node = self.ring.get(arguments[0]); + let multi = multis[node]; if (!multi) { multi = multis[node] = clients[node].multi(); } interlachen.push(node); - multi[command].apply(multi, arguments); + multi[command](...arguments); return self; }; }); - UNSHARDABLE.forEach(function(command) { - self[command] = function() { - throw new Error(command + " is not supported"); + UNSHARDABLE.forEach((command) => { + self[command] = function () { + throw new Error(`${command} is not supported`); }; }); // Exec the pipeline and interleave the results - self.exec = function(callback) { - var nodes = Object.keys(multis); + self.exec = function (callback) { + const nodes = Object.keys(multis); step( function run() { - var group = this.group(); - nodes.forEach(function(node) { + const group = this.group(); + nodes.forEach((node) => { multis[node].exec(group()); }); }, - function done(error, groups) { + (error, groups) => { if (error) { return callback(error); } - assert(nodes.length === groups.length, "wrong number of responses"); - var results = []; - interlachen.forEach(function(node) { - var index = nodes.indexOf(node); - assert(groups[index].length > 0, node + " is missing a result"); + assert(nodes.length === groups.length, 'wrong number of responses'); + const results = []; + interlachen.forEach((node) => { + const index = nodes.indexOf(node); + assert(groups[index].length > 0, `${node} is missing a result`); results.push(groups[index].shift()); }); callback(null, results); @@ -118,26 +196,26 @@ module.exports = function RedisShard(options) { }; - self.on = function(event, listener) { - options.servers.forEach(function(server) { - clients[server].on(event, function() { + self.on = function (event, listener) { + options.servers.forEach((server) => { + clients[server].on(event, function () { // append server as last arg passed to listener - var args = Array.prototype.slice.call(arguments).concat(server); - listener.apply(undefined, args); + const args = Array.prototype.slice.call(arguments).concat(server); + listener(...args); }); }); }; // Note: listener will fire once per shard, not once per cluster - self.once = function(event, listener) { - options.servers.forEach(function(server) { - clients[server].once(event, function() { + self.once = function (event, listener) { + options.servers.forEach((server) => { + clients[server].once(event, function () { // append server as last arg passed to listener - var args = Array.prototype.slice.call(arguments).concat(server); - listener.apply(undefined, args); + const args = Array.prototype.slice.call(arguments).concat(server); + listener(...args); }); }); - } + }; return self; // RedisShard() }; diff --git a/package.json b/package.json index fb127ef..b73769b 100644 --- a/package.json +++ b/package.json @@ -1,14 +1,17 @@ { - "name": "redis-shard", - "version": "0.3.0", + "name": "redis-shard-optimized", + "version": "0.3.4", "description": "A consistent hashing library for redis in node", "main": "index.js", "scripts": { - "test": "echo \"Error: no test specified\" && exit 1" + "tests": "node test.js" }, "dependencies": { + "async": "^1.5.2", "hashring": "^3.2.0", + "lodash": "^4.17.5", "redis": "^2.1.0", + "redis-server": "^1.2.0", "step": "0.0.6" }, "engines": { @@ -16,12 +19,21 @@ }, "repository": { "type": "git", - "url": "git://github.com/blindsey/node-redis-shard.git" + "url": "git://github.com/ppsreejith/node-redis-shard.git" }, "author": "Ben Lindsey ", "contributors": [ + "Sreejith Pp ", "Toby Fox " ], "license": "MIT", - "readmeFilename": "README.md" + "readmeFilename": "README.md", + "devDependencies": { + "babel-eslint": "^8.2.1", + "eslint": "^4.18.0", + "eslint-config-airbnb": "^16.1.0", + "eslint-plugin-import": "^2.8.0", + "eslint-plugin-jsx-a11y": "^6.0.3", + "eslint-plugin-react": "^7.7.0" + } } diff --git a/test.js b/test.js new file mode 100644 index 0000000..8cd89c7 --- /dev/null +++ b/test.js @@ -0,0 +1,110 @@ +const RedisServer = require('redis-server'); +const async = require('async'); +const _ = require('lodash'); + +const RedisShard = require('./index'); + +const PORT = 7000; +const NUM = 5; +const TESTS = {}; + +TESTS.keysCommandTests = ({KEY_NUM, redis}, callback) => { + const keys = _.times(KEY_NUM, n => `fooregexx${n}`); + const values = _.times(KEY_NUM, n => `barregexx${n}`); + + const keys2 = _.times(KEY_NUM, n => `fooregex${n}`); + const values2 = _.times(KEY_NUM, n => `barregex${n}`); + + const allKeys = _.concat(keys, keys2); + const allValues = _.concat(values, values2); + + const msetArgs = _ + .chain(allKeys) + .zip(allValues) + .flatten() + .value(); + + const regex1 = 'fooregexx*'; // meant to match keys + const regex2 = 'fooregex*'; // meant to match allKeys + + async.auto({ + mset: callback => redis.mset(msetArgs, callback), + keys1: ['mset', callback => redis.keys(regex1, callback)], + keys2: ['mset', callback => redis.keys(regex2, callback)], + }, (err, { mset, keys1, keys2 }) => { + if (err) { + return callback(err); + } + if (mset !== 'OK' ) { + return callback('mset was not successful'); + } + + const sortedKeys1 = _.sortBy(keys); + const sortedKeys2 = _.sortBy(allKeys); + + const retrievedKeys1 = _.sortBy(keys1); + const retrievedKeys2 = _.sortBy(keys2); + + if (!_.isEqual(sortedKeys1, retrievedKeys1) || !_.isEqual(sortedKeys2, retrievedKeys2)) { + return callback('Keys response doesn\'t match values'); + } + return callback(); + }); +}; + +TESTS.multiSetGetCommandTests = ({KEY_NUM, redis}, callback) => { + const keys = _.times(KEY_NUM, n => `foo${n}`); + const values = _.times(KEY_NUM, n => `bar${n}`); + + const msetArgs = _ + .chain(keys) + .zip(values) + .flatten() + .value(); + + async.auto({ + mset: callback => redis.mset(msetArgs, callback), + mget: ['mset', callback => redis.mget(keys, callback)] + }, (err, { mset, mget }) => { + if (err) { + return callback(err); + } + if (mset !== 'OK') { + return callback('mset was not successful'); + } + if (!_.isEqual(mget, values)) { + return callback('Mget response doesn\'t match values'); + } + return callback(); + }); +} + +async.times(NUM, (index, next) => { + const server = new RedisServer(PORT + index); + server.open(next); +}, (err) => { + const options = { + servers: _.map(_.times(NUM, n => `127.0.0.1:${PORT + n}`)) + }; + const redis = new RedisShard(options); + const KEY_NUM = 1000; + + const args = { KEY_NUM, redis }; + + const asyncArgs = _.reduce( + TESTS, + (acc, value, key) => { + acc[key] = _.partial(value, args); + return acc; + }, + {} + ); + + async.auto(asyncArgs, (err) => { + if (err) { + throw err; + } + console.log('Success. All tests passed.'); + process.exit(0); + }); +}); \ No newline at end of file