vectordotdev/vector · error

producer unexpectedly dropped

Error message

producer unexpectedly dropped

What it means

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.

Solutions

  1. Ensure the sink shuts down cleanly: cancel in-flight events before dropping the producer task.
  2. Check earlier logs for producer/rdkafka initialization errors that would kill the producer task.
  3. Upgrade Vector — this class of shutdown-race panic is typically fixed in newer releases.
  4. Reduce topology reload churn (frequent config reloads) while the kafka sink has queued events.
Defensive patterns

Strategy: retry

Try / catch

// topology-level: supervise Vector and restart on panic (systemd)
[Service]
Restart=on-failure
RestartSec=5

Prevention

When it happens

Trigger: 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.

Common situations: Sink shutdown races during topology reload/restart; an rdkafka producer creation failure earlier poisoned the task; process teardown while in-flight events are pending.

Understand the failure class

Background: "This is a bug, please report it": internal invariant violations, unreachable panics, and SNH errors explained — this error's family across 47 libraries.

Related errors


AI-assisted analysis of vectordotdev/vector@bdb87aeaa4 (2026-09-16). Data as JSON: /api/errors/41602a70fcc4a659. Report an issue: GitHub.

Appendix: source

Thrown at src/sinks/kafka/service.rs:154

            if let Some(timestamp) = request.metadata.timestamp_millis {
                record = record.timestamp(timestamp);
            }
            if let Some(headers) = request.metadata.headers {
                record = record.headers(headers);
            }

            // Manually poll [FutureProducer::send_result] instead of [FutureProducer::send] to track
            // records that fail to be enqueued on the producer.
            let mut blocked_state: Option<BlockedRecordState> = None;
            loop {
                match this.kafka_producer.send_result(record) {
                    // Record was successfully enqueued on the producer.
                    Ok(fut) => {
                        // Drop the blocked state (if any), as the producer is no longer blocked.
                        drop(blocked_state.take());
                        return fut
                            .await
                            .expect("producer unexpectedly dropped")
                            .map(|_| KafkaResponse {
                                event_byte_size,
                                raw_byte_size,
                                event_status: EventStatus::Delivered,
                            })
                            .map_err(|(err, _)| err);
                    }
                    // Producer queue is full or a policy has been violated and the request should
                    // be retried
                    Err((
                        KafkaError::MessageProduction(
                            RDKafkaErrorCode::QueueFull | RDKafkaErrorCode::PolicyViolation,
                        ),
                        original_record,
                    )) => {
                        if blocked_state.is_none() {
                            blocked_state =
                                Some(BlockedRecordState::new(Arc::clone(&this.records_blocked)));

View on GitHub (pinned to bdb87aeaa4)