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
- Read the inner DataFusionError message (printed after 'datafusion error:') to identify the root cause.
- Fix the query or RecordBatch schema that DataFusion rejected.
- Ensure influxdb3 and datafusion crate versions are aligned (no mixed versions in the build).
- 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
- Validate RecordBatch schemas against the table definition before persisting.
- Pin compatible datafusion/arrow versions in Cargo.toml.
- Test persistence paths with minimal RecordBatches in CI.
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)