Class: Dalli::PipelinedGetter

Inherits:
Object
  • Object
show all
Defined in:
lib/dalli/pipelined_getter.rb

Overview

Contains logic for the pipelined gets implemented by the client.

Constant Summary collapse

INTERLEAVE_THRESHOLD =

For large batches, interleave sends with response draining to prevent socket buffer deadlock. Only kicks in above this threshold.

10_000
CHUNK_SIZE =

Number of keys to send before draining responses during interleaved mode

10_000

Instance Method Summary collapse

Constructor Details

#initialize(ring, key_manager) ⇒ PipelinedGetter

Returns a new instance of PipelinedGetter.



15
16
17
18
# File 'lib/dalli/pipelined_getter.rb', line 15

def initialize(ring, key_manager)
  @ring = ring
  @key_manager = key_manager
end

Instance Method Details

#process(keys, req_options = nil, &block) ⇒ Object

Yields, one at a time, keys and their values+attributes.

req_options accepts :p_token/:l_token, applied to every key in the batch.

A transient network error is retried automatically. If a server remains unreachable after retrying, raises Dalli::NetworkError.



28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
# File 'lib/dalli/pipelined_getter.rb', line 28

def process(keys, req_options = nil, &block)
  return {} if keys.empty?

  @ring.lock do
    # Stores partial results collected during interleaved send phase
    @partial_results = {}
    servers = setup_requests(keys, req_options)
    start_time = Process.clock_gettime(Process::CLOCK_MONOTONIC)

    # First yield any partial results collected during interleaved send
    yield_partial_results(&block)

    servers = fetch_responses(servers, start_time, @ring.socket_timeout, &block) until servers.empty?
  end
rescue Dalli::RetryableNetworkError => e
  Dalli.logger.debug { e.inspect }
  Dalli.logger.debug { 'retrying pipelined gets because of timeout' }
  retry
end

#process_with_metadata(keys, req_options = nil) ⇒ Object

Stale-aware bulk get across servers. Returns { key => metadata Hash } for the keys that were found; see Protocol::Meta#read_multi_with_metadata_req for why misses are absent rather than present with miss: true.

Unlike #process this issues one request per server and reads its full response before moving on, rather than pipelining across servers: the metadata path has no interleaving support, and the stale-aware callers it serves fetch far smaller batches than get_multi does.

req_options accepts :p_token/:l_token, applied to every key in the batch.



60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
# File 'lib/dalli/pipelined_getter.rb', line 60

def (keys, req_options = nil)
  return {} if keys.empty?

  @ring.lock do
    results = {}
    groups_for_keys(keys).each do |server, keys_for_server|
      results.merge!(server.request(:read_multi_with_metadata_req, keys_for_server, req_options))
    rescue Dalli::RetryableNetworkError
      raise
    rescue DalliError, NetworkError => e
      Dalli.logger.debug { e.inspect }
      Dalli.logger.debug { "unable to get keys for server #{server.name}" }
    end
    results.transform_keys! { |key| @key_manager.key_without_namespace(key) }
    results
  end
rescue Dalli::RetryableNetworkError => e
  Dalli.logger.debug { e.inspect }
  Dalli.logger.debug { 'retrying pipelined get with metadata because of network error' }
  retry
end