influxdata/influxdb · error · PersisterError

datafusion error

Error message

datafusion error: {0}

What it means

PersisterError::DataFusion wraps a DataFusionError produced while executing DataFusion queries or converting RecordBatches inside the Persister. It is surfaced with `#[from]`, so any DataFusion operation performed by the persister (query planning, execution, schema/batch handling) that fails is automatically converted into this variant.

Solutions

  1. Read the inner DataFusionError message (printed after 'datafusion error:') to identify the root cause.
  2. Fix the query or RecordBatch schema that DataFusion rejected.
  3. Ensure influxdb3 and datafusion crate versions are aligned (no mixed versions in the build).
  4. Add tests around the failing persister call to reproduce with a minimal batch.

Example fix

// before: mismatched schema batch passed to persister
let batch = RecordBatch::try_new(bad_schema, cols)?;
persister.write(db, table, batch).await?;
// after: use the schema derived from the table definition
let batch = RecordBatch::try_new(table_schema.clone(), cols)?;
Defensive patterns

Strategy: try-catch

Type guard

fn is_datafusion_error(e: &PersisterError) -> bool {
    matches!(e, PersisterError::DataFusion(_))
}

Try / catch

match persister.write(db, table, batch).await {
    Err(PersisterError::DataFusion(e)) => eprintln!("datafusion failed: {e}"),
    Err(e) => return Err(e.into()),
    Ok(v) => v,
}

Prevention

When it happens

Trigger: Calling Persister methods (e.g. persist of query results, write of RecordBatches) that internally run or consume DataFusion APIs; a DataFusion plan, schema mismatch, or execution failure propagates via `?` into DataFusion(PersisterError).

Common situations: Querying data whose schema changed between snapshots; DataFusion SQL errors from bad queries issued internally; incompatible record batch schemas passed to write; DataFusion version upgrades introducing behavioral changes.

Understand the failure class

Background: Database query failed: Internal Server Error 500s wrapping SQL, Prisma, and connection failures — what to check first — this error's family across 16 libraries.

Related errors


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

Appendix: source

Thrown at influxdb3_write/src/persister.rs:50

use futures_util::stream::TryStreamExt;
use futures_util::stream::{FuturesOrdered, StreamExt};
use influxdb3_cache::parquet_cache::ParquetFileDataToCache;
use influxdb3_wal::SnapshotSequenceNumber;
use iox_time::TimeProvider;
use object_store::path::Path as ObjPath;
use object_store::{ObjectMeta, ObjectStore};
use object_store_utils::{AdaptiveGetExt, AdaptivePutExt};
use observability_deps::tracing::{debug, error, info, trace, warn};
use parking_lot::RwLock;
use parquet::arrow::ArrowWriter;
use parquet::basic::Compression;
use parquet::file::metadata::ParquetMetaData;
use parquet::file::properties::WriterProperties;
use tokio::sync::Semaphore;

#[derive(Debug, thiserror::Error)]
pub enum PersisterError {
    #[error("datafusion error: {0}")]
    DataFusion(#[from] DataFusionError),

    #[error("serde_json error: {0}")]
    SerdeJson(#[from] serde_json::Error),

    #[error("object_store error: {0}")]
    ObjectStore(#[from] object_store::Error),

    #[error("parquet error: {0}")]
    ParquetError(#[from] parquet::errors::ParquetError),

    #[error("tried to serialize a parquet file with no rows")]
    NoRows,

    #[error("parse int error: {0}")]
    ParseInt(#[from] std::num::ParseIntError),

    #[error("unexpected persister error: {0:?}")]

View on GitHub (pinned to 06200ef96b)