knut-thund 0.1.7

Þund — a Rust-native, Arrow-centric streaming dataflow engine (batch + streaming) with a pluggable execution backend: native Arrow/DataFusion or lower-to-Spark-Declarative-Pipelines via Spark Connect. The 'Airflow killer' authoring+runtime for knut.
Documentation
//! Arrow-typed dataset schemas.
//!
//! The IR carries an optional, fully Arrow-typed schema for each dataset so
//! the same type information is available from authoring through lowering. We
//! wrap [`arrow_schema::Schema`] behind a small newtype rather than exposing
//! it directly so that (a) the type is `serde`-serializable for the DSL and
//! warehouse-historized defs, and (b) callers on the default (no-`native`)
//! build still get a meaningful, comparable schema description without pulling
//! the full `arrow`/`datafusion` tree.

use serde::{Deserialize, Serialize};

/// A dataset's Arrow schema, captured as `(name, arrow-type-string,
/// nullable)` triples.
///
/// We store the *display* form of each [`arrow_schema::DataType`] (e.g.
/// `"Int64"`, `"Utf8"`, `"Timestamp(Microsecond, None)"`) rather than the
/// enum so the schema round-trips through `serde`/RON and is comparable across
/// the no-`native` build. [`DatasetSchema::to_arrow`] (feature-gated) rebuilds
/// a real [`arrow_schema::Schema`] for the native backend.
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct DatasetSchema {
    /// Ordered `(field_name, arrow_type, nullable)` columns.
    pub fields: Vec<SchemaField>,
}

/// One column in a [`DatasetSchema`].
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct SchemaField {
    /// Column name.
    pub name: String,
    /// The Arrow `DataType`, rendered as its display string.
    pub arrow_type: String,
    /// Whether the column admits nulls.
    pub nullable: bool,
}

impl DatasetSchema {
    /// An empty (to-be-inferred) schema.
    pub fn new() -> Self {
        Self::default()
    }

    /// Add a column. `arrow_type` is the display form of the Arrow type, e.g.
    /// `"Int64"` or `"Utf8"`.
    pub fn field(
        mut self,
        name: impl Into<String>,
        arrow_type: impl Into<String>,
        nullable: bool,
    ) -> Self {
        self.fields.push(SchemaField {
            name: name.into(),
            arrow_type: arrow_type.into(),
            nullable,
        });
        self
    }

    /// `true` when the schema is unspecified (to be inferred by the backend).
    pub fn is_empty(&self) -> bool {
        self.fields.is_empty()
    }

    /// Render as a one-line SDP-style schema string (`col TYPE, ...`) for the
    /// Spark backend's [`knut_pipelines::Dataset::schema`] field.
    pub fn to_sdp_schema_string(&self) -> Option<String> {
        if self.fields.is_empty() {
            return None;
        }
        Some(
            self.fields
                .iter()
                .map(|f| format!("{} {}", f.name, f.arrow_type))
                .collect::<Vec<_>>()
                .join(", "),
        )
    }

    // NOTE(scaffold): `to_arrow(&self) -> arrow_schema::Schema` is provided
    // under the `native` feature; it parses each `arrow_type` display string
    // back into an `arrow_schema::DataType`. Stubbed here to keep the default
    // build free of the arrow type-parser; see `backend::native`.
}