influxdata/influxdb · error · anyhow::Error

unexpected batch schema mismatch: expected

Error message

unexpected batch schema mismatch: expected {} columns, got {}

What it means

record_batches_to_py_rows converts Arrow RecordBatches to Python rows using a column-name list derived from the first/schema fields. Before processing each batch it verifies the batch's column count matches that field list; a mismatch means the batch's schema diverged from what was expected. This is an internal invariant violation raised via anyhow's ensure! macro.

Solutions

  1. Ensure all data written to the table conforms to one schema; fix writers producing mismatched column counts.
  2. Re-check the query/table schema (SHOW/modern schema APIs) and align clients with the current schema.
  3. If batches are legitimately heterogeneous, split them per schema before conversion instead of one call.
Defensive patterns

Strategy: validation

Validate before calling

// verify all batches share the expected schema before conversion
let expected = batches[0].schema();
for b in &batches {
    assert_eq!(b.num_columns(), expected.fields().len());
}

Try / catch

match record_batches_to_py_rows(py, &batches, &field_names) {
    Ok(rows) => ..., 
    Err(e) => log::error!("schema mismatch during row conversion: {e:#}"),
}

Prevention

When it happens

Trigger: Calling query() or wal_flush_to_py() when result batches have differing column counts (e.g. schema evolved mid-query, batches from mixed schemas concatenated).

Common situations: Writing to a table with altered schema between flushes so WAL batches differ from query schema; mixed-version writers producing inconsistent batch schemas; upstream query engine bugs.

Related errors


AI-assisted analysis of influxdata/influxdb@06200ef96b (2026-09-19). Data as JSON: /api/errors/0aaf93ac248d6ced. Report an issue: GitHub.

Appendix: source

Thrown at influxdb3_py_api/src/py_conversion.rs:39

    batches: &[RecordBatch],
) -> Result<Bound<'py, PyList>, anyhow::Error> {
    // Pre-create Python strings for field/tag names once for all batches;
    // schema must be the same across batches.
    let Some(first_batch) = batches.first() else {
        return Ok(PyList::empty(py));
    };
    let field_names: Vec<Bound<'_, PyString>> = first_batch
        .schema()
        .fields()
        .iter()
        .map(|f| PyString::new(py, f.name().as_str()))
        .collect();

    let total_rows: usize = batches.iter().map(|b| b.num_rows()).sum();
    let mut rows: Vec<Py<PyAny>> = Vec::with_capacity(total_rows);

    for batch in batches {
        ensure!(
            batch.num_columns() == field_names.len(),
            "unexpected batch schema mismatch: expected {} columns, got {}",
            field_names.len(),
            batch.num_columns()
        );
        let num_rows = batch.num_rows();
        for row_idx in 0..num_rows {
            let row = PyDict::new(py);
            for (col_idx, field_name) in field_names.iter().enumerate() {
                let array = batch.column(col_idx);
                let value = extract_arrow_value_to_py(py, array, row_idx)?;
                row.set_item(field_name, value).context("set dict item")?;
            }
            rows.push(row.into());
        }
    }

    let list = PyList::new(py, rows)?;

View on GitHub (pinned to 06200ef96b)