{"record":{"id":"03d894af11ede291","repo":"instructure/canvas-lms","slug":"kafka-delivery-failed-event-error-message","errorCode":null,"errorMessage":"Kafka delivery failed: #{event[:error]&.message}","messagePattern":"Kafka delivery failed: #(.+?)","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"lib/canvas/kafka_events/producer.rb","lineNumber":55,"sourceCode":"      nil\n    end\n\n    def close\n      @producer&.close\n    rescue => e\n      Rails.logger.warn(\"Kafka producer close error: #{e.message}\")\n    end\n\n    private\n\n    def build_waterdrop\n      wd = WaterDrop::Producer.new do |c|\n        c.kafka = @config.kafka_options\n        c.logger = Rails.logger\n      end\n      wd.monitor.subscribe(\"error.occurred\") do |event|\n        InstStatsd::Statsd.distributed_increment(\"kafka_events.delivery_errors\")\n        Rails.logger.warn(\"Kafka delivery failed: #{event[:error]&.message}\")\n      end\n      wd.monitor.subscribe(\"message.purged\") do |event|\n        # Fires when a buffered message hits message.timeout.ms without broker ack —\n        # the event is silently dropped on the floor. Worth an alert.\n        InstStatsd::Statsd.distributed_increment(\"kafka_events.messages_purged\")\n        Rails.logger.warn(\"Kafka message purged: #{event[:error]&.message}\")\n      end\n      wd\n    end\n\n    def log_ready\n      resolved = Events.topic_keys.map { |key| \"#{key}=#{@config.topic_for(key)}\" }.join(\" \")\n      Rails.logger.info(\"Kafka events producer ready: brokers=#{@config.brokers} #{resolved}\")\n    end\n  end\nend\n","sourceCodeStart":37,"sourceCodeEnd":72,"githubUrl":"https://github.com/instructure/canvas-lms/blob/1c9f0bb8013ed69c4f2efe11fd483025469b7e6c/lib/canvas/kafka_events/producer.rb#L37-L72","documentation":"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.","triggerScenarios":"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.","commonSituations":"Kafka cluster down or DNS misconfigured in an environment; auth/SSL mismatch with brokers; oversized messages; consumers down causing retention/broker pressure.","solutions":["Check the logged error message under 'Kafka delivery failed' for the root cause (connection refused vs timeout vs auth)","Verify Kafka broker connectivity (host/port, TLS, SASL credentials) for the environment","Raise message.timeout.ms / retry settings in the WaterDrop config if timeouts are transient","Monitor the kafka_events.delivery_errors and kafka_events.messages_purged Statsd metrics and alert on spikes"],"exampleFix":"// before\nc.kafka = @config.kafka_options\n// after\nc.kafka = @config.kafka_options.merge(\n  'message.timeout.ms' => 30_000,\n  'message.max.in.flight' => 1_000_000\n)","handlingStrategy":"try-catch","validationCode":"raise 'kafka config missing' unless @config&.kafka_options&.key?('bootstrap.servers')","typeGuard":"def kafka_configured?(config) = config&.kafka_options.is_a?(Hash) && config.kafka_options['bootstrap.servers'].present?","tryCatchPattern":"wd.monitor.subscribe('error.occurred') do |event|\n  err = event[:error]\n  Rails.logger.warn(\"Kafka delivery failed: #{err&.message}\")\n  retry_or_dead_letter(err)\nend","preventionTips":["Alert on kafka_events.delivery_errors metrics","Verify broker connectivity, TLS, and SASL before deploy","Tune message.timeout.ms for transient broker latency","Implement dead-letter handling for purged messages"],"tags":["kafka","waterdrop","logging"],"backgroundTag":"network-request-failed","analyzedSha":"1c9f0bb8013ed69c4f2efe11fd483025469b7e6c","analyzedAt":"2026-09-15T20:33:18.891Z","contentChangedAt":"2026-09-15T20:33:18.891Z","schemaVersion":2},"datasetVersion":"2026-09-23T02:17:17.105Z"}