diff --git a/README.md b/README.md index 589deb7..9ca58a5 100644 --- a/README.md +++ b/README.md @@ -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. diff --git a/lib/fluent/plugin/out_rdkafka2.rb b/lib/fluent/plugin/out_rdkafka2.rb index 0f0c7d5..4b0c469 100644 --- a/lib/fluent/plugin/out_rdkafka2.rb +++ b/lib/fluent/plugin/out_rdkafka2.rb @@ -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 @@ -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 diff --git a/test/plugin/test_out_rdkafka2.rb b/test/plugin/test_out_rdkafka2.rb index 5e81feb..017bf2a 100644 --- a/test/plugin/test_out_rdkafka2.rb +++ b/test/plugin/test_out_rdkafka2.rb @@ -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 @@ -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 @@ -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?