instructure/canvas-lms · warning

Kafka delivery failed: #

Error message

Kafka delivery failed: #{event[:error]&.message}

What it means

Canvas::KafkaEvents::Producer builds a WaterDrop producer and subscribes to the 'error.occurred' monitor event. This is not a raise: when a Kafka message fails delivery (broker unreachable, timeouts, ack failures), the callback logs 'Kafka delivery failed: <message>' with the error's message and increments a Statsd counter.

Solutions

  1. Check the logged error message under 'Kafka delivery failed' for the root cause (connection refused vs timeout vs auth)
  2. Verify Kafka broker connectivity (host/port, TLS, SASL credentials) for the environment
  3. Raise message.timeout.ms / retry settings in the WaterDrop config if timeouts are transient
  4. Monitor the kafka_events.delivery_errors and kafka_events.messages_purged Statsd metrics and alert on spikes

Example fix

// before
c.kafka = @config.kafka_options
// after
c.kafka = @config.kafka_options.merge(
  'message.timeout.ms' => 30_000,
  'message.max.in.flight' => 1_000_000
)
Defensive patterns

Strategy: try-catch

Validate before calling

raise 'kafka config missing' unless @config&.kafka_options&.key?('bootstrap.servers')

Type guard

def kafka_configured?(config) = config&.kafka_options.is_a?(Hash) && config.kafka_options['bootstrap.servers'].present?

Try / catch

wd.monitor.subscribe('error.occurred') do |event|
  err = event[:error]
  Rails.logger.warn("Kafka delivery failed: #{err&.message}")
  retry_or_dead_letter(err)
end

Prevention

When it happens

Trigger: Publishing an event while the Kafka brokers are unreachable, message.timeout.ms elapses without a broker ack, or max in-flight/queue capacity is exceeded — WaterDrop fires error.occurred and the callback logs this line.

Common situations: Kafka cluster down or DNS misconfigured in an environment; auth/SSL mismatch with brokers; oversized messages; consumers down causing retention/broker pressure.

Related errors


AI-assisted analysis of instructure/canvas-lms@1c9f0bb801 (2026-09-15). Data as JSON: /api/errors/03d894af11ede291. Report an issue: GitHub.

Appendix: source

Thrown at lib/canvas/kafka_events/producer.rb:55

      nil
    end

    def close
      @producer&.close
    rescue => e
      Rails.logger.warn("Kafka producer close error: #{e.message}")
    end

    private

    def build_waterdrop
      wd = WaterDrop::Producer.new do |c|
        c.kafka = @config.kafka_options
        c.logger = Rails.logger
      end
      wd.monitor.subscribe("error.occurred") do |event|
        InstStatsd::Statsd.distributed_increment("kafka_events.delivery_errors")
        Rails.logger.warn("Kafka delivery failed: #{event[:error]&.message}")
      end
      wd.monitor.subscribe("message.purged") do |event|
        # Fires when a buffered message hits message.timeout.ms without broker ack —
        # the event is silently dropped on the floor. Worth an alert.
        InstStatsd::Statsd.distributed_increment("kafka_events.messages_purged")
        Rails.logger.warn("Kafka message purged: #{event[:error]&.message}")
      end
      wd
    end

    def log_ready
      resolved = Events.topic_keys.map { |key| "#{key}=#{@config.topic_for(key)}" }.join(" ")
      Rails.logger.info("Kafka events producer ready: brokers=#{@config.brokers} #{resolved}")
    end
  end
end

View on GitHub (pinned to 1c9f0bb801)