{"record":{"id":"5789a027d75a38de","repo":"risingwavelabs/risingwave","slug":"pulsar-error-non-connection-error","errorCode":null,"errorMessage":"{pulsar::Error (non-connection error)}","messagePattern":"\\{pulsar::Error \\(non-connection error\\)\\}","errorType":"exception","errorClass":"SinkError::Pulsar","httpStatus":null,"severity":"error","filePath":"src/connector/src/sink/pulsar.rs","lineNumber":394,"sourceCode":"                // a SendFuture holding the message receipt\n                // or error after sending is returned\n                Ok(send_future) => {\n                    self.add_future\n                        .add_future_may_await(may_delivery_future(send_future))\n                        .await?;\n                    success_flag = true;\n                    break;\n                }\n                // error upon sending\n                Err(e) => match e {\n                    pulsar::Error::Connection(_)\n                    | pulsar::Error::Producer(_)\n                    | pulsar::Error::Consumer(_) => {\n                        connection_err = Some(e);\n                        tokio::time::sleep(self.config.retry_interval).await;\n                        continue;\n                    }\n                    _ => return Err(SinkError::Pulsar(anyhow!(e))),\n                },\n            }\n        }\n\n        if !success_flag {\n            Err(SinkError::Pulsar(anyhow!(connection_err.unwrap())))\n        } else {\n            Ok(())\n        }\n    }\n\n    async fn write_inner(\n        &mut self,\n        event_key_object: Option<String>,\n        event_object: Option<Vec<u8>>,\n    ) -> Result<()> {\n        let message = Message {\n            partition_key: event_key_object,","sourceCodeStart":376,"sourceCodeEnd":412,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/sink/pulsar.rs#L376-L412","documentation":"During send_message, pulsar::Error variants that are NOT recognized connection errors (Connection, Producer, Consumer, etc.) cause an immediate return of SinkError::Pulsar wrapping the original error, without further retry. Connection-class errors are retried with retry_interval and stored in connection_err, but everything else aborts the send right away.","triggerScenarios":"A send attempt in send_message's retry loop receives a pulsar::Error that is not Connection/Producer/Consumer — e.g. Authentication (bad JWT/token), PermissionDenied (topic ACL), TopicNotFound/NotFound, InvalidConfiguration, or client-side serialization errors.","commonSituations":"Expired or wrong Pulsar auth token; sink user lacking produce permission on the topic; topic name typo (tenant/namespace/topic); Pulsar client version incompatibility with broker.","solutions":["Read the wrapped pulsar::Error text to identify the non-connection failure class.","Fix auth credentials (token/auth plugin) in the sink's service.url/config if authentication failed.","Grant the sink's role produce permission on the topic if PermissionDenied.","Correct the topic name/namespace in the sink config if NotFound.","For genuinely transient conditions, recreate the sink or upgrade to a version with broader retry classification."],"exampleFix":null,"handlingStrategy":"try-catch","validationCode":"// Verify topic exists and role can produce before creating the sink\n// pulsar-admin topics inspect <tenant>/<ns>/<topic>\n// pulsar-admin namespaces grant-permission <ns> --role rw-sink --actions produce","typeGuard":"fn is_non_retryable_pulsar_error(e: &SinkError) -> bool {\n    matches!(e, SinkError::Pulsar(err) if !format!(\"{err:#}\").contains(\"Connection\"))\n}","tryCatchPattern":"match send_result {\n    Err(SinkError::Pulsar(e)) => {\n        // Non-connection errors do not retry; inspect and fix config, then recreate sink\n        error!(\"pulsar sink fatal: {e:#}\");\n    }\n    _ => {}\n}","preventionTips":["Grant produce permission to the sink role on the topic","Use valid, non-expired auth tokens","Double-check topic/tenant/namespace spelling","Keep broker ACLs and TLS config consistent"],"tags":["rust","pulsar","sink","authentication","error-handling"],"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"}