Class: Fluent::HashForwardOutput
- Inherits:
-
ForwardOutput
- Object
- ForwardOutput
- Fluent::HashForwardOutput
- Defined in:
- lib/fluent/plugin/out_hash_forward.rb
Defined Under Namespace
Classes: NonHeartbeatNode
Instance Attribute Summary collapse
-
#hash_key_slice_lindex ⇒ Object
Returns the value of attribute hash_key_slice_lindex.
-
#hash_key_slice_rindex ⇒ Object
Returns the value of attribute hash_key_slice_rindex.
-
#regular_nodes ⇒ Object
readonly
for test.
-
#regular_weight_array ⇒ Object
readonly
Returns the value of attribute regular_weight_array.
-
#standby_nodes ⇒ Object
readonly
Returns the value of attribute standby_nodes.
-
#standby_weight_array ⇒ Object
readonly
Returns the value of attribute standby_weight_array.
-
#watcher_interval ⇒ Object
Returns the value of attribute watcher_interval.
Instance Method Summary collapse
-
#build_weight_array(nodes) ⇒ Object
This is just a partial copy from ForwardOuput#rebuild_weight_array.
- #cache_sock(node, sock) ⇒ Object
- #configure(conf) ⇒ Object
-
#get_index(key, size) ⇒ Object
hashing(key) mod N.
- #get_mutex(node) ⇒ Object
- #get_sock ⇒ Object
- #get_sock_expired_at ⇒ Object
-
#nodes(tag) ⇒ Object
Get nodes (a regular_node and a standby_node if available) using hash algorithm.
- #perform_hash_key_slice(tag) ⇒ Object
- #primary_available?(nodes) ⇒ Boolean
-
#rebuild_weight_array ⇒ Object
Override: I change weight algorithm.
- #reconnect(node) ⇒ Object
-
#run ⇒ Object
Override to disable heartbeat.
-
#send_data(node, tag, chunk) ⇒ Object
Override for keepalive.
- #shutdown ⇒ Object
- #sock_close(node) ⇒ Object
- #sock_write(sock, tag, chunk) ⇒ Object
- #start ⇒ Object
- #start_watcher ⇒ Object
- #stop_watcher ⇒ Object
-
#str_hash(key) ⇒ Object
the simplest hashing ever gist.github.com/sonots/7263495.
-
#watch_keepalive_time ⇒ Object
watcher thread callback.
-
#write_objects(tag, chunk) ⇒ Object
Override.
Instance Attribute Details
#hash_key_slice_lindex ⇒ Object
Returns the value of attribute hash_key_slice_lindex.
60 61 62 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 60 def hash_key_slice_lindex @hash_key_slice_lindex end |
#hash_key_slice_rindex ⇒ Object
Returns the value of attribute hash_key_slice_rindex.
61 62 63 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 61 def hash_key_slice_rindex @hash_key_slice_rindex end |
#regular_nodes ⇒ Object (readonly)
for test
56 57 58 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 56 def regular_nodes @regular_nodes end |
#regular_weight_array ⇒ Object (readonly)
Returns the value of attribute regular_weight_array.
58 59 60 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 58 def regular_weight_array @regular_weight_array end |
#standby_nodes ⇒ Object (readonly)
Returns the value of attribute standby_nodes.
57 58 59 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 57 def standby_nodes @standby_nodes end |
#standby_weight_array ⇒ Object (readonly)
Returns the value of attribute standby_weight_array.
59 60 61 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 59 def standby_weight_array @standby_weight_array end |
#watcher_interval ⇒ Object
Returns the value of attribute watcher_interval.
62 63 64 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 62 def watcher_interval @watcher_interval end |
Instance Method Details
#build_weight_array(nodes) ⇒ Object
This is just a partial copy from ForwardOuput#rebuild_weight_array
156 157 158 159 160 161 162 163 164 165 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 156 def build_weight_array(nodes) weight_array = [] gcd = nodes.map {|n| n.weight }.inject(0) {|r,w| r.gcd(w) } nodes.each {|n| (n.weight / gcd).times { weight_array << n } } weight_array end |
#cache_sock(node, sock) ⇒ Object
301 302 303 304 305 306 307 308 309 310 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 301 def cache_sock(node, sock) if sock get_sock[node] = sock get_sock_expired_at[node] = Time.now + @keepalive_time if @keepalive_time log.info "out_hash_forward: keepalive connection opened", :host=>node.host, :port=>node.port else get_sock[node] = nil get_sock_expired_at[node] = nil end end |
#configure(conf) ⇒ Object
28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 28 def configure(conf) super if @hash_key_slice lindex, rindex = @hash_key_slice.split('..', 2) if lindex.nil? or rindex.nil? or lindex !~ /^-?\d+$/ or rindex !~ /^-?\d+$/ raise Fluent::ConfigError, "out_hash_forard: hash_key_slice must be formatted like [num]..[num]" else @hash_key_slice_lindex = lindex.to_i @hash_key_slice_rindex = rindex.to_i end end if @heartbeat_type == :none @nodes = @nodes.map {|node| NonHeartbeatNode.new(node) } end @standby_nodes, @regular_nodes = @nodes.partition {|n| n.standby? } @regular_weight_array = build_weight_array(@regular_nodes) @standby_weight_array = build_weight_array(@standby_nodes) @cache_nodes = {} @sock = {} @sock_expired_at = {} @mutex = {} @watcher_interval = 1 end |
#get_index(key, size) ⇒ Object
hashing(key) mod N
184 185 186 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 184 def get_index(key, size) str_hash(key) % size end |
#get_mutex(node) ⇒ Object
295 296 297 298 299 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 295 def get_mutex(node) thread_id = Thread.current.object_id @mutex[thread_id] ||= {} @mutex[thread_id][node] ||= Mutex.new end |
#get_sock ⇒ Object
312 313 314 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 312 def get_sock @sock[Thread.current.object_id] ||= {} end |
#get_sock_expired_at ⇒ Object
316 317 318 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 316 def get_sock_expired_at @sock_expired_at[Thread.current.object_id] ||= {} end |
#nodes(tag) ⇒ Object
Get nodes (a regular_node and a standby_node if available) using hash algorithm
168 169 170 171 172 173 174 175 176 177 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 168 def nodes(tag) if nodes = @cache_nodes[tag] return nodes end hash_key = @hash_key_slice ? perform_hash_key_slice(tag) : tag regular_index = @regular_weight_array.size > 0 ? get_index(hash_key, @regular_weight_array.size) : 0 standby_index = @standby_weight_array.size > 0 ? get_index(hash_key, @standby_weight_array.size) : 0 nodes = [@regular_weight_array[regular_index], @standby_weight_array[standby_index]].compact @cache_nodes[tag] = nodes end |
#perform_hash_key_slice(tag) ⇒ Object
194 195 196 197 198 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 194 def perform_hash_key_slice(tag) = tag.split('.') sliced = [@hash_key_slice_lindex..@hash_key_slice_rindex] return sliced.nil? ? "" : sliced.join('.') end |
#primary_available?(nodes) ⇒ Boolean
179 180 181 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 179 def primary_available?(nodes) nodes.size > 1 && nodes.first.available? end |
#rebuild_weight_array ⇒ Object
Override: I change weight algorithm
152 153 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 152 def rebuild_weight_array end |
#reconnect(node) ⇒ Object
229 230 231 232 233 234 235 236 237 238 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 229 def reconnect(node) sock = connect(node) opt = [1, @send_timeout.to_i].pack('I!I!') # { int l_onoff; int l_linger; } sock.setsockopt(Socket::SOL_SOCKET, Socket::SO_LINGER, opt) opt = [@send_timeout.to_i, 0].pack('L!L!') # struct timeval sock.setsockopt(Socket::SOL_SOCKET, Socket::SO_SNDTIMEO, opt) sock end |
#run ⇒ Object
Override to disable heartbeat
92 93 94 95 96 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 92 def run unless @heartbeat_type == :none super end end |
#send_data(node, tag, chunk) ⇒ Object
Override for keepalive
201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 201 def send_data(node, tag, chunk) sock = nil get_mutex(node).synchronize do sock = get_sock[node] if @keepalive unless sock sock = reconnect(node) cache_sock(node, sock) if @keepalive end begin sock_write(sock, tag, chunk) node.heartbeat(false) log.debug "out_hash_forward: write to", :host=>node.host, :port=>node.port rescue Errno::EPIPE, Errno::ECONNRESET, Errno::ECONNABORTED, Errno::ETIMEDOUT => e log.warn "out_hash_forward: send_data failed #{e.class} #{e.}", :host=>node.host, :port=>node.port if @keepalive sock.close rescue IOError cache_sock(node, nil) end raise e ensure unless @keepalive sock.close if sock end end end end |
#shutdown ⇒ Object
69 70 71 72 73 74 75 76 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 69 def shutdown @finished = true @loop.watchers.each {|w| w.detach } @loop.stop unless @heartbeat_type == :none # custom @thread.join @usock.close if @usock stop_watcher end |
#sock_close(node) ⇒ Object
284 285 286 287 288 289 290 291 292 293 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 284 def sock_close(node) get_mutex(node).synchronize do if sock = get_sock[node] sock.close rescue IOError log.info "out_hash_forward: keepalive connection closed", :host=>node.host, :port=>node.port end get_sock[node] = nil get_sock_expired_at[node] = nil end end |
#sock_write(sock, tag, chunk) ⇒ Object
240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 240 def sock_write(sock, tag, chunk) # beginArray(2) sock.write FORWARD_HEADER # writeRaw(tag) sock.write tag.to_msgpack # tag # beginRaw(size) sz = chunk.size #if sz < 32 # # FixRaw # sock.write [0xa0 | sz].pack('C') #elsif sz < 65536 # # raw 16 # sock.write [0xda, sz].pack('Cn') #else # raw 32 sock.write [0xdb, sz].pack('CN') #end # writeRawBody(packed_es) chunk.write_to(sock) end |
#start ⇒ Object
64 65 66 67 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 64 def start super start_watcher end |
#start_watcher ⇒ Object
78 79 80 81 82 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 78 def start_watcher if @keepalive and @keepalive_time @watcher = Thread.new(&method(:watch_keepalive_time)) end end |
#stop_watcher ⇒ Object
84 85 86 87 88 89 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 84 def stop_watcher if @watcher @watcher.terminate @watcher.join end end |
#str_hash(key) ⇒ Object
the simplest hashing ever gist.github.com/sonots/7263495
190 191 192 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 190 def str_hash(key) key.bytes.inject(&:+) end |
#watch_keepalive_time ⇒ Object
watcher thread callback
265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 265 def watch_keepalive_time while true sleep @watcher_interval thread_ids = @sock.keys thread_ids.each do |thread_id| @sock[thread_id].each do |node, sock| @mutex[thread_id][node].synchronize do next unless sock_expired_at = @sock_expired_at[thread_id][node] next unless Time.now >= sock_expired_at sock.close rescue IOError if sock @sock[thread_id][node] = nil @sock_expired_at[thread_id][node] = nil log.debug "out_hash_forward: keepalive connection closed", :host=>node.host, :port=>node.port, :thread_id=>thread_id end end end end end |
#write_objects(tag, chunk) ⇒ Object
Override
123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 |
# File 'lib/fluent/plugin/out_hash_forward.rb', line 123 def write_objects(tag, chunk) return if chunk.empty? error = nil nodes = nodes(tag) if @keepalive and primary_available?(nodes) sock_close(nodes.last) # close standby end # below is just copy from out_forward nodes.each do |node| if node.available? begin send_data(node, tag, chunk) return rescue # for load balancing during detecting crashed servers error = $! # use the latest error end end end if error raise error else raise "no nodes are available" # TODO message end end |