Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,11 @@ Performance:
- Each connection now records when a quiet request is written, and the block's drain skips the others; servers that were never connected are no longer connected just to be drained
- With ~300us of network round trip, `multi { set }` on a 16-server ring goes from 4.9ms to 0.32ms; an empty block no longer touches the network
- Thanks to Julian Richard Contreras for this contribution
- Reduce Ruby overhead in `get_multi` reply parsing and key routing (#1169)
- `VA` reply headers are read in place instead of being split into tokens, on both the multi-server (pipelined) and single-server paths; results are unchanged, and unusual headers still take the token path
- A key's server is found through a bucket table over the ring's continuum instead of a binary search, and the key's own server is tried before the failover loop; the server chosen for every key is unchanged
- Allocations for a 100-key `get_multi` on 4 servers drop from 1,500 to 1,104; over loopback it goes from 301us to 243us
- Thanks to Julian Richard Contreras for this contribution

Bug fixes:

Expand Down
19 changes: 9 additions & 10 deletions lib/dalli/protocol/meta.rb
Original file line number Diff line number Diff line change
Expand Up @@ -335,26 +335,25 @@ def read_multi_metadata_responses(is_raw)

def read_multi_get_responses(is_raw)
hash = {}
key_index = is_raw ? 2 : 3
while (line = @connection_manager.read_line)
break if line.start_with?('MN')
next unless line.start_with?('VA ')

key, value = parse_multi_get_value(line, key_index, is_raw)
key, value = parse_multi_get_value(line, is_raw)
hash[key] = value if key
end
hash
end

def parse_multi_get_value(line, key_index, is_raw)
tokens = line.chomp!(TERMINATOR).split
value = @connection_manager.read(tokens[1].to_i + TERMINATOR.bytesize)&.chomp!(TERMINATOR)
raw_key = tokens[key_index]
return unless raw_key
# Reads the "VA <size> [f<flags>] k<key> ..." line in place rather than
# splitting it into tokens
def parse_multi_get_value(line, is_raw)
processor = response_processor
value = @connection_manager.read(processor.size_from_va_line(line) + TERMINATOR.bytesize)&.chomp!(TERMINATOR)
key = processor.key_from_va_line(line)
return unless key

key = raw_key[1..]
key = KeyRegularizer.decode(key) if tokens.include?('b')
bitflags = is_raw ? 0 : response_processor.bitflags_from_tokens(tokens)
bitflags = is_raw ? 0 : processor.bitflags_from_va_line(line)
[key, @value_marshaller.retrieve(value, bitflags)]
end

Expand Down
59 changes: 59 additions & 0 deletions lib/dalli/protocol/response_processor.rb
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ class ResponseProcessor
HD_PREFIX = 'HD '
FLAGS_TOKEN_PREFIX = ' f'
CAS_TOKEN_PREFIX = ' c'
KEY_TOKEN_PREFIX = ' k'
BYTE_B = 'b'.ord
BYTE_C = 'c'.ord
BYTE_F = 'f'.ord
Expand Down Expand Up @@ -231,6 +232,11 @@ def getk_response_from_buffer(buf, offset = 0)
term_idx = buf.byteindex(TERMINATOR, offset)
return [0] unless term_idx

if buf.byteslice(offset, VA_PREFIX.bytesize) == VA_PREFIX
response = va_response_from_buffer(buf, offset, term_idx)
return response if response
end

header = buf.byteslice(offset, term_idx - offset)
tokens = header.split
header_len = header.bytesize + TERMINATOR.length
Expand Down Expand Up @@ -272,6 +278,43 @@ def getk_response_from_buffer(buf, offset = 0)
[tokens.first == VA, flag_int(cas, 'c'), key, value, resp_size]
end

# getk_response_from_buffer for a "VA ..." header, read in place
# without allocating the header, its tokens or the flag strings.
# Gives the same result as the token path. Returns nil to fall back to
# it when the header has no s flag or a zero size.
def va_response_from_buffer(buf, offset, term_idx)
size = bitflags = cas = key = nil
base64 = false
pos = buf.byteindex(' ', offset + VA_PREFIX.bytesize)
while pos && pos < term_idx
start = pos + 1
pos = buf.byteindex(' ', start)
stop = pos && pos < term_idx ? pos : term_idx
case buf.getbyte(start)
when BYTE_S then size ||= flag_int_at(buf, start, stop)
when BYTE_F then bitflags ||= flag_int_at(buf, start, stop)
when BYTE_C then cas ||= flag_int_at(buf, start, stop)
when BYTE_K then key ||= buf.byteslice(start + 1, stop - start - 1)
when BYTE_B then base64 ||= stop == start + 1
end
end
return nil if size.nil? || size.zero?

header_len = term_idx - offset + TERMINATOR.length
resp_size = header_len + size + TERMINATOR.length
return [0] unless buf.bytesize >= offset + resp_size

value = @value_marshaller.retrieve(buf.byteslice(offset + header_len, size), bitflags || 0)
key = KeyRegularizer.decode(key) if base64 && key
[true, cas || 0, key || 0, value, resp_size]
end

# Integer value of a flag's token in buf, after its one-byte prefix,
# as flag_int would read it
def flag_int_at(buf, start, stop)
buf.byteslice(start + 1, stop - start - 1).to_i
end

# Integer value of a flag token such as "f123", or 0 when absent.
# Strips the prefix in place, like value_from_tokens.
def flag_int(token, flag)
Expand Down Expand Up @@ -306,6 +349,22 @@ def bitflags_from_va_line(line)
flag_from_line(line, FLAGS_TOKEN_PREFIX, VA_PREFIX.bytesize)
end

# The key from a VA line's k flag, decoded when the line has the b flag,
# or nil when there is no k flag
def key_from_va_line(line)
idx = line.index(KEY_TOKEN_PREFIX, VA_PREFIX.bytesize)
return unless idx

start = idx + KEY_TOKEN_PREFIX.bytesize
stop = line.index(' ', start) || line.index(TERMINATOR, start) || line.bytesize
key = line.byteslice(start, stop - start)
base64_flag_in_line?(line) ? KeyRegularizer.decode(key) : key
end

def base64_flag_in_line?(line)
line.include?(' b ') || line.end_with?(" b#{TERMINATOR}", ' b')
end

# Integer value of the first " <flag><digits>" token at or after start, or 0
def flag_from_line(line, token_prefix, start)
idx = line.index(token_prefix, start)
Expand Down
56 changes: 39 additions & 17 deletions lib/dalli/ring.rb
Original file line number Diff line number Diff line change
Expand Up @@ -36,9 +36,10 @@ def initialize(servers_arg, options)

def continuum=(entries)
@continuum = entries
# Plain integers, so the binary search in server_for_hash_key needn't
# call Entry#value at each step
# Plain integers, so the lookup in server_for_hash_key needn't call
# Entry#value at each step
@continuum_values = entries&.map(&:value)
build_buckets
end

# alive_cache (optional) remembers each server's alive? result, so a
Expand All @@ -55,15 +56,17 @@ def server_for_key(key, alive_cache = nil)
raise Dalli::RingError, 'No server available'
end

# Tries the key's own server first, then (with failover) up to 19
# rehashes of "<try><key>". The first try is outside the loop because it
# is almost always the one that's used.
def server_from_continuum(key, alive_cache = nil)
hkey = hash_for(key)
20.times do |try|
server = server_for_hash_key(hkey)
server = server_for_hash_key(hash_for(key))
return server if server_alive?(server, alive_cache)
return nil unless @failover

19.times do |try|
server = server_for_hash_key(hash_for("#{try}#{key}"))
return server if server_alive?(server, alive_cache)
break unless @failover

hkey = hash_for("#{try}#{key}")
end
nil
end
Expand Down Expand Up @@ -126,22 +129,41 @@ def hash_for(key)
def server_alive?(server, alive_cache)
return server.alive? unless alive_cache

alive_cache.fetch(server) { alive_cache[server] = server.alive? }
alive = alive_cache[server]
alive.nil? ? (alive_cache[server] = server.alive?) : alive
end

def entry_count_for(server, total_servers, total_weight)
((total_servers * POINTS_PER_SERVER * server.weight) / Float(total_weight)).floor
end

def server_for_hash_key(hash_key)
# Find the closest index in the Ring with value <= the given value
entryidx = @continuum_values.bsearch_index { |value| value > hash_key }
if entryidx.nil?
entryidx = @continuum.size - 1
else
entryidx -= 1
# Hash values are 32 bits. Splitting that range into about one bucket per
# continuum entry, and noting the first entry at or after each bucket's
# start, lets server_for_hash_key skip the binary search: it starts at the
# key's bucket and steps over the few entries before the key's hash.
def build_buckets
return unless @continuum_values

bits = [@continuum_values.size.bit_length, 1].max
@bucket_shift = 32 - bits
@bucket_first = Array.new(1 << bits)
idx = 0
size = @continuum_values.size
@bucket_first.each_index do |bucket|
floor = bucket << @bucket_shift
idx += 1 while idx < size && @continuum_values[idx] < floor
@bucket_first[bucket] = idx
end
@continuum[entryidx].server
end

def server_for_hash_key(hash_key)
# Find the last entry with value <= hash_key, wrapping around to the
# last entry when every value is greater (idx - 1 is then -1)
idx = @bucket_first[hash_key >> @bucket_shift]
values = @continuum_values
size = values.size
idx += 1 while idx < size && values[idx] <= hash_key
@continuum[idx - 1].server
end

def build_continuum(servers)
Expand Down
48 changes: 48 additions & 0 deletions test/protocol/test_response_processor.rb
Original file line number Diff line number Diff line change
Expand Up @@ -463,6 +463,54 @@ def expect_read_data(data, size)
assert_equal [0], processor.getk_response_from_buffer("VA 0 f0 kfoo s0\r\n".b)
end

it 'parses VA headers in place with the same results as the token path' do
tokens_only = Dalli::Protocol::Meta::ResponseProcessor.new(io_source, value_marshaller)
tokens_only.define_singleton_method(:va_response_from_buffer) { |*| nil }
value = Marshal.dump('hello')
b64 = ['k€y'].pack('m0')
headers = [
"VA #{value.bytesize} f1 kfoo s#{value.bytesize}\r\n#{value}\r\n",
"VA #{value.bytesize} f1 c42 kfoo s#{value.bytesize}\r\n#{value}\r\n",
"VA 5 f4 kbar s5\r\nhello\r\n",
"VA 5 s5 f0 kbar\r\nhello\r\n",
"VA 5 f0 k#{b64} b s5\r\nhello\r\n",
"VA 5 b f0 k#{b64} s5\r\nhello\r\n",
"VA 5 f0 kbar s5 f7 c1 c2 kother\r\nhello\r\n",
"VA 5 kbar s5\r\nhello\r\n",
"VA 5 f0 s5\r\nhello\r\n",
"VA 5 f0 kbar s5 bx W Z\r\nhello\r\n",
"VA 5 f4x kbar s5\r\nhello\r\n",
"VA 5 f0 kbar s5\r\nhel",
'VA 5 f0 kbar s5',
"VA 0 f0 kbar s0\r\n\r\n",
"VA 5 f0 kbar\r\nhello\r\n"
]

headers.each do |response|
# A non-zero offset, as when earlier responses are still in the buffer
buf = "HD\r\n#{response}MN\r\n".b

assert_equal tokens_only.getk_response_from_buffer(buf.dup, 4),
processor.getk_response_from_buffer(buf.dup, 4), response.inspect
end
# and the in-place parse is what handles a normal hit
hit = headers.first.b

assert_equal 'hello', processor.send(:va_response_from_buffer, hit, 0, hit.index("\r\n"))[3]
end

it 'reads the key from a VA line' do
b64 = ['k€y'].pack('m0')

assert_equal 'foo', processor.key_from_va_line("VA 5 f0 kfoo s5\r\n")
assert_equal 'foo', processor.key_from_va_line("VA 5 kfoo s5\r\n")
assert_equal 'foo', processor.key_from_va_line("VA 5 f0 kfoo\r\n")
assert_equal 'b', processor.key_from_va_line("VA 5 f0 kb s5\r\n")
assert_equal 'k€y', processor.key_from_va_line("VA 5 f0 k#{b64} b s5\r\n").force_encoding(Encoding::UTF_8)
assert_equal 'k€y', processor.key_from_va_line("VA 5 f0 k#{b64} s5 b\r\n").force_encoding(Encoding::UTF_8)
assert_nil processor.key_from_va_line("VA 5 f0 s5\r\n")
end

it 'returns [0] when the body has not fully arrived' do
buf = "VA 5 f0 c1 kfoo s5\r\nhel".b

Expand Down
35 changes: 35 additions & 0 deletions test/test_ring.rb
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,41 @@ def counting_ring(alive: true)
assert_predicate ring.server_for_key('test'), :alive?
end
end

# The server a binary search of the continuum picks: the last entry with
# value <= hash, or the last entry when every value is greater
def searched_server(ring, hash)
idx = ring.continuum.bsearch_index { |entry| entry.value > hash }
ring.continuum[(idx || ring.continuum.size) - 1].server
end

it 'maps hashes to the same server as a binary search of the continuum' do
[(1..2), (1..16), (1..50)].each do |ports|
ring = Dalli::Ring.new(ports.map { |p| "localhost:#{p}:#{(p % 3) + 1}" }, {})
values = ring.continuum.map(&:value)
hashes = values + values.map { |v| v - 1 } + values.map { |v| v + 1 } +
[0, (2**32) - 1] + Array.new(2000) { rand(2**32) }
hashes.select! { |h| h.between?(0, (2**32) - 1) }

hashes.each do |hash|
assert_same searched_server(ring, hash), ring.send(:server_for_hash_key, hash), "hash #{hash}"
end
end
end

it 'fails over to the server the rehashed key maps to' do
ring = Dalli::Ring.new(%w[localhost:12345 localhost:12346 localhost:12347], {})
dead = ring.servers.first
ring.servers.each { |s| s.define_singleton_method(:alive?) { !equal?(dead) } }

100.times do |i|
key = "key#{i}"
attempts = [key] + Array.new(19) { |try| "#{try}#{key}" }
expected = attempts.map { |k| searched_server(ring, Zlib.crc32(k)) }.find { |s| !s.equal?(dead) }

assert_same expected, ring.server_for_key(key)
end
end
end

it 'detect when a dead server is up again' do
Expand Down
Loading