Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -571,6 +571,11 @@ You need to install rdkafka gem.
`rdkafka2` supports `discard_kafka_delivery_failed_regex` parameter:
- `discard_kafka_delivery_failed_regex` - default: nil - discard the record where the Kafka::DeliveryFailed occurred and the emitted message matches the given regex pattern, such as `/unknown_topic/`.

`rdkafka2` supports `unrecoverable_error_codes` parameter:
- `unrecoverable_error_codes` - default: `["topic_authorization_failed", "msg_size_too_large"]` - librdkafka error codes to treat as unrecoverable errors. When a listed code is returned, `Fluent::UnrecoverableError` is raised, so Fluentd hands the whole chunk to the secondary output or to the backup directory instead of retrying it until `retry_timeout` expires. Both the produce call and the delivery report are covered, the latter only while `rdkafka_delivery_handle_poll_timeout` is not 0. Set the parameter to an empty value keeping retrying every error, as this plugin did before v0.19.x.
- A code is spelled the way `Rdkafka::RdkafkaError#code` reports it, that is the librdkafka name lowercased with the leading underscore removed, so `RD_KAFKA_RESP_ERR__UNKNOWN_PARTITION` is written as `unknown_partition`. A name that matches no error code is accepted and never matches anything.
- `discard_kafka_delivery_failed` and `discard_kafka_delivery_failed_regex` take precedence over this parameter. The regex is matched against the librdkafka message, such as `Broker: Message size too large (msg_size_too_large)`.

If you use v0.12, use `rdkafka` instead.

<match kafka.**>
Expand Down
12 changes: 6 additions & 6 deletions lib/fluent/plugin/out_rdkafka2.rb
Original file line number Diff line number Diff line change
Expand Up @@ -469,6 +469,7 @@ def write(chunk)
log.warn "Delivery failed and matched regexp pattern #{@discard_kafka_delivery_failed_regex}. Discard events:", :error => e.to_s, :error_class => e.class.to_s, :tag => tag
else
log.warn "Send exception occurred: #{e} at #{e.backtrace.first}"
raise Fluent::UnrecoverableError, "Rejected due to #{e}" if unrecoverable_error?(e)
# Raise exception to retry sendind messages
raise e
end
Expand Down Expand Up @@ -525,15 +526,14 @@ def enqueue_with_retry(producer, topic, record_buf, message_key, partition, head

raise e
else
if unrecoverable_error_codes.include?(e.code.to_s)
# some of the errors should be handled as an unrecoverable error
raise Fluent::UnrecoverableError, "Rejected due to #{e}"
else
raise e
end
raise e
end
end
end
end

def unrecoverable_error?(e)
e.is_a?(Rdkafka::RdkafkaError) && unrecoverable_error_codes.include?(e.code.to_s)
end
end
end
126 changes: 124 additions & 2 deletions test/plugin/test_out_rdkafka2.rb
Original file line number Diff line number Diff line change
Expand Up @@ -38,15 +38,22 @@ def create_driver(conf = config, tag='test')
end

class DummyDeliveryHandle
def initialize(error = nil)
@error = error
end

def wait(**)
raise @error if @error
end
end

class DummyProducer
attr_reader :produced

def initialize
def initialize(produce_error: nil, delivery_error: nil)
@produced = []
@produce_error = produce_error
@delivery_error = delivery_error
end

# Raises what the rdkafka FFI binding raises: the partition must be an
Expand All @@ -57,9 +64,19 @@ def produce(topic:, payload:, key:, partition:, headers:, timestamp:)
raise RangeError, "integer #{partition} too big to convert to 'int'" unless (-2**31..2**31 - 1).cover?(partition)
end
key.bytesize unless key.nil?
raise @produce_error if @produce_error

@produced << {payload: payload, key: key, partition: partition}
DummyDeliveryHandle.new
DummyDeliveryHandle.new(@delivery_error)
end
end

RD_KAFKA_RESP_ERR__UNKNOWN_PARTITION = -190
RD_KAFKA_RESP_ERR_MSG_SIZE_TOO_LARGE = 10

class DummyErrorWithCode < StandardError
def code
:msg_size_too_large
end
end

Expand Down Expand Up @@ -276,6 +293,111 @@ def test_write_converts_non_string_message_key
assert_equal ["12345"], producer.produced.collect { |message| message[:key] }
end

def test_rdkafka_error_codes_under_test
assert_equal [:unknown_partition, :msg_size_too_large],
[Rdkafka::RdkafkaError.new(RD_KAFKA_RESP_ERR__UNKNOWN_PARTITION).code,
Rdkafka::RdkafkaError.new(RD_KAFKA_RESP_ERR_MSG_SIZE_TOO_LARGE).code]
end

data("produce" => :produce_error,
"delivery report" => :delivery_error)
def test_write_raises_unrecoverable_error_for_listed_error_code(path)
d = create_driver
producer = DummyProducer.new(path => Rdkafka::RdkafkaError.new(RD_KAFKA_RESP_ERR_MSG_SIZE_TOO_LARGE))
stub(d.instance).get_producer { producer }

assert_raise(Fluent::UnrecoverableError) {
d.instance.write(create_chunk([{"a" => "b"}]))
}
end

data("produce" => :produce_error,
"delivery report" => :delivery_error)
def test_write_reraises_unlisted_error_code(path)
d = create_driver
producer = DummyProducer.new(path => Rdkafka::RdkafkaError.new(RD_KAFKA_RESP_ERR__UNKNOWN_PARTITION))
stub(d.instance).get_producer { producer }

assert_raise(Rdkafka::RdkafkaError) {
d.instance.write(create_chunk([{"a" => "b"}]))
}
end

data("produce" => :produce_error,
"delivery report" => :delivery_error)
def test_write_raises_unrecoverable_error_for_unknown_partition_when_listed(path)
d = create_driver(config + config_element('ROOT', '', {"unrecoverable_error_codes" => "unknown_partition"}))
producer = DummyProducer.new(path => Rdkafka::RdkafkaError.new(RD_KAFKA_RESP_ERR__UNKNOWN_PARTITION))
stub(d.instance).get_producer { producer }

assert_raise(Fluent::UnrecoverableError) {
d.instance.write(create_chunk([{"a" => "b"}]))
}
end

data("produce" => :produce_error,
"delivery report" => :delivery_error)
def test_write_discards_listed_error_code_when_discard_kafka_delivery_failed(path)
d = create_driver(config + config_element('ROOT', '', {"discard_kafka_delivery_failed" => "true"}))
producer = DummyProducer.new(path => Rdkafka::RdkafkaError.new(RD_KAFKA_RESP_ERR_MSG_SIZE_TOO_LARGE))
stub(d.instance).get_producer { producer }

assert_nothing_raised {
d.instance.write(create_chunk([{"a" => "b"}]))
}
assert_true d.logs.any? { |log| log.include?("Delivery failed. Discard events") }
end

data("produce" => :produce_error,
"delivery report" => :delivery_error)
def test_write_reraises_error_that_is_not_a_kafka_error(path)
d = create_driver
producer = DummyProducer.new(path => DummyErrorWithCode.new)
stub(d.instance).get_producer { producer }

assert_raise(DummyErrorWithCode) {
d.instance.write(create_chunk([{"a" => "b"}]))
}
end

data("produce" => :produce_error,
"delivery report" => :delivery_error)
def test_write_discards_listed_error_code_matching_discard_kafka_delivery_failed_regex(path)
d = create_driver(config + config_element('ROOT', '', {"discard_kafka_delivery_failed_regex" => "/msg_size_too_large/"}))
producer = DummyProducer.new(path => Rdkafka::RdkafkaError.new(RD_KAFKA_RESP_ERR_MSG_SIZE_TOO_LARGE))
stub(d.instance).get_producer { producer }

assert_nothing_raised {
d.instance.write(create_chunk([{"a" => "b"}]))
}
assert_true d.logs.any? { |log| log.include?("Delivery failed and matched regexp pattern") }
end

def test_write_produces_every_record_before_raising_unrecoverable_error
d = create_driver
producer = DummyProducer.new(delivery_error: Rdkafka::RdkafkaError.new(RD_KAFKA_RESP_ERR_MSG_SIZE_TOO_LARGE))
stub(d.instance).get_producer { producer }

assert_raise(Fluent::UnrecoverableError) {
d.instance.write(create_chunk([{"a" => "b"}, {"a" => "c"}, {"a" => "d"}]))
}

assert_equal ['{"a":"b"}', '{"a":"c"}', '{"a":"d"}'], producer.produced.collect { |message| message[:payload] }
end

data("produce" => :produce_error,
"delivery report" => :delivery_error)
def test_write_reraises_listed_error_code_when_unrecoverable_error_codes_is_empty(path)
d = create_driver(config + config_element('ROOT', '', {"unrecoverable_error_codes" => ""}))
producer = DummyProducer.new(path => Rdkafka::RdkafkaError.new(RD_KAFKA_RESP_ERR_MSG_SIZE_TOO_LARGE))
stub(d.instance).get_producer { producer }

assert_equal [], d.instance.unrecoverable_error_codes
assert_raise(Rdkafka::RdkafkaError) {
d.instance.write(create_chunk([{"a" => "b"}]))
}
end

def test_mutli_worker_support
d = create_driver
assert_equal true, d.instance.multi_workers_ready?
Expand Down