risingwavelabs/risingwave · error · ConnectorError

illegal message id string

Error message

illegal message id string {}

What it means

Pulsar message IDs are expected in the '{ledger_id}:{entry_id}:{partition}:{batch_index}' format (partition and batch_index optional). parse_message_id fails when the string does not contain 2–4 colon-separated segments, which would otherwise panic or corrupt the MessageIdData.

Solutions

  1. Provide the message id as at least 'ledger:entry', e.g. '3:17' or '3:17:-1:-1'.
  2. Verify the id was copied fully from Pulsar output (ledger_id:entry_id:partition:batch_index).
  3. Strip whitespace, surrounding brackets, or trailing characters before passing the id.
  4. Note the numeric segments must parse as u64; a malformed number fails with 'illegal ledger id'/'illegal entry id' context.

Example fix

// before
parse_message_id("1043-29")?; // wrong separator
// after
parse_message_id("1043:29")?; // ledger_id:entry_id
Defensive patterns

Strategy: validation

Validate before calling

fn is_valid_message_id(id: &str) -> bool {
    let parts: Vec<&str> = id.split(':').collect();
    (2..=4).contains(&parts.len()) && parts[..2].iter().all(|p| p.parse::<u64>().is_ok())
}

Type guard

fn parse_pulsar_message_id(id: &str) -> Option<(u64, u64)> {
    let mut it = id.split(':');
    Some((it.next()?.parse().ok()?, it.next()?.parse().ok()?))
}

Try / catch

match parse_message_id(id) {
    Err(e) if e.to_string().contains("illegal message id") => eprintln!("expected ledger:entry[:partition][:batch_index]"),
    other => other,
}

Prevention

When it happens

Trigger: Calling parse_message_id with a start-offset string that has fewer than 2 or more than 4 ':'-separated parts — e.g. '12345', 'a:b:c:d:e', or an empty string.

Common situations: Typing a Pulsar message id by hand into source start-offset config; copying an id from another system; including extra separators or pasting with whitespace/suffixes.

Understand the failure class

Background: "Invalid ... format", "must be in format X", "does not look like a ..." — invalid argument format errors across CLI tools and libraries — this error's family across 17 libraries.

Related errors


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

Appendix: source

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

pub struct PulsarBrokerReader {
    #[expect(dead_code)]
    pulsar: Pulsar<TokioExecutor>,
    consumer: Consumer<Vec<u8>, TokioExecutor>,
    split: PulsarSplit,
    split_id: SplitId,
    parser_config: ParserConfig,
    source_ctx: SourceContextRef,

    // for filter out already read messages
    already_read_offset: Option<PulsarFilterOffset>,
}

// {ledger_id}:{entry_id}:{partition}:{batch_index}
fn parse_message_id(id: &str) -> ConnectorResult<MessageIdData> {
    let splits = id.split(':').collect_vec();

    if splits.len() < 2 || splits.len() > 4 {
        bail!("illegal message id string {}", id);
    }

    let ledger_id = splits[0].parse::<u64>().context("illegal ledger id")?;
    let entry_id = splits[1].parse::<u64>().context("illegal entry id")?;

    let mut message_id = MessageIdData {
        ledger_id,
        entry_id,
        partition: None,
        batch_index: None,
        ack_set: vec![],
        batch_size: None,
        first_chunk_message_id: None,
    };

    if splits.len() > 2 {
        let partition = splits[2].parse::<i32>().context("illegal partition")?;
        message_id.partition = Some(partition);

View on GitHub (pinned to 6469eb736d)