risingwavelabs/risingwave · error · SinkError::Pulsar

{pulsar::Error (non-connection error)}

Error message

{pulsar::Error (non-connection error)}

What it means

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.

Solutions

  1. Read the wrapped pulsar::Error text to identify the non-connection failure class.
  2. Fix auth credentials (token/auth plugin) in the sink's service.url/config if authentication failed.
  3. Grant the sink's role produce permission on the topic if PermissionDenied.
  4. Correct the topic name/namespace in the sink config if NotFound.
  5. For genuinely transient conditions, recreate the sink or upgrade to a version with broader retry classification.
Defensive patterns

Strategy: try-catch

Validate before calling

// Verify topic exists and role can produce before creating the sink
// pulsar-admin topics inspect <tenant>/<ns>/<topic>
// pulsar-admin namespaces grant-permission <ns> --role rw-sink --actions produce

Type guard

fn is_non_retryable_pulsar_error(e: &SinkError) -> bool {
    matches!(e, SinkError::Pulsar(err) if !format!("{err:#}").contains("Connection"))
}

Try / catch

match send_result {
    Err(SinkError::Pulsar(e)) => {
        // Non-connection errors do not retry; inspect and fix config, then recreate sink
        error!("pulsar sink fatal: {e:#}");
    }
    _ => {}
}

Prevention

When it happens

Trigger: 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.

Common situations: 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.

Related errors


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

Appendix: source

Thrown at src/connector/src/sink/pulsar.rs:394

                // a SendFuture holding the message receipt
                // or error after sending is returned
                Ok(send_future) => {
                    self.add_future
                        .add_future_may_await(may_delivery_future(send_future))
                        .await?;
                    success_flag = true;
                    break;
                }
                // error upon sending
                Err(e) => match e {
                    pulsar::Error::Connection(_)
                    | pulsar::Error::Producer(_)
                    | pulsar::Error::Consumer(_) => {
                        connection_err = Some(e);
                        tokio::time::sleep(self.config.retry_interval).await;
                        continue;
                    }
                    _ => return Err(SinkError::Pulsar(anyhow!(e))),
                },
            }
        }

        if !success_flag {
            Err(SinkError::Pulsar(anyhow!(connection_err.unwrap())))
        } else {
            Ok(())
        }
    }

    async fn write_inner(
        &mut self,
        event_key_object: Option<String>,
        event_object: Option<Vec<u8>>,
    ) -> Result<()> {
        let message = Message {
            partition_key: event_key_object,

View on GitHub (pinned to 6469eb736d)