risingwavelabs/risingwave · error · ConnectorError

only support single split

Error message

only support single split

What it means

The Pulsar source reader supports exactly one split (topic partition) per reader instance; constructing PulsarSplitReader with more than one split is rejected. Pulsar consumers in this connector are created per topic, so multi-split readers are not implemented.

Solutions

  1. Limit the source to a single split/partition, or rely on the framework to distribute splits across parallel reader instances.
  2. If reading a partitioned topic, partition the workload so each reader gets one partition's split.
  3. File a feature request / check newer versions for multi-split Pulsar reader support.

Example fix

// before
reader.new(props, vec![split_a, split_b], ...).await?;
// after
for split in splits.into_iter() {
    reader.new(props.clone(), vec![split], ...).await?;
}
Defensive patterns

Strategy: validation

Validate before calling

if splits.len() != 1 {
    return Err(format!("pulsar reader accepts exactly 1 split, got {}", splits.len()));
}

Try / catch

match reader.new(props, splits, ...).await {
    Err(e) if e.to_string().contains("only support single split") => distribute_splits_across_readers(),
    other => other,
}

Prevention

When it happens

Trigger: Calling PulsarSplitReader::new (broker reader path) with a splits Vec containing 2 or more PulsarSplits.

Common situations: Configuring a partitioned Pulsar topic that produces multiple splits; upstream enumerator changes returning several splits; hand-constructing splits in tests or tooling.

Understand the failure class

Background: UnsupportedOperationException and "is not supported" errors: when a library deliberately refuses a call — this error's family across 30 libraries.

Related errors


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

Appendix: source

Thrown at src/connector/src/source/pulsar/source/reader.rs:95

const PULSAR_DEFAULT_SUBSCRIPTION_PREFIX: &str = "rw-consumer";

pub enum PulsarSplitReader {
    Broker(PulsarBrokerReader),
}

#[async_trait]
impl SplitReader for PulsarSplitReader {
    type Properties = PulsarProperties;
    type Split = PulsarSplit;

    async fn new(
        props: PulsarProperties,
        splits: Vec<PulsarSplit>,
        parser_config: ParserConfig,
        source_ctx: SourceContextRef,
        _columns: Option<Vec<Column>>,
    ) -> ConnectorResult<Self> {
        ensure!(splits.len() == 1, "only support single split");
        let split = splits.into_iter().next().unwrap();
        let topic = split.topic.to_string();

        tracing::debug!("creating consumer for pulsar split topic {}", topic,);

        if props.iceberg_loader_enabled.unwrap_or(false) {
            bail!("PulsarIcebergReader has already been deprecated");
        } else {
            Ok(Self::Broker(
                PulsarBrokerReader::new(props, vec![split], parser_config, source_ctx, None)
                    .await?,
            ))
        }
    }

    fn into_stream(self) -> BoxSourceChunkStream {
        match self {
            Self::Broker(reader) => {

View on GitHub (pinned to 6469eb736d)