risingwavelabs/risingwave · error

max.num.messages expect usize

Error message

max.num.messages expect usize

What it means

In `KafkaSplitReader::new` (src/connector/src/source/kafka/source/reader.rs:186), the developer-only tuning option `max.num.messages` must parse as a `usize`. Like `bytes.per.second`, it is used only for performance testing, so invalid values deliberately panic with `.expect("max.num.messages expect usize")` rather than producing a normal connector error.

Solutions

  1. Set `max.num.messages` to a plain integer, e.g. `max.num.messages = '1000'`.
  2. Remove the `max.num.messages` option if not needed (defaults to usize::MAX).
  3. Use scientific-free decimal notation (no `1e6`, no suffixes).

Example fix

// before (WITH option)
max.num.messages = '1e6'
// after
max.num.messages = '1000000'
Defensive patterns

Strategy: validation

Validate before calling

// Pre-validate the WITH option before CREATE SOURCE
let max_msgs: Option<usize> = props.max_num_messages.as_ref().map(|s| {
    s.parse().expect("max.num.messages must be a plain integer")
});

Try / catch

// This path panics (expect), not returns Result; guard at construction
let max_msgs = props.max_num_messages.as_deref()
    .map(str::parse::<usize>)
    .transpose()
    .map_err(|_| anyhow!("max.num.messages expect usize"))?;

Prevention

When it happens

Trigger: Creating a Kafka source with `max.num.messages` set to a non-integer string (e.g. `'100k'`, `'ten'`) or an integer exceeding usize range, while the option is Some.

Common situations: Developers throttling/testing Kafka consumption who pass abbreviated values like `'1000 messages'` or `'1e6'` instead of a plain decimal integer.

Understand the failure class

Background: "Invalid value" and "allowed values are" config errors: what your library rejected and how to fix it — this error's family across 41 libraries.

Related errors


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

Appendix: source

Thrown at src/connector/src/source/kafka/source/reader.rs:186

            "backfill_info: {:?}",
            backfill_info
        );

        consumer.assign(&tpl)?;

        // The two parameters below are only used by developers for performance testing purposes,
        // so we panic here on purpose if the input is not correctly recognized.
        let bytes_per_second = match properties.bytes_per_second {
            None => usize::MAX,
            Some(number) => number
                .parse::<usize>()
                .expect("bytes.per.second expect usize"),
        };
        let max_num_messages = match properties.max_num_messages {
            None => usize::MAX,
            Some(number) => number
                .parse::<usize>()
                .expect("max.num.messages expect usize"),
        };

        Ok(Self {
            consumer,
            offsets,
            splits,
            backfill_info,
            known_eof_offsets,
            bytes_per_second,
            sync_call_timeout: properties.common.sync_call_timeout,
            max_num_messages,
            parser_config,
            source_ctx,
        })
    }

    fn into_stream(self) -> BoxSourceChunkStream {
        let parser_config = self.parser_config.clone();

View on GitHub (pinned to 6469eb736d)