Class: Fluent::KeepForwardOutput
- Inherits:
-
ForwardOutput
- Object
- ForwardOutput
- Fluent::KeepForwardOutput
- Defined in:
- lib/fluent/plugin/out_keep_forward.rb
Instance Attribute Summary collapse
-
#watcher_interval ⇒ Object
for test.
Instance Method Summary collapse
- #cache_node(tag, node) ⇒ Object
- #cache_sock(node, sock) ⇒ Object
- #configure(conf) ⇒ Object
- #get_mutex(node) ⇒ Object
- #get_node(tag) ⇒ Object
- #get_sock ⇒ Object
- #get_sock_expired_at ⇒ Object
- #keepforward(tag) ⇒ Object
- #reconnect(node) ⇒ Object
-
#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
-
#watch_keepalive_time ⇒ Object
watcher thread callback.
- #weight_send_data(tag, chunk, error_node = nil) ⇒ Object
-
#write_objects(tag, chunk) ⇒ Object
Override.
Instance Attribute Details
#watcher_interval ⇒ Object
for test
26 27 28 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 26 def watcher_interval @watcher_interval end |
Instance Method Details
#cache_node(tag, node) ⇒ Object
46 47 48 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 46 def cache_node(tag, node) @node[keepforward(tag)] = node end |
#cache_sock(node, sock) ⇒ Object
217 218 219 220 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 217 def cache_sock(node, sock) get_sock[node] = sock get_sock_expired_at[node] = Time.now + @keepalive_time if @keepalive_time end |
#configure(conf) ⇒ Object
28 29 30 31 32 33 34 35 36 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 28 def configure(conf) super @node = {} @sock = {} @sock_expired_at = {} @mutex = {} @watcher_interval = 1 end |
#get_mutex(node) ⇒ Object
211 212 213 214 215 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 211 def get_mutex(node) thread_id = Thread.current.object_id @mutex[thread_id] ||= {} @mutex[thread_id][node] ||= Mutex.new end |
#get_node(tag) ⇒ Object
38 39 40 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 38 def get_node(tag) @node[keepforward(tag)] end |
#get_sock ⇒ Object
222 223 224 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 222 def get_sock @sock[Thread.current.object_id] ||= {} end |
#get_sock_expired_at ⇒ Object
226 227 228 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 226 def get_sock_expired_at @sock_expired_at[Thread.current.object_id] ||= {} end |
#keepforward(tag) ⇒ Object
42 43 44 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 42 def keepforward(tag) @keepforward == :one ? :one : tag end |
#reconnect(node) ⇒ Object
148 149 150 151 152 153 154 155 156 157 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 148 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 |
#send_data(node, tag, chunk) ⇒ Object
Override for keepalive
125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 125 def send_data(node, tag, chunk) get_mutex(node).synchronize do sock = get_sock[node] unless sock sock = reconnect(node) cache_sock(node, sock) if @keepalive end begin sock_write(sock, tag, chunk) node.heartbeat(false) rescue Errno::EPIPE, Errno::ECONNRESET, Errno::ECONNABORTED, Errno::ETIMEDOUT => e log.warn "out_keep_forward: send_data failed #{e.class} #{e.}, try to reconnect", :host=>node.host, :port=>node.port sock.close rescue IOError sock = reconnect(node) cache_sock(node, sock) if @keepalive retry ensure sock.close if sock and !@keepalive end end end |
#shutdown ⇒ Object
55 56 57 58 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 55 def shutdown super stop_watcher end |
#sock_close(node) ⇒ Object
202 203 204 205 206 207 208 209 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 202 def sock_close(node) get_mutex(node).synchronize do sock = get_sock[node] sock.close rescue IOError if sock get_sock[node] = nil get_sock_expired_at[node] = nil end end |
#sock_write(sock, tag, chunk) ⇒ Object
159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 159 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
50 51 52 53 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 50 def start super start_watcher end |
#start_watcher ⇒ Object
60 61 62 63 64 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 60 def start_watcher if @keepalive and @keepalive_time @watcher = Thread.new(&method(:watch_keepalive_time)) end end |
#stop_watcher ⇒ Object
66 67 68 69 70 71 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 66 def stop_watcher if @watcher @watcher.terminate @watcher.join end end |
#watch_keepalive_time ⇒ Object
watcher thread callback
184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 184 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 end end end end end |
#weight_send_data(tag, chunk, error_node = nil) ⇒ Object
93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 93 def weight_send_data(tag, chunk, error_node = nil) error = nil if error_node sock_close(error_node) if @keepalive and @keepforward == :one end wlen = @weight_array.length wlen.times do @rr = (@rr + 1) % wlen node = @weight_array[@rr] if node.available? begin send_data(node, tag, chunk) return node rescue # for load balancing during detecting crashed servers error = $! # use the latest error end end end cache_node(tag, nil) if error raise error else raise "no nodes are available" # TODO message end end |
#write_objects(tag, chunk) ⇒ Object
Override
74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 |
# File 'lib/fluent/plugin/out_keep_forward.rb', line 74 def write_objects(tag, chunk) return if chunk.empty? error = nil node = get_node(tag) if node and node.available? and (!@prefer_recover or @weight_array.include?(node)) begin send_data(node, tag, chunk) return rescue node = weight_send_data(tag, chunk, error_node = node) cache_node(tag, node) end else node = weight_send_data(tag, chunk, error_node = node) cache_node(tag, node) end end |