Skip to main content

OxideSession

Struct OxideSession 

Source
pub struct OxideSession { /* private fields */ }
Expand description

An OxideLake session.

Implementations§

Source§

impl OxideSession

Source

pub fn local() -> Result<OxideSession, EngineError>

Creates an embedded session: Parquet pruning on, the local object store registered, the SQL UDFs, and the placement rule targeting the detected backend.

Source

pub fn local_with_target( target: BackendKind, ) -> Result<OxideSession, EngineError>

An embedded session whose placement rule targets target instead of the detected backend. Placement is a planning decision: operators still select the real local backend at execute() time and fall back to the CPU reference, so planning for an absent GPU is safe — it is exactly what cluster executors do with the scheduler’s plans.

Source

pub fn local_with_options( options: &SessionOptions, ) -> Result<OxideSession, EngineError>

An embedded session built with options.

Source

pub async fn connect(scheduler_url: &str) -> Result<OxideSession, EngineError>

Connects to a Ballista scheduler (df://host:port). The session carries OxideLake’s plan codec so Gpu*Exec nodes survive the trip to executors, and the SQL UDFs so queries plan client-side; placement itself happens on the scheduler.

Source

pub async fn connect_with_options( scheduler_url: &str, options: &SessionOptions, ) -> Result<OxideSession, EngineError>

Connects to a Ballista scheduler with options. target is ignored: on a cluster the scheduler’s OXIDE_CLUSTER_BACKEND decides placement, so a client-side target would be a knob that quietly does nothing.

Source

pub fn mode(&self) -> &SessionMode

The execution mode.

Source

pub fn ctx(&self) -> &SessionContext

The underlying DataFusion context.

Source

pub fn telemetry(&self) -> &Arc<TelemetryHub>

The telemetry hub for this session.

Embedded sessions only. A cluster session’s operators run on executors, in other processes; this hub is created for the shape of the type and stays empty, so a snapshot of it is not “no work happened” but “the work happened somewhere else” (#33). Each worker reports into its own process-wide hub, which oxide-worker --metrics-port exposes as Prometheus text.

Source

pub async fn sql(&self, query: &str) -> Result<DataFrame, EngineError>

Plans a SQL statement into a lazily executed DataFrame.

Nothing runs here, so nothing is logged here: see Self::collect for the one-line-per-query record.

Source

pub async fn collect( &self, query: &str, ) -> Result<(Arc<Schema>, Vec<RecordBatch>), EngineError>

Runs query to completion and logs one INFO line describing it (#33): the backend it was planned for, the rows it produced, how long it took, and how many batches fell back to the CPU reference.

This is the line an operator reads to answer “is the GPU being used and how long did the query take” without attaching a dashboard. It is on collect rather than on Self::sql because a DataFrame has not run yet: a line logged at planning time could only report the plan, and the interesting half is what the plan then did. The schema comes back beside the batches because an empty result has no batch to take it from, and a caller rendering CSV still has to print the header.

Source

pub async fn register_parquet( &self, name: &str, path: &str, ) -> Result<(), EngineError>

Registers a Parquet file or directory as name.

Source

pub async fn explain(&self, query: &str) -> Result<String, EngineError>

The indented physical plan for query, with placement tags in embedded mode. In cluster mode this is the client-side plan (DistributedQueryExec); the scheduler’s plan is what carries the tags there.

Trait Implementations§

Source§

impl Debug for OxideSession

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result<(), Error>

Formats the value using the given formatter. Read more
Source§

impl OxideSessionExt for OxideSession

Source§

async fn read_parquet(&self, path: &str) -> Result<OxideFrame, EngineError>

Reads a Parquet file or directory as a frame.
Source§

async fn table(&self, name: &str) -> Result<OxideFrame, EngineError>

A frame over a registered table.
Source§

async fn sql_frame(&self, query: &str) -> Result<OxideFrame, EngineError>

Plans a SQL statement as a frame.

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

impl<T> IntoRequest<T> for T

Source§

fn into_request(self) -> Request<T>

Wrap the input message T in a tonic::Request
Source§

impl<L> LayerExt<L> for L

Source§

fn named_layer<S>(&self, service: S) -> Layered<<L as Layer<S>>::Service, S>
where L: Layer<S>,

Applies the layer to a service and wraps it in Layered.
Source§

impl<T> Pointable for T

Source§

const ALIGN: usize

The alignment of pointer.
Source§

type Init = T

The type for initializers.
Source§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
Source§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
Source§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
Source§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. Read more
Source§

impl<T> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more