use std::sync::Arc;
use datafusion::execution::SessionStateBuilder;
use datafusion::execution::context::SessionState;
use openlineage_client::{ClientError, LineageContext, OpenLineageClient, Transport};
use crate::config::{DataFusionConfig, OpenLineageConfig};
use crate::context::{LineageContextProvider, StaticContextProvider};
use crate::rule::OpenLineageQueryPlanner;
#[derive(Debug)]
pub struct OpenLineage;
impl OpenLineage {
pub fn builder() -> OpenLineageBuilder {
OpenLineageBuilder::default()
}
}
#[derive(Default)]
pub struct OpenLineageBuilder {
client: Option<OpenLineageClient>,
transport: Option<Arc<dyn Transport>>,
context: Option<Arc<dyn LineageContextProvider>>,
config: Option<OpenLineageConfig>,
}
impl std::fmt::Debug for OpenLineageBuilder {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("OpenLineageBuilder")
.field("has_client", &self.client.is_some())
.field("has_transport", &self.transport.is_some())
.field("has_context", &self.context.is_some())
.field("config", &self.config)
.finish()
}
}
impl OpenLineageBuilder {
pub fn from_env(mut self) -> Result<Self, ClientError> {
let config = self
.config
.take()
.unwrap_or_else(OpenLineageConfig::for_datafusion_from_env);
if self.client.is_none() && self.transport.is_none() {
self.client = Some(OpenLineageClient::from_env()?);
}
if self.context.is_none() {
let ctx = LineageContext::from_env(&config);
self.context = Some(Arc::new(StaticContextProvider(ctx)));
}
self.config = Some(config);
Ok(self)
}
pub fn client(mut self, client: OpenLineageClient) -> Self {
self.client = Some(client);
self
}
pub fn transport(mut self, transport: Arc<dyn Transport>) -> Self {
self.transport = Some(transport);
self
}
pub fn context(mut self, context: Arc<dyn LineageContextProvider>) -> Self {
self.context = Some(context);
self
}
pub fn config(mut self, config: OpenLineageConfig) -> Self {
self.config = Some(config);
self
}
pub fn instrument(self, state: SessionState) -> SessionState {
let client = self.client.unwrap_or_else(|| match self.transport {
Some(transport) => OpenLineageClient::new(transport),
None => OpenLineageClient::noop(),
});
let context = self
.context
.unwrap_or_else(|| Arc::new(StaticContextProvider::default()));
let config = self
.config
.unwrap_or_else(OpenLineageConfig::for_datafusion);
let planner = Arc::new(OpenLineageQueryPlanner::new(
client,
context,
config,
Vec::new(),
));
SessionStateBuilder::from(state)
.with_query_planner(planner)
.build()
}
}
pub fn instrument_session_state(
state: SessionState,
client: OpenLineageClient,
context: Arc<dyn LineageContextProvider>,
config: OpenLineageConfig,
) -> SessionState {
OpenLineage::builder()
.client(client)
.context(context)
.config(config)
.instrument(state)
}
pub fn instrument_session_state_simple(
state: SessionState,
client: OpenLineageClient,
config: OpenLineageConfig,
) -> SessionState {
OpenLineage::builder()
.client(client)
.config(config)
.instrument(state)
}