{"record":{"id":"fd6c1cc89e4c81b4","repo":"risingwavelabs/risingwave","slug":"es-sink-only-supports-single-pk-or-pk-with-delimit","errorCode":null,"errorMessage":"Es sink only supports single pk or pk with delimiter option","messagePattern":"Es sink only supports single pk or pk with delimiter option","errorType":"validation","errorClass":"anyhow","httpStatus":null,"severity":"error","filePath":"src/connector/src/sink/remote.rs","lineNumber":192,"sourceCode":"    }\n\n    async fn validate(&self) -> Result<()> {\n        validate_remote_sink(&self.param, Self::SINK_NAME).await?;\n        Ok(())\n    }\n}\n\nasync fn validate_remote_sink(param: &SinkParam, sink_name: &str) -> ConnectorResult<()> {\n    // if sink_name == OpenSearchJavaSink::SINK_NAME {\n    //     risingwave_common::license::Feature::OpenSearchSink\n    //         .check_available()\n    //         .map_err(|e| anyhow::anyhow!(e))?;\n    // }\n    if is_remote_es_sink(sink_name)\n        && param.downstream_pk_or_empty().len() > 1\n        && !param.properties.contains_key(ES_OPTION_DELIMITER)\n    {\n        bail!(\"Es sink only supports single pk or pk with delimiter option\");\n    }\n    // FIXME: support struct and array in stream sink\n    param.columns.iter().try_for_each(|col| {\n        match &col.data_type {\n            DataType::Int16\n                    | DataType::Int32\n                    | DataType::Int64\n                    | DataType::Float32\n                    | DataType::Float64\n                    | DataType::Boolean\n                    | DataType::Decimal\n                    | DataType::Timestamp\n                    | DataType::Timestamptz\n                    | DataType::Varchar\n                    | DataType::Date\n                    | DataType::Time\n                    | DataType::Interval\n                    | DataType::Jsonb","sourceCodeStart":174,"sourceCodeEnd":210,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/sink/remote.rs#L174-L210","documentation":"When validating a remote Elasticsearch sink, `validate_remote_sink` requires that the downstream primary key either be a single column or, if it is a multi-column PK, that the `delimiter` option be provided so ES document ids can be built by joining PK values. Otherwise validation bails with this error.","triggerScenarios":"Create a remote Elasticsearch sink whose table has a composite (2+) primary key and no `delimiter` option in the sink properties.","commonSituations":"Users sink tables with natural composite keys (e.g. (user_id, event_id)) to ES without realizing document id generation needs a delimiter; frameworks auto-generating sinks from wide PK schemas.","solutions":["Add `delimiter = '<char>'` to the sink properties so multi-column PKs can be joined into document ids.","Restructure the sink to have a single-column primary key (e.g. via a view with a synthetic id).","If the sink is not actually Elasticsearch (`is_remote_es_sink` check), verify the sink name/type to confirm which validation path applies."],"exampleFix":"// before\nCREATE SINK es_sink FROM t WITH (connector='elasticsearch', type='append-only');  -- composite pk\n// after\nCREATE SINK es_sink FROM t WITH (connector='elasticsearch', type='append-only', delimiter='-');","handlingStrategy":"validation","validationCode":"if (connector === 'elasticsearch' && pkColumns.length > 1 && !opts.delimiter) throw new Error('composite pk requires delimiter option');","typeGuard":null,"tryCatchPattern":"try { await createSink(cfg); } catch (e) { if (/single pk or pk with delimiter/.test(String(e))) { cfg.delimiter = '-'; await createSink(cfg); } else throw e; }","preventionTips":["Always set delimiter when sinking composite-key tables to Elasticsearch","Inspect downstream_pk before sink creation","Prefer single-column synthetic ids for ES sinks"],"tags":["elasticsearch","sink","primary-key","config-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-23T08:17:48.524Z"}