risingwavelabs/risingwave · error · SinkError::Mongodb

cannot serialize primary key

Error message

cannot serialize primary key

What it means

Raised in UpsertCommandBuilder::add_upsert when mongodb::bson::to_vec fails to serialize the primary key document into BSON bytes. This indicates the generated primary key document contains values unrepresentable in BSON (e.g. unsupported types), a programming/encoding bug rather than user misconfiguration.

Solutions

  1. Check that the primary key document contains only BSON-serializable values.
  2. Verify the primary key schema does not contain unsupported types (e.g., NaN/Infinity floats, invalid datetime).
  3. Fix upstream data types or the sink PK definition so the key can be serialized.
Defensive patterns

Strategy: try-catch

When it happens

Trigger: Thrown at src/connector/src/sink/mongodb.rs:551 when the library encounters an invalid state.

Common situations: See trigger scenarios.


AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11). Data as JSON: /api/errors/97e808b994d9a8b9. Report an issue: GitHub.

Appendix: source

Thrown at src/connector/src/sink/mongodb.rs:551

struct UpsertCommandBuilder {
    coll: String,
    updates: Array,
    deletes: HashMap<Vec<u8>, Document>,
}

impl UpsertCommandBuilder {
    fn new(coll: String) -> Self {
        Self {
            coll,
            updates: Array::new(),
            deletes: HashMap::new(),
        }
    }

    fn add_upsert(&mut self, pk: Document, row: Document) -> Result<()> {
        let pk_data = mongodb::bson::to_vec(&pk).map_err(|err| {
            SinkError::Mongodb(anyhow!(err).context("cannot serialize primary key"))
        })?;
        // under same pk, if the record currently being upserted was marked for deletion previously, we should
        // revert the deletion, otherwise, the upserting record may be accidentally deleted.
        // see https://github.com/risingwavelabs/risingwave/pull/17102#discussion_r1630684160 for more information.
        self.deletes.remove(&pk_data);

        self.updates.push(bson!( {
            "q": pk,
            "u": bson!( {
                "$set": row,
            }),
            "upsert": true,
            "multi": false,
        }));

        Ok(())
    }

View on GitHub (pinned to 6469eb736d)