{"record":{"id":"41602a70fcc4a659","repo":"vectordotdev/vector","slug":"producer-unexpectedly-dropped","errorCode":null,"errorMessage":"producer unexpectedly dropped","messagePattern":"producer unexpectedly dropped","errorType":"panic","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/sinks/kafka/service.rs","lineNumber":154,"sourceCode":"            if let Some(timestamp) = request.metadata.timestamp_millis {\n                record = record.timestamp(timestamp);\n            }\n            if let Some(headers) = request.metadata.headers {\n                record = record.headers(headers);\n            }\n\n            // Manually poll [FutureProducer::send_result] instead of [FutureProducer::send] to track\n            // records that fail to be enqueued on the producer.\n            let mut blocked_state: Option<BlockedRecordState> = None;\n            loop {\n                match this.kafka_producer.send_result(record) {\n                    // Record was successfully enqueued on the producer.\n                    Ok(fut) => {\n                        // Drop the blocked state (if any), as the producer is no longer blocked.\n                        drop(blocked_state.take());\n                        return fut\n                            .await\n                            .expect(\"producer unexpectedly dropped\")\n                            .map(|_| KafkaResponse {\n                                event_byte_size,\n                                raw_byte_size,\n                                event_status: EventStatus::Delivered,\n                            })\n                            .map_err(|(err, _)| err);\n                    }\n                    // Producer queue is full or a policy has been violated and the request should\n                    // be retried\n                    Err((\n                        KafkaError::MessageProduction(\n                            RDKafkaErrorCode::QueueFull | RDKafkaErrorCode::PolicyViolation,\n                        ),\n                        original_record,\n                    )) => {\n                        if blocked_state.is_none() {\n                            blocked_state =\n                                Some(BlockedRecordState::new(Arc::clone(&this.records_blocked)));","sourceCodeStart":136,"sourceCodeEnd":172,"githubUrl":"https://github.com/vectordotdev/vector/blob/bdb87aeaa4c4ff27c0ba643c1c77b21bf2ef4013/src/sinks/kafka/service.rs#L136-L172","documentation":"Kafka sink's per-request future polls a oneshot/stream that carries the enqueue result from the shared producer task. If that channel yields None, the producer task was dropped without replying, which the code treats as an impossible internal state — hence expect(\"producer unexpectedly dropped\"). It signals a broken internal channel between the sink service and the Kafka producer worker, not a Kafka broker problem.","triggerScenarios":"The Kafka producer future/task holding the rdkafka producer is dropped or shut down while service calls are still awaiting enqueue results; awaiting fut after the producer background task has exited.","commonSituations":"Sink shutdown races during topology reload/restart; an rdkafka producer creation failure earlier poisoned the task; process teardown while in-flight events are pending.","solutions":["Ensure the sink shuts down cleanly: cancel in-flight events before dropping the producer task.","Check earlier logs for producer/rdkafka initialization errors that would kill the producer task.","Upgrade Vector — this class of shutdown-race panic is typically fixed in newer releases.","Reduce topology reload churn (frequent config reloads) while the kafka sink has queued events."],"exampleFix":null,"handlingStrategy":"retry","validationCode":null,"typeGuard":null,"tryCatchPattern":"// topology-level: supervise Vector and restart on panic (systemd)\n[Service]\nRestart=on-failure\nRestartSec=5","preventionTips":["Avoid frequent config reloads while the kafka sink has a full buffer.","Monitor Vector logs for kafka producer initialization warnings.","Keep Vector upgraded to pick up shutdown-race fixes.","Use graceful shutdown (SIGTERM, not SIGKILL) for the Vector process."],"tags":["rust","panic","kafka","concurrency","sink"],"backgroundTag":"internal-invariant-violation","analyzedSha":"bdb87aeaa4c4ff27c0ba643c1c77b21bf2ef4013","analyzedAt":"2026-09-16T02:53:35.741Z","contentChangedAt":"2026-09-16T02:53:35.741Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}