{"record":{"id":"4288c1591e64d4c5","repo":"risingwavelabs/risingwave","slug":"missing-format-encode-4288c1","errorCode":null,"errorMessage":"missing FORMAT ... ENCODE ...","messagePattern":"missing FORMAT \\.\\.\\. ENCODE \\.\\.\\.","errorType":"validation","errorClass":"SinkError::Config","httpStatus":null,"severity":"error","filePath":"src/connector/src/sink/kafka.rs","lineNumber":339,"sourceCode":"        }\n        Ok(())\n    }\n}\n\nimpl TryFrom<SinkParam> for KafkaSink {\n    type Error = SinkError;\n\n    fn try_from(param: SinkParam) -> std::result::Result<Self, Self::Error> {\n        let schema = param.schema();\n        let pk_indices = param.downstream_pk_or_empty();\n        let config = KafkaConfig::from_btreemap(param.properties)?;\n        Ok(Self {\n            config,\n            schema,\n            pk_indices,\n            format_desc: param\n                .format_desc\n                .ok_or_else(|| SinkError::Config(anyhow!(\"missing FORMAT ... ENCODE ...\")))?,\n            db_name: param.db_name,\n            sink_from_name: param.sink_from_name,\n        })\n    }\n}\n\nimpl Sink for KafkaSink {\n    type LogSinker = AsyncTruncateLogSinkerOf<KafkaSinkWriter>;\n\n    const SINK_NAME: &'static str = KAFKA_SINK;\n\n    crate::impl_validate_sink_unknown_fields!();\n\n    async fn new_log_sinker(&self, _writer_param: SinkWriterParam) -> Result<Self::LogSinker> {\n        let formatter = SinkFormatterImpl::new(\n            &self.format_desc,\n            self.schema.clone(),\n            self.pk_indices.clone(),","sourceCodeStart":321,"sourceCodeEnd":357,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/sink/kafka.rs#L321-L357","documentation":"KafkaSink is constructed from SinkParam and requires a parsed format descriptor (FORMAT ... ENCODE ... from the DDL). If param.format_desc is None, construction fails with this SinkError::Config, because the sink cannot know how to serialize messages without a declared format/encode.","triggerScenarios":"Creating a Kafka sink whose DDL lacks FORMAT ... ENCODE ... clauses, or an internal path building KafkaSink from a SinkParam where format_desc was never parsed (e.g. malformed or missing FORMAT/ENCODE in CREATE SINK).","commonSituations":"Writing CREATE SINK statements without `FORMAT PLAIN/DEBEZIUM/UPSERT ENCODE JSON/AVRO/...`; constructing sinks programmatically with a hand-built SinkParam; upgrading from older syntax without format declarations.","solutions":["Add FORMAT ... ENCODE ... to the CREATE SINK statement (e.g. FORMAT PLAIN ENCODE JSON)","Verify the statement parses the format by checking desc/param before sink creation","When building SinkParam in code, populate format_desc from parsed options"],"exampleFix":"// before\nCREATE SINK s FROM t WITH (connector='kafka', properties.bootstrap.server='b:9092');\n// after\nCREATE SINK s FROM t WITH (connector='kafka', properties.bootstrap.server='b:9092') FORMAT PLAIN ENCODE JSON;","handlingStrategy":"validation","validationCode":"fn ensure_format_desc(param: &SinkParam) -> Result<()> {\n    if param.format_desc.is_none() { bail!(\"sink requires FORMAT ... ENCODE ...\"); }\n    Ok(())\n}","typeGuard":null,"tryCatchPattern":"let format_desc = param.format_desc.as_ref()\n    .ok_or_else(|| anyhow!(\"sink created without FORMAT ... ENCODE ...\"))?;","preventionTips":["Always include FORMAT ... ENCODE ... in CREATE SINK for kafka/kinesis sinks","Copy the format clause from a known-working sink DDL template","Validate statements with EXPLAIN/parse before deployment"],"tags":["kafka","config","format","sql"],"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"}