{"record":{"id":"cfe4762c633aba27","repo":"risingwavelabs/risingwave","slug":"primary-key-mismatch-postgres-table-has-primary-k-cfe476","errorCode":null,"errorMessage":"Primary key mismatch: Postgres table has primary key on column `{}`, but sink schema does not define it as a primary key","messagePattern":"Primary key mismatch: Postgres table has primary key on column `(.+?)`, but sink schema does not define it as a primary key","errorType":"validation","errorClass":"SinkError::Config","httpStatus":null,"severity":"error","filePath":"src/connector/src/sink/postgres.rs","lineNumber":323,"sourceCode":"\n            // check that pk matches\n            {\n                let pg_pk_names = pg_table.pk_names();\n                let sink_pk_names = self\n                    .pk_indices\n                    .iter()\n                    .map(|i| &self.schema.fields()[*i].name)\n                    .collect::<HashSet<_>>();\n                if pg_pk_names.len() != sink_pk_names.len() {\n                    return Err(SinkError::Config(anyhow!(\n                        \"Primary key mismatch: Postgres table has primary key on columns {:?}, but sink schema defines primary key on columns {:?}\",\n                        pg_pk_names,\n                        sink_pk_names\n                    )));\n                }\n                for name in pg_pk_names {\n                    if !sink_pk_names.contains(name) {\n                        return Err(SinkError::Config(anyhow!(\n                            \"Primary key mismatch: Postgres table has primary key on column `{}`, but sink schema does not define it as a primary key\",\n                            name\n                        )));\n                    }\n                }\n            }\n        }\n\n        Ok(())\n    }\n\n    async fn new_log_sinker(&self, writer_param: SinkWriterParam) -> Result<Self::LogSinker> {\n        let writer = PostgresSinkWriter::new(\n            self.config.clone(),\n            self.schema.clone(),\n            self.pk_indices.clone(),\n            self.is_append_only,\n            &writer_param,","sourceCodeStart":305,"sourceCodeEnd":341,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/sink/postgres.rs#L305-L341","documentation":"After the PK count check passes, validation iterates the Postgres table's PK column names and requires each to appear in the sink's declared PK set. Any PG PK column missing from the sink's primary_key definition aborts CREATE SINK for that single column.","triggerScenarios":"PG table PK is composite (a,b) and the sink declares primary_key='a,b' but with a renamed column, or a name differs by case/alias so sink_pk_names lacks the PG name.","commonSituations":"Aliasing a PK column in the sink query so its name no longer matches the PG PK column; case-sensitivity mismatch; PG PK changed since the sink config was written.","solutions":["Rename the column in the sink query (AS <pg_pk_name>) so its name matches the Postgres primary key column.","Update the sink's `primary_key` option to use the exact PG PK column names, then recreate the sink.","ALTER the Postgres table PK to match the names used in the sink schema."],"exampleFix":"-- before (pg PK on 'user_id')\nCREATE SINK s AS SELECT user_id AS key FROM mv WITH (connector='postgres', table='t', primary_key='key');\n-- after\nCREATE SINK s AS SELECT user_id FROM mv WITH (connector='postgres', table='t', primary_key='user_id');","handlingStrategy":"validation","validationCode":"-- every PG PK column name must equal a sink schema column name and be in primary_key\nSELECT conkey, conname FROM pg_constraint WHERE conrelid='t'::regclass AND contype='p';","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Do not alias away PK column names in the sink query.","Keep identifier casing consistent between RW and PG."],"tags":["rust","postgres","sink","primary-key"],"backgroundTag":"schema-validation-failed","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"}