{"record":{"id":"009578cfa012e327","repo":"risingwavelabs/risingwave","slug":"register-v3-sink-worker","errorCode":null,"errorMessage":"register v3 sink worker","messagePattern":"register v3 sink worker","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/meta/src/rpc/ddl_controller.rs","lineNumber":1384,"sourceCode":"                }\n                // Validate the sink on the connector node.\n                validate_sink(sink).await?;\n                // For Iceberg pk-index sinks, spawn the per-sink commit worker now\n                // so it's ready to receive epoch reports from the very first\n                // barrier instead of relying on lazy registration on every\n                // commit.\n                if crate::manager::iceberg_pk_index_sink::is_iceberg_pk_index_sink(&sink.properties)\n                {\n                    let iceberg_config =\n                        crate::manager::iceberg_pk_index_sink::build_iceberg_config(sink)?;\n                    self.iceberg_pk_index_sink_manager\n                        .register_sink(\n                            sink.id,\n                            crate::barrier::to_partial_graph_id(sink.database_id, None),\n                            iceberg_config,\n                        )\n                        .await\n                        .map_err(|e| anyhow!(e).context(\"register v3 sink worker\"))?;\n                }\n                let connector_name = sink.get_properties().get(UPSTREAM_SOURCE_KEY).cloned();\n                let attr = sink.format_desc.as_ref().map(|sink_info| {\n                    jsonbb::json!({\n                        \"format\": sink_info.format().as_str_name(),\n                        \"encode\": sink_info.encode().as_str_name(),\n                    })\n                });\n                report_create_object(\n                    streaming_job.id(),\n                    \"sink\",\n                    PbTelemetryDatabaseObject::Sink,\n                    connector_name,\n                    attr,\n                );\n            }\n            StreamingJob::Source(source) => {\n                // Register the source on the connector node.","sourceCodeStart":1366,"sourceCodeEnd":1402,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/meta/src/rpc/ddl_controller.rs#L1366-L1402","documentation":"For sink v3 (stream sink framework), generate_streaming_job registers the sink worker with the sink manager; a registration failure is wrapped with context 'register v3 sink worker'. This usually means the sink manager refused or failed to register the sink (e.g. connector coordinator unavailable or invalid config).","triggerScenarios":"Creating a v3 sink where meta's register_sink call into the sink manager returns an error — inner cause is preserved in the chained error message.","commonSituations":"Iceberg sink with misconfigured catalog/warehouse; sink coordinator/worker service not running; network issues between meta and sink workers; unsupported sink properties for v3.","solutions":["Read the chained inner error for the root cause (e.g. invalid Iceberg catalog config) and fix the sink properties.","Ensure the sink manager / coordinator service is healthy and reachable from meta.","Retry creation after fixing connector configuration or infrastructure issues."],"exampleFix":"// before\nCREATE SINK s FROM mv WITH (connector='iceberg', warehouse='bad-path');\n// after\nCREATE SINK s FROM mv WITH (\n  connector='iceberg',\n  warehouse='s3://bucket/warehouse',\n  catalog='glue', ...);","handlingStrategy":"try-catch","validationCode":"// validate iceberg/sink properties before creating\nlet props = sink.get_properties();\nif props.get(\"warehouse\").map_or(true, |w| w.is_empty()) {\n    return Err(\"warehouse must be set for iceberg v3 sink\".into());\n}","typeGuard":null,"tryCatchPattern":"match create_streaming_job(..).await {\n    Err(e) if e.to_string().contains(\"register v3 sink worker\") => {\n        // inspect chained root cause, fix sink config/infrastructure, retry\n    }\n    other => other?,\n}","preventionTips":["Validate sink connector properties (catalog, warehouse, endpoints) before DDL","Monitor sink manager / coordinator health","Keep the full error chain — the inner error names the real cause"],"tags":["sink","registration","connector"],"backgroundTag":"api-request-failed","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"}