{"record":{"id":"13987594c3fadd2c","repo":"risingwavelabs/risingwave","slug":"connector-error-0-139875","errorCode":null,"errorMessage":"Connector error: {0}","messagePattern":"Connector error: (.+?)","errorType":"error_code","errorClass":"StreamExecutorError","httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/error.rs","lineNumber":98,"sourceCode":"        #[from]\n        #[backtrace]\n        RpcError,\n    ),\n\n    #[error(\"Channel closed: {0}\")]\n    ChannelClosed(String),\n\n    #[error(transparent)]\n    ExchangeChannelClosed(\n        #[from]\n        #[backtrace]\n        ExchangeChannelClosed,\n    ),\n\n    #[error(\"Failed to align barrier: expected `{0:?}` but got `{1:?}`\")]\n    AlignBarrier(Box<Barrier>, Box<Barrier>),\n\n    #[error(\"Connector error: {0}\")]\n    ConnectorError(\n        #[source]\n        #[backtrace]\n        BoxedError,\n    ),\n\n    #[error(transparent)]\n    DmlError(\n        #[from]\n        #[backtrace]\n        DmlError,\n    ),\n\n    #[error(transparent)]\n    NotImplemented(#[from] NotImplemented),\n\n    #[error(transparent)]\n    Uncategorized(","sourceCodeStart":80,"sourceCodeEnd":116,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/error.rs#L80-L116","documentation":"The ConnectorError variant wraps a BoxedError (boxed source error) under the message 'Connector error: {0}'. It surfaces failures raised by the connector subsystem (sources/sinks: Kafka, CDC, etc.) when data is read or written by a streaming executor. Via From<ConnectorError> the connector's own error is boxed and reported as the source.","triggerScenarios":"Thrown when an executor converting ConnectorError into StreamExecutorError (src/stream/src/executor/error.rs:136) hits a connector failure — e.g. a source poll, message parsing, or sink write fails inside the stream pipeline.","commonSituations":"Kafka broker unavailability or auth failure; CDC connector losing connection to the upstream database; malformed upstream messages; connector rate limits or quota exhaustion.","solutions":["Read the boxed source error for the connector-specific cause.","Check the external system (broker, database) availability, credentials, and network reachability.","Validate the source/sink configuration (topic, bootstrap servers, auth settings).","Retry/recover — connectors typically reconnect automatically; persistent errors need config or external-system fixes."],"exampleFix":"// before\nCREATE SOURCE s (...) WITH (connector='kafka', properties.bootstrap.servers='wrong-host:9092');\n// after\nCREATE SOURCE s (...) WITH (connector='kafka', properties.bootstrap.servers='broker:9092');","handlingStrategy":"fallback","validationCode":"// validate connector config before creating the source\nfn validate_kafka_props(props: &KafkaProps) -> Result<(), String> {\n    if props.bootstrap_servers.is_empty() { return Err(\"bootstrap.servers required\".into()); }\n    if !props.brokers_reachable() { return Err(\"brokers unreachable\".into()); }\n    Ok(())\n}","typeGuard":"fn is_connector_error(e: &StreamExecutorError) -> bool {\n    e.variant_name() == \"ConnectorError\"\n}","tryCatchPattern":"if let Err(e) = poll_result {\n    if e.variant_name() == \"ConnectorError\" {\n        tracing::warn!(source = %e, \"connector failure; will retry after backoff\");\n        return Ok(()); // stay retryable instead of crashing the actor\n    }\n    return Err(e);\n}","preventionTips":["Pre-validate connector options (hosts, topics, credentials) at DDL time.","Monitor external system health (brokers, upstream DB replication slots).","Use retry-friendly connector configs for transient outages.","Keep connector crate and external system versions compatible."],"tags":["connector","source","sink","stream-executor"],"backgroundTag":"upstream-api-error","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"}