Class: RedisClient::Cluster::Pipeline
- Inherits:
-
Object
- Object
- RedisClient::Cluster::Pipeline
- Defined in:
- lib/redis_client/cluster/pipeline.rb
Defined Under Namespace
Classes: Extended, RedirectionNeeded
Constant Summary collapse
- ReplySizeError =
Class.new(::RedisClient::Error)
Instance Method Summary collapse
- #blocking_call(timeout, *args, **kwargs, &block) ⇒ Object
- #blocking_call_v(timeout, args, &block) ⇒ Object
- #call(*args, **kwargs, &block) ⇒ Object
- #call_once(*args, **kwargs, &block) ⇒ Object
- #call_once_v(args, &block) ⇒ Object
- #call_v(args, &block) ⇒ Object
- #empty? ⇒ Boolean
-
#execute ⇒ Object
rubocop:disable Metrics/AbcSize, Metrics/CyclomaticComplexity, Metrics/PerceivedComplexity.
-
#initialize(router, command_builder, concurrent_worker, seed: Random.new_seed) ⇒ Pipeline
constructor
A new instance of Pipeline.
Constructor Details
#initialize(router, command_builder, concurrent_worker, seed: Random.new_seed) ⇒ Pipeline
Returns a new instance of Pipeline.
101 102 103 104 105 106 107 108 |
# File 'lib/redis_client/cluster/pipeline.rb', line 101 def initialize(router, command_builder, concurrent_worker, seed: Random.new_seed) @router = router @command_builder = command_builder @concurrent_worker = concurrent_worker @seed = seed @pipelines = nil @size = 0 end |
Instance Method Details
#blocking_call(timeout, *args, **kwargs, &block) ⇒ Object
134 135 136 137 138 |
# File 'lib/redis_client/cluster/pipeline.rb', line 134 def blocking_call(timeout, *args, **kwargs, &block) command = @command_builder.generate(args, kwargs) node_key = @router.find_node_key(command, seed: @seed) append_pipeline(node_key).blocking_call_v(timeout, command, &block) end |
#blocking_call_v(timeout, args, &block) ⇒ Object
140 141 142 143 144 |
# File 'lib/redis_client/cluster/pipeline.rb', line 140 def blocking_call_v(timeout, args, &block) command = @command_builder.generate(args) node_key = @router.find_node_key(command, seed: @seed) append_pipeline(node_key).blocking_call_v(timeout, command, &block) end |
#call(*args, **kwargs, &block) ⇒ Object
110 111 112 113 114 |
# File 'lib/redis_client/cluster/pipeline.rb', line 110 def call(*args, **kwargs, &block) command = @command_builder.generate(args, kwargs) node_key = @router.find_node_key(command, seed: @seed) append_pipeline(node_key).call_v(command, &block) end |
#call_once(*args, **kwargs, &block) ⇒ Object
122 123 124 125 126 |
# File 'lib/redis_client/cluster/pipeline.rb', line 122 def call_once(*args, **kwargs, &block) command = @command_builder.generate(args, kwargs) node_key = @router.find_node_key(command, seed: @seed) append_pipeline(node_key).call_once_v(command, &block) end |
#call_once_v(args, &block) ⇒ Object
128 129 130 131 132 |
# File 'lib/redis_client/cluster/pipeline.rb', line 128 def call_once_v(args, &block) command = @command_builder.generate(args) node_key = @router.find_node_key(command, seed: @seed) append_pipeline(node_key).call_once_v(command, &block) end |
#call_v(args, &block) ⇒ Object
116 117 118 119 120 |
# File 'lib/redis_client/cluster/pipeline.rb', line 116 def call_v(args, &block) command = @command_builder.generate(args) node_key = @router.find_node_key(command, seed: @seed) append_pipeline(node_key).call_v(command, &block) end |
#empty? ⇒ Boolean
146 147 148 |
# File 'lib/redis_client/cluster/pipeline.rb', line 146 def empty? @size.zero? end |
#execute ⇒ Object
rubocop:disable Metrics/AbcSize, Metrics/CyclomaticComplexity, Metrics/PerceivedComplexity
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 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 |
# File 'lib/redis_client/cluster/pipeline.rb', line 150 def execute # rubocop:disable Metrics/AbcSize, Metrics/CyclomaticComplexity, Metrics/PerceivedComplexity return if @pipelines.nil? || @pipelines.empty? work_group = @concurrent_worker.new_group(size: @pipelines.size) @pipelines.each do |node_key, pipeline| work_group.push(node_key, @router.find_node(node_key), pipeline) do |cli, pl| replies = do_pipelining(cli, pl) raise ReplySizeError, "commands: #{pl._size}, replies: #{replies.size}" if pl._size != replies.size replies end end all_replies = errors = required_redirections = nil work_group.each do |node_key, v| case v when ::RedisClient::Cluster::Pipeline::RedirectionNeeded required_redirections ||= {} required_redirections[node_key] = v when StandardError errors ||= {} errors[node_key] = v else all_replies ||= Array.new(@size) @pipelines[node_key].outer_indices.each_with_index { |outer, inner| all_replies[outer] = v[inner] } end end work_group.close raise ::RedisClient::Cluster::ErrorCollection, errors unless errors.nil? required_redirections&.each do |node_key, v| all_replies ||= Array.new(@size) pipeline = @pipelines[node_key] v.indices.each { |i| v.replies[i] = handle_redirection(v.replies[i], pipeline, i) } pipeline.outer_indices.each_with_index { |outer, inner| all_replies[outer] = v.replies[inner] } end all_replies end |