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
- Provide the message id as at least 'ledger:entry', e.g. '3:17' or '3:17:-1:-1'.
- Verify the id was copied fully from Pulsar output (ledger_id:entry_id:partition:batch_index).
- Strip whitespace, surrounding brackets, or trailing characters before passing the id.
- 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
- Copy message ids verbatim from Pulsar output; they use ':' separators.
- Validate the id with a regex like ^\d+:\d+(:-?\d+){0,2}$ before use.
- Trim whitespace and brackets when extracting ids from logs.
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
- invalid topic name ' '; it must be in the format / /
- only support single split
- unrecognized format() type specifier
- unterminated format() type specifier
- {0}
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)