Class: Dalli::PipelinedGetter
- Inherits:
-
Object
- Object
- Dalli::PipelinedGetter
- 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
-
#initialize(ring, key_manager) ⇒ PipelinedGetter
constructor
A new instance of PipelinedGetter.
-
#process(keys, req_options = nil, &block) ⇒ Object
Yields, one at a time, keys and their values+attributes.
-
#process_with_metadata(keys, req_options = nil) ⇒ Object
Stale-aware bulk get across servers.
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, = nil, &block) return {} if keys.empty? @ring.lock do # Stores partial results collected during interleaved send phase @partial_results = {} servers = setup_requests(keys, ) 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, = 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, )) 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 |