diff --git a/beacon-query/src/lib.rs b/beacon-query/src/lib.rs index 6c6bd634..73c4978e 100644 --- a/beacon-query/src/lib.rs +++ b/beacon-query/src/lib.rs @@ -49,11 +49,17 @@ pub struct QueryBody { #[serde(default)] from: Option, sort_by: Option>, - distinct: Option>, + distinct: Option, offset: Option, limit: Option, } +#[derive(Clone, Debug, serde::Serialize, serde::Deserialize, ToSchema)] +pub struct Distinct { + pub on: Vec, +} + #[derive(Clone, Debug, serde::Serialize, serde::Deserialize, ToSchema)] #[serde(untagged)] pub enum Select { diff --git a/beacon-query/src/parser.rs b/beacon-query/src/parser.rs index 41b0bf40..1f3a58b7 100644 --- a/beacon-query/src/parser.rs +++ b/beacon-query/src/parser.rs @@ -1,17 +1,10 @@ -use std::sync::Arc; - use beacon_data_lake::DataLake; use datafusion::{ - datasource::file_format::{csv::CsvFormatFactory, format_as_file_type, FileFormat}, - logical_expr::{Analyze, LogicalPlan, LogicalPlanBuilder}, + logical_expr::LogicalPlan, prelude::{SQLOptions, SessionContext}, }; -use crate::{ - output::{Output, OutputFormat, QueryOutputFile}, - plan::ParsedPlan, - InnerQuery, QueryBody, -}; +use crate::{plan::ParsedPlan, InnerQuery, QueryBody}; use super::Query; @@ -107,12 +100,12 @@ impl Parser { let df_schema = builder.schema().clone(); let schema = df_schema.as_arrow(); if let Some(filter) = query_body.filter { - builder = builder.filter(filter.parse(&session_state, &schema)?)?; + builder = builder.filter(filter.parse(&session_state, schema)?)?; } if let Some(filters) = query_body.filters { for filter in filters { - builder = builder.filter(filter.parse(&session_state, &schema)?)?; + builder = builder.filter(filter.parse(&session_state, schema)?)?; } } @@ -120,6 +113,22 @@ impl Parser { builder = builder.sort(sort_by.iter().map(|s| s.to_expr()))?; } + if let Some(distinct) = query_body.distinct { + let on_exprs = distinct + .on + .iter() + .map(|s| s.to_expr(&session.state())) + .collect::>>()?; + + let select_exprs = distinct + .select + .iter() + .map(|s| s.to_expr(&session.state())) + .collect::>>()?; + + builder = builder.distinct_on(on_exprs, select_exprs, None)?; + } + let offset = query_body.offset.unwrap_or(0); builder = builder.limit(offset, query_body.limit)?;