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
- 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.
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
- 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
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
- {connection_err (pulsar::Error) after retries exhausted}
- missing FORMAT ... ENCODE ...
- primary key not defined for
- {pulsar::Error}
- {pulsar::Error from delivery future}
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)