{"record":{"id":"51398d2357c98f90","repo":"risingwavelabs/risingwave","slug":"primary-key-not-defined-for-kafka-sink-pleas","errorCode":null,"errorMessage":"primary key not defined for {:?} kafka sink (please define in `primary_key` field)","messagePattern":"primary key not defined for (.+?) kafka sink \\(please define in `primary_key` field\\)","errorType":"validation","errorClass":"SinkError::Config","httpStatus":null,"severity":"error","filePath":"src/connector/src/sink/kafka.rs","lineNumber":380,"sourceCode":"        .await?;\n        let max_delivery_buffer_size = (self\n            .config\n            .rdkafka_properties_producer\n            .queue_buffering_max_messages\n            .as_ref()\n            .cloned()\n            .unwrap_or(KAFKA_WRITER_MAX_QUEUE_SIZE) as f32\n            * KAFKA_WRITER_MAX_QUEUE_SIZE_RATIO) as usize;\n\n        Ok(KafkaSinkWriter::new(self.config.clone(), formatter)\n            .await?\n            .into_log_sinker(max_delivery_buffer_size))\n    }\n\n    async fn validate(&self) -> Result<()> {\n        // For non-append-only Kafka sink, the primary key must be defined.\n        if self.format_desc.format != SinkFormat::AppendOnly && self.pk_indices.is_empty() {\n            return Err(SinkError::Config(anyhow!(\n                \"primary key not defined for {:?} kafka sink (please define in `primary_key` field)\",\n                self.format_desc.format\n            )));\n        }\n        // Check for formatter constructor error, before it is too late for error reporting.\n        SinkFormatterImpl::new(\n            &self.format_desc,\n            self.schema.clone(),\n            self.pk_indices.clone(),\n            self.db_name.clone(),\n            self.sink_from_name.clone(),\n            &self.config.common.topic,\n        )\n        .await?;\n\n        // Try Kafka connection.\n        // There is no such interface for kafka producer to validate a connection\n        // use enumerator to validate broker reachability and existence of topic","sourceCodeStart":362,"sourceCodeEnd":398,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/sink/kafka.rs#L362-L398","documentation":"During Kafka sink validation, any non-AppendOnly format (e.g. UPSERT/DEBEZIUM) requires a primary key so change records can be keyed correctly. If pk_indices is empty for such a sink, validate() returns this SinkError::Config naming the offending format.","triggerScenarios":"Calling Sink::validate() on a KafkaSink whose format is UPSERT or DEBEZIUM and whose CREATE SINK statement did not define a primary_key column list.","commonSituations":"Creating an upsert Kafka sink over a source/table without declaring `primary_key` in the WITH/options; sink over a table whose primary key was dropped; planner passing empty pk_indices.","solutions":["Add a `primary_key` definition to the CREATE SINK statement","Use FORMAT PLAIN (append-only) if no key is needed","Ensure the sinked relation actually has a primary key when relying on implicit propagation"],"exampleFix":"// before\nCREATE SINK s FROM t WITH (connector='kafka', type='upsert', ...) FORMAT UPSERT ENCODE JSON;\n// after\nCREATE SINK s FROM t WITH (connector='kafka', type='upsert', primary_key='id', ...) FORMAT UPSERT ENCODE JSON;","handlingStrategy":"validation","validationCode":"fn kafka_pk_ok(format: SinkFormat, pk_indices: &[usize]) -> bool {\n    format == SinkFormat::AppendOnly || !pk_indices.is_empty()\n}","typeGuard":null,"tryCatchPattern":"if let Err(e) = sink.validate().await {\n    return Err(anyhow!(\"kafka sink validation failed: {e}\"));\n}","preventionTips":["Declare `primary_key` for UPSERT/DEBEZIUM kafka sinks","Choose FORMAT PLAIN if the stream is truly append-only","Ensure the sinked relation has a primary key when key propagation is expected"],"tags":["kafka","config","primary-key","validation"],"backgroundTag":"missing-required-config-field","analyzedSha":"6469eb736d691e8e9b8a419a57edd6429ca77417","analyzedAt":"2026-09-11T21:06:21.487Z","contentChangedAt":"2026-09-11T21:06:21.487Z","schemaVersion":2},"datasetVersion":"2026-09-14T16:17:12.679Z"}