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.
98 99 100 101 102 103 104 105 |
# File 'lib/redis_client/cluster/pipeline.rb', line 98 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
131 132 133 134 135 |
# File 'lib/redis_client/cluster/pipeline.rb', line 131 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
137 138 139 140 141 |
# File 'lib/redis_client/cluster/pipeline.rb', line 137 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
107 108 109 110 111 |
# File 'lib/redis_client/cluster/pipeline.rb', line 107 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
119 120 121 122 123 |
# File 'lib/redis_client/cluster/pipeline.rb', line 119 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
125 126 127 128 129 |
# File 'lib/redis_client/cluster/pipeline.rb', line 125 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
113 114 115 116 117 |
# File 'lib/redis_client/cluster/pipeline.rb', line 113 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
143 144 145 |
# File 'lib/redis_client/cluster/pipeline.rb', line 143 def empty? @size.zero? end |
#execute ⇒ Object
rubocop:disable Metrics/AbcSize, Metrics/CyclomaticComplexity, Metrics/PerceivedComplexity
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 175 176 177 178 179 180 181 182 183 184 |
# File 'lib/redis_client/cluster/pipeline.rb', line 147 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 = nil work_group.each do |node_key, v| case v when ::RedisClient::Cluster::Pipeline::RedirectionNeeded 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] } 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? all_replies end |