This repository has been archived by the owner on Sep 18, 2021. It is now read-only.
forked from freels/kestrel-client
/
client.rb
174 lines (147 loc) · 5.42 KB
/
client.rb
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
require 'forwardable'
module Kestrel
class Client
require 'kestrel/client/stats_helper'
require 'kestrel/client/retry_helper'
autoload :Proxy, 'kestrel/client/proxy'
autoload :Envelope, 'kestrel/client/envelope'
autoload :Blocking, 'kestrel/client/blocking'
autoload :Partitioning, 'kestrel/client/partitioning'
autoload :Unmarshal, 'kestrel/client/unmarshal'
autoload :Namespace, 'kestrel/client/namespace'
autoload :Json, 'kestrel/client/json'
autoload :Transactional, "kestrel/client/transactional"
KESTREL_OPTIONS = [:gets_per_server, :exception_retry_limit, :get_timeout_ms].freeze
DEFAULT_OPTIONS = {
:retry_timeout => 0,
:exception_retry_limit => 5,
:timeout => 0.25,
:gets_per_server => 100,
:get_timeout_ms => 10
}.freeze
# Exceptions which are connection failures we retry after
RECOVERABLE_ERRORS = [
Memcached::ServerIsMarkedDead,
Memcached::ATimeoutOccurred,
Memcached::ConnectionBindFailure,
Memcached::ConnectionFailure,
Memcached::ConnectionSocketCreateFailure,
Memcached::Failure,
Memcached::MemoryAllocationFailure,
Memcached::ReadFailure,
Memcached::ServerError,
Memcached::SystemError,
Memcached::UnknownReadFailure,
Memcached::WriteFailure,
Memcached::NotFound
]
extend Forwardable
include StatsHelper
include RetryHelper
attr_accessor :servers, :options
attr_reader :current_queue, :kestrel_options, :current_server
def_delegators :@write_client, :add, :append, :cas, :decr, :incr, :get_orig, :prepend
def initialize(*servers)
opts = servers.last.is_a?(Hash) ? servers.pop : {}
opts = DEFAULT_OPTIONS.merge(opts)
@kestrel_options = extract_kestrel_options!(opts)
@default_get_timeout = kestrel_options[:get_timeout_ms]
@gets_per_server = kestrel_options[:gets_per_server]
@exception_retry_limit = kestrel_options[:exception_retry_limit]
@counter = 0
# we handle our own retries so that we can apply different
# policies to sets and gets, so set memcached limit to 0
opts[:exception_retry_limit] = 0
opts[:distribution] = :random # force random distribution
self.servers = Array(servers).flatten.compact
self.options = opts
@server_count = self.servers.size # Minor optimization.
@read_client = Memcached.new(self.servers[rand(@server_count)], opts)
@write_client = Memcached.new(self.servers, opts)
end
def delete(key, expiry=0)
with_retries { @write_client.delete key }
rescue Memcached::NotFound, Memcached::ServerEnd
end
def set(key, value, ttl=0, raw=false)
with_retries { @write_client.set key, value, ttl, !raw }
true
rescue Memcached::NotStored
false
end
# ==== Parameters
# key<String>:: Queue name
# opts<Boolean,Hash>:: True/false toggles Marshalling. A Hash
# allows collision-avoiding options support.
#
# ==== Options (opts)
# :open<Boolean>:: Begins a transactional read.
# :close<Boolean>:: Ends a transactional read.
# :abort<Boolean>:: Cancels an existing transactional read
# :peek<Boolean>:: Return the head of the queue, without removal
# :timeout<Integer>:: Milliseconds to block for a new item
# :raw<Boolean>:: Toggles Marshalling. Equivalent to the "old
# style" second argument.
#
def get(key, opts = {})
raw = opts.delete(:raw) || false
commands = extract_queue_commands(opts)
val =
begin
shuffle_if_necessary! key
# NOP rehashing for a single server on #get seems broken in
# libmemcached. Stick to get_from_last, for now.
# FIXME: When `rake benchmark` on #get matches #get_from_last, switch.
@read_client.get_from_last key + commands, !raw
rescue *RECOVERABLE_ERRORS
# we can't tell the difference between a server being down
# and an empty queue, so just return nil. our sticky server
# logic should eliminate piling on down servers
nil
end
# nil result, force next get to jump from current server
@counter = @gets_per_server unless val
val
end
def flush(queue)
count = 0
while sizeof(queue) > 0
count += 1 while get queue, :raw => true
end
count
end
def peek(queue)
get queue, :peek => true
end
private
def extract_kestrel_options!(opts)
kestrel_opts, memcache_opts = opts.inject([{}, {}]) do |(kestrel, memcache), (key, opt)|
(KESTREL_OPTIONS.include?(key) ? kestrel : memcache)[key] = opt
[kestrel, memcache]
end
opts.replace(memcache_opts)
kestrel_opts
end
def shuffle_if_necessary!(key)
# Don't reset servers on the first request:
# i.e. @counter == 0 && @current_queue == nil
if (@counter > 0 && key != @current_queue) || @counter >= @gets_per_server
@counter = 0
@current_queue = key
@read_client.quit
@read_client.set_servers(servers[rand(@server_count)])
else
@counter +=1
end
end
def extract_queue_commands(opts)
commands = [:open, :close, :abort, :peek].select do |key|
opts[key]
end
if timeout = (opts[:timeout] || @default_get_timeout)
commands << "t=#{timeout}"
end
commands.map { |c| "/#{c}" }.join('')
end
end
end