Class: LogStash::Outputs::RedisList
- Inherits:
-
Base
- Object
- Base
- LogStash::Outputs::RedisList
- Includes:
- Stud::Buffer
- Defined in:
- lib/logstash/outputs/redis_list.rb
Overview
This output will send events to a Redis queue using RPUSH. The RPUSH command is supported in Redis v0.0.7+. Using PUBLISH to a channel requires at least v1.3.8+. While you may be able to make these Redis versions work, the best performance and stability will be found in more recent stable versions. Versions 2.6.0+ are recommended.
For more information, see redis.io/[the Redis homepage]
Instance Method Summary collapse
- #close ⇒ Object
-
#congestion_check(key) ⇒ Object
def receive.
-
#flush(events, key, close = false) ⇒ Object
called from Stud::Buffer#buffer_flush when there are events to flush.
-
#on_flush_error(e) ⇒ Object
called from Stud::Buffer#buffer_flush when an error occurs.
-
#receive(event) ⇒ Object
def register.
- #register ⇒ Object
Instance Method Details
#close ⇒ Object
199 200 201 202 203 204 205 206 207 |
# File 'lib/logstash/outputs/redis_list.rb', line 199 def close if @batch buffer_flush(:final => true) end if @data_type == 'channel' and @redis @redis.quit @redis = nil end end |
#congestion_check(key) ⇒ Object
def receive
170 171 172 173 174 175 176 177 178 179 |
# File 'lib/logstash/outputs/redis_list.rb', line 170 def congestion_check(key) return if @congestion_threshold == 0 if (Time.now.to_i - @congestion_check_times[key]) >= @congestion_interval # Check congestion only if enough time has passed since last check. while @redis.llen(key) > @congestion_threshold # Don't push event to Redis key which has reached @congestion_threshold. @logger.warn? and @logger.warn("Redis key size has hit a congestion threshold #{@congestion_threshold} suspending output for #{@congestion_interval} seconds") sleep @congestion_interval end @congestion_check_times[key] = Time.now.to_i end end |
#flush(events, key, close = false) ⇒ Object
called from Stud::Buffer#buffer_flush when there are events to flush
182 183 184 185 186 187 188 |
# File 'lib/logstash/outputs/redis_list.rb', line 182 def flush(events, key, close=false) @redis ||= connect # we should not block due to congestion on close # to support this Stud::Buffer#buffer_flush should pass here the :final boolean value. congestion_check(key) unless close @redis.rpush(key, events) end |
#on_flush_error(e) ⇒ Object
called from Stud::Buffer#buffer_flush when an error occurs
190 191 192 193 194 195 196 197 |
# File 'lib/logstash/outputs/redis_list.rb', line 190 def on_flush_error(e) @logger.warn("Failed to send backlog of events to Redis", :identity => identity, :exception => e, :backtrace => e.backtrace ) @redis = connect end |
#receive(event) ⇒ Object
def register
152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 |
# File 'lib/logstash/outputs/redis_list.rb', line 152 def receive(event) # TODO(sissel): We really should not drop an event, but historically # we have dropped events that fail to be converted to json. # TODO(sissel): Find a way to continue passing events through even # if they fail to convert properly. begin @codec.encode(event) rescue LocalJumpError # This LocalJumpError rescue clause is required to test for regressions # for https://github.com/logstash-plugins/logstash-output-redis/issues/26 # see specs. Without it the LocalJumpError is rescued by the StandardError raise rescue StandardError => e @logger.warn("Error encoding event", :exception => e, :event => event) end end |
#register ⇒ Object
106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 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 150 |
# File 'lib/logstash/outputs/redis_list.rb', line 106 def register require 'redis' # TODO remove after setting key and data_type to true if @queue if @key or @data_type raise RuntimeError.new( "Cannot specify queue parameter and key or data_type" ) end @key = @queue @data_type = 'list' end if not @key or not @data_type raise RuntimeError.new( "Must define queue, or key and data_type parameters" ) end # end TODO if @batch if @data_type != "list" raise RuntimeError.new( "batch is not supported with data_type #{@data_type}" ) end buffer_initialize( :max_items => @batch_events, :max_interval => @batch_timeout, :logger => @logger ) end @redis = nil if @shuffle_hosts @host.shuffle! end @host_idx = 0 @congestion_check_times = Hash.new { |h,k| h[k] = Time.now.to_i - @congestion_interval } @codec.on_event(&method(:send_to_redis)) end |