{"record":{"id":"0a41ea8fd71e3ddc","repo":"instructure/canvas-lms","slug":"kafka-message-purged-event-error-message","errorCode":null,"errorMessage":"Kafka message purged: #{event[:error]&.message}","messagePattern":"Kafka message purged: #(.+?)","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"lib/canvas/kafka_events/producer.rb","lineNumber":61,"sourceCode":"      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":43,"sourceCodeEnd":72,"githubUrl":"https://github.com/instructure/canvas-lms/blob/1c9f0bb8013ed69c4f2efe11fd483025469b7e6c/lib/canvas/kafka_events/producer.rb#L43-L72","documentation":"This is a warning emitted by the waterdrop Kafka producer monitoring subscriber for the 'message.purged' event. It fires when a buffered message is dropped because it hit 'message.timeout.ms' without receiving a broker acknowledgment; the message is silently discarded. Canvas logs it and increments the kafka_events.messages_purged Statsd counter, and the code comment explicitly says it is worth alerting on.","triggerScenarios":"Producing a Canvas event via the waterdrop producer whose buffered message exceeds message.timeout.ms waiting for a broker ack, e.g. Kafka brokers unreachable, slow disk flush on the broker, or producer-side retries exhausted before ack timeout.","commonSituations":"Kafka cluster outage or broker restart during event production; misconfigured message.timeout.ms too low relative to broker latency; network partition between app and Kafka; broker acks settings (e.g. acks=all) with an under-replicated slow topic.","solutions":["Check Kafka broker health and connectivity from the app host (kafka-topics --describe, broker logs) to find why acks were not delivered.","Raise message.timeout.ms (and delivery-related timeouts) in the waterdrop producer config so normal latency does not trigger purging.","Verify topic replication/ISR is healthy; acks=all with degraded ISR can exceed the timeout.","Add an alert on kafka_events.messages_purged as the code comment suggests, and replay the lost events from the source if data loss matters."],"exampleFix":"# before\nconfig[:message_timeout_ms] = 5000\n\n# after\nconfig[:message_timeout_ms] = 30000","handlingStrategy":"retry","validationCode":"if wd.config[:message_timeout_ms] < 30_000\n  Rails.logger.warn(\"kafka message.timeout.ms very low; messages may be purged\")\nend\n# also verify broker reachability before producing\n# Kafka::Admin or TCP probe to bootstrap servers","typeGuard":"def broker_reachable?(seed_servers)\n  seed_servers.all? { |host, port| TCPSocket.new(host, port).close; true }\nrescue Errno::ECONNREFUSED, SocketError\n  false\nend","tryCatchPattern":"begin\n  producer.produce_async(payload)\nrescue WaterDrop::Errors::ProduceError => e\n  InstStatsd::Statsd.increment(\"kafka_events.produce_failed\")\n  # buffer/replay payload later; message.purged means data loss, not an exception\nend","preventionTips":["Set message.timeout.ms well above worst-case broker latency","Monitor kafka_events.messages_purged and alert on any nonzero rate","Ensure topic min.insync.replicas and ISR health before enabling acks=all","Keep a replayable log of produced events so purged messages can be re-sent"],"tags":["kafka","message-loss","timeout","event-producer"],"backgroundTag":"request-timeout","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"}