Skip to main content

laminar_connectors/lakehouse/delta_source/
mod.rs

1//! Delta Lake source connector.
2
3use std::collections::VecDeque;
4use std::sync::Arc;
5
6use arrow_array::RecordBatch;
7use arrow_schema::SchemaRef;
8use async_trait::async_trait;
9#[cfg(feature = "delta-lake")]
10use std::time::Instant;
11#[cfg(feature = "delta-lake")]
12use tracing::debug;
13use tracing::info;
14#[cfg(feature = "delta-lake")]
15use tracing::warn;
16
17#[cfg(feature = "delta-lake")]
18use deltalake::DeltaTable;
19
20use crate::checkpoint::SourceCheckpoint;
21use crate::config::{ConnectorConfig, ConnectorState};
22use crate::connector::{
23    SourceBatch, SourceConnector, SourceConsistency, SourceContract, SourceInputMode,
24    SourceTopology,
25};
26use crate::connector::{SourcePosition, SourceStart};
27use crate::error::ConnectorError;
28
29use super::delta_source_config::DeltaSourceConfig;
30
31#[cfg(feature = "delta-lake")]
32const MAX_CDF_COMMIT_ROWS: usize = 262_144;
33#[cfg(feature = "delta-lake")]
34const MAX_CDF_COMMIT_BYTES: usize = 64 * 1024 * 1024;
35
36/// Delta Lake source connector.
37///
38/// Reads incremental Change Data Feed commits from Delta Lake tables.
39///
40/// # Lifecycle
41///
42/// ```text
43/// new() -> start() -> [poll_batch()]* -> close()
44///                           |
45///                      checkpoint()
46/// ```
47pub struct DeltaSource {
48    /// Source configuration.
49    config: DeltaSourceConfig,
50    /// Connector lifecycle state.
51    state: ConnectorState,
52    /// Arrow schema (set from table metadata on start).
53    schema: Option<SchemaRef>,
54    /// Current Delta Lake version cursor — the last *fully consumed* version.
55    /// Only advanced after all buffered batches for a version are drained.
56    current_version: i64,
57    /// The version currently being drained. While `pending_batches` is
58    /// non-empty this holds the version they came from. Once drained,
59    /// `current_version` is advanced to this value and the field is cleared.
60    #[cfg(feature = "delta-lake")]
61    inflight_version: Option<i64>,
62    /// The latest version known at the table. Used in incremental mode to
63    /// walk versions one-by-one without re-calling `get_latest_version` for
64    /// each step.
65    #[cfg(feature = "delta-lake")]
66    known_latest_version: i64,
67    /// Buffered batches from the last version load.
68    pending_batches: VecDeque<RecordBatch>,
69    /// Total records read so far.
70    records_read: u64,
71    /// Delta Lake table handle.
72    #[cfg(feature = "delta-lake")]
73    table: Option<DeltaTable>,
74    /// Catalog-resolved location used by initial open and subsequent reopens.
75    #[cfg(feature = "delta-lake")]
76    resolved_table_path: String,
77    /// Explicit options paired with the catalog-resolved location.
78    #[cfg(feature = "delta-lake")]
79    stable_storage_options: std::collections::HashMap<String, String>,
80    /// Last time we checked for new Delta versions. Used to throttle
81    /// `get_latest_version()` calls to `poll_interval` instead of
82    /// hammering every source-adapter tick (10ms).
83    #[cfg(feature = "delta-lake")]
84    last_version_check: Option<Instant>,
85}
86
87#[cfg(feature = "delta-lake")]
88fn initial_current_version(starting_version: Option<i64>, latest_version: i64) -> i64 {
89    starting_version.map_or(latest_version, |first_version| first_version - 1)
90}
91
92#[cfg(feature = "delta-lake")]
93fn cdf_output_matches(expected: &SchemaRef, batch: &RecordBatch) -> bool {
94    use laminar_core::changelog::WEIGHT_COLUMN;
95
96    let actual = batch.schema();
97    actual.fields().len() == expected.fields().len() + 1
98        && actual
99            .fields()
100            .iter()
101            .take(expected.fields().len())
102            .eq(expected.fields().iter())
103        && actual
104            .fields()
105            .last()
106            .is_some_and(|field| field.name() == WEIGHT_COLUMN)
107}
108
109impl DeltaSource {
110    /// Creates a new Delta Lake source with the given configuration.
111    #[must_use]
112    pub fn new(config: DeltaSourceConfig, _registry: Option<&prometheus::Registry>) -> Self {
113        Self {
114            config,
115            state: ConnectorState::Created,
116            schema: None,
117            current_version: -1,
118            #[cfg(feature = "delta-lake")]
119            inflight_version: None,
120            #[cfg(feature = "delta-lake")]
121            known_latest_version: -1,
122            pending_batches: VecDeque::new(),
123            records_read: 0,
124            #[cfg(feature = "delta-lake")]
125            table: None,
126            #[cfg(feature = "delta-lake")]
127            resolved_table_path: String::new(),
128            #[cfg(feature = "delta-lake")]
129            stable_storage_options: std::collections::HashMap::new(),
130            #[cfg(feature = "delta-lake")]
131            last_version_check: None,
132        }
133    }
134
135    /// Returns the current connector state.
136    #[must_use]
137    pub fn state(&self) -> ConnectorState {
138        self.state
139    }
140
141    /// Returns the current Delta Lake version cursor.
142    #[must_use]
143    pub fn current_version(&self) -> i64 {
144        self.current_version
145    }
146
147    /// Returns the source configuration.
148    #[must_use]
149    pub fn config(&self) -> &DeltaSourceConfig {
150        &self.config
151    }
152
153    /// Re-opens the Delta Lake table (e.g., after a connection failure).
154    #[cfg(feature = "delta-lake")]
155    async fn reopen_table(&mut self) -> Result<(), ConnectorError> {
156        use super::delta_io;
157
158        if self.resolved_table_path.is_empty() {
159            return Err(ConnectorError::InvalidState {
160                expected: "catalog-resolved table location".into(),
161                actual: "table location not resolved".into(),
162            });
163        }
164        let storage_options = crate::storage::StorageCredentialResolver::resolve(
165            &self.resolved_table_path,
166            &self.stable_storage_options,
167        )
168        .options;
169        let table =
170            delta_io::open_or_create_table(&self.resolved_table_path, storage_options, None)
171                .await?;
172
173        self.table = Some(table);
174        Ok(())
175    }
176}
177
178#[async_trait]
179#[allow(clippy::too_many_lines)]
180impl SourceConnector for DeltaSource {
181    fn contract(&self, config: &ConnectorConfig) -> Result<SourceContract, ConnectorError> {
182        if config.properties().is_empty() {
183            self.config.validate()?;
184        } else {
185            DeltaSourceConfig::from_config(config)?.validate()?;
186        }
187
188        Ok(SourceContract::new(
189            SourceConsistency::Ephemeral,
190            SourceTopology::Singleton,
191            SourceInputMode::FullChangelog,
192        ))
193    }
194
195    async fn start(&mut self, request: SourceStart) -> Result<(), ConnectorError> {
196        let (config, position, _) = request.into_parts();
197        if let SourcePosition::Resume { attempt, .. } = position {
198            return Err(ConnectorError::ConfigurationError(format!(
199                "Delta Lake is an ephemeral source and cannot resume checkpoint attempt {attempt:?}"
200            )));
201        }
202        let config = &config;
203
204        if !config.properties().is_empty() {
205            self.config = DeltaSourceConfig::from_config(config)?;
206        }
207        self.config.validate()?;
208        self.state = ConnectorState::Initializing;
209
210        #[cfg(feature = "delta-lake")]
211        {
212            use super::delta_io;
213
214            let stable_options = self.config.stable_storage_options();
215            let (resolved_path, stable_options) = delta_io::resolve_catalog_options(
216                &self.config.catalog_type,
217                self.config.catalog_database.as_deref(),
218                self.config.catalog_name.as_deref(),
219                self.config.catalog_schema.as_deref(),
220                &self.config.table_path,
221                &stable_options,
222            )
223            .await?;
224            let resolved_storage =
225                crate::storage::StorageCredentialResolver::resolve(&resolved_path, &stable_options);
226            info!(
227                operation_class = "connector-open",
228                storage_provider = %resolved_storage.provider,
229                storage_endpoint_class = %resolved_storage.endpoint_class(),
230                storage_auth_source = %resolved_storage.auth_source,
231                starting_version = ?self.config.starting_version,
232                "opening Delta Lake source connector"
233            );
234            let resolved_options = resolved_storage.options;
235
236            let table =
237                delta_io::open_or_create_table(&resolved_path, resolved_options, None).await?;
238
239            self.schema = Some(delta_io::get_table_schema(&table)?);
240            let table_version = table.version().ok_or_else(|| {
241                ConnectorError::ReadError("opened Delta table has no committed version".into())
242            })?;
243            let table_version = delta_io::from_delta_version(table_version)?;
244            self.current_version =
245                initial_current_version(self.config.starting_version, table_version);
246            self.known_latest_version = table_version;
247            self.last_version_check = Some(Instant::now());
248            self.resolved_table_path = resolved_path;
249            self.stable_storage_options = stable_options;
250
251            info!(
252                table_version,
253                current_version = self.current_version,
254                "Delta Lake source: resolved starting version"
255            );
256
257            self.table = Some(table);
258        }
259
260        #[cfg(not(feature = "delta-lake"))]
261        {
262            self.state = ConnectorState::Failed;
263            return Err(ConnectorError::ConfigurationError(
264                "Delta Lake source requires the 'delta-lake' feature to be enabled. \
265                 Build with: cargo build --features delta-lake"
266                    .into(),
267            ));
268        }
269
270        #[cfg(feature = "delta-lake")]
271        {
272            self.state = ConnectorState::Running;
273            info!("Delta Lake source connector opened successfully");
274            Ok(())
275        }
276    }
277
278    #[allow(unused_variables)]
279    async fn poll_batch(
280        &mut self,
281        max_records: usize,
282    ) -> Result<Option<SourceBatch>, ConnectorError> {
283        if self.state != ConnectorState::Running {
284            return Err(ConnectorError::InvalidState {
285                expected: "Running".into(),
286                actual: self.state.to_string(),
287            });
288        }
289
290        // Return buffered batches first. When the buffer drains
291        // completely, advance current_version to the inflight version
292        // so that checkpoint() reports the fully-consumed position.
293        if let Some(batch) = self.pending_batches.pop_front() {
294            self.records_read += batch.num_rows() as u64;
295
296            #[cfg(feature = "delta-lake")]
297            if self.pending_batches.is_empty() {
298                if let Some(v) = self.inflight_version.take() {
299                    self.current_version = v;
300                }
301            }
302
303            return Ok(Some(SourceBatch::new(batch)));
304        }
305
306        // Check for new versions, throttled by poll_interval.
307        #[cfg(feature = "delta-lake")]
308        {
309            use super::delta_io;
310
311            // Recover from lost table handle (e.g., connection failure).
312            if self.table.is_none() {
313                match self.reopen_table().await {
314                    Ok(()) => {
315                        info!("Delta Lake source: re-opened table after lost handle");
316                    }
317                    Err(e) => {
318                        warn!(error = %e, "Delta Lake source: reopen failed, will retry");
319                        return Ok(None);
320                    }
321                }
322            }
323
324            // Throttle version checks: skip if less than poll_interval has
325            // elapsed since the last check. This prevents hammering
326            // get_latest_version() on every source-adapter tick (10ms).
327            // In incremental mode, skip the throttle if we already know
328            // there are more versions to process (catch-up).
329            let needs_refresh = self.known_latest_version <= self.current_version;
330            if needs_refresh {
331                if let Some(last_check) = self.last_version_check {
332                    if last_check.elapsed() < self.config.poll_interval {
333                        return Ok(None);
334                    }
335                }
336                self.last_version_check = Some(Instant::now());
337
338                let table = self
339                    .table
340                    .as_mut()
341                    .ok_or_else(|| ConnectorError::InvalidState {
342                        expected: "table initialized".into(),
343                        actual: "table not initialized".into(),
344                    })?;
345                let latest_version = match delta_io::get_latest_version(table).await {
346                    Ok(v) => v,
347                    Err(e) => {
348                        warn!(error = %e, "Delta Lake source: version check failed, will retry");
349                        return Ok(None);
350                    }
351                };
352                self.known_latest_version = latest_version;
353
354                if latest_version <= self.current_version {
355                    return Ok(None); // No new data
356                }
357
358                debug!(
359                    current_version = self.current_version,
360                    latest_version, "Delta Lake source: new version(s) available"
361                );
362            }
363
364            let target_version = self.current_version.checked_add(1).ok_or_else(|| {
365                ConnectorError::ConfigurationError("Delta source version cursor overflowed".into())
366            })?;
367
368            let table = self
369                .table
370                .as_ref()
371                .ok_or_else(|| ConnectorError::InvalidState {
372                    expected: "table initialized".into(),
373                    actual: "table not initialized".into(),
374                })?;
375            let log_store = table.log_store();
376            let delta_version = delta_io::to_delta_version(target_version)?;
377            match log_store.read_commit_entry(delta_version).await {
378                Ok(Some(_)) => {}
379                Ok(None) => {
380                    return Err(ConnectorError::ConfigurationError(format!(
381                        "Delta commit {target_version} is unavailable; incremental streaming cannot fall back to a snapshot"
382                    )));
383                }
384                Err(error) => {
385                    return Err(ConnectorError::ReadError(format!(
386                        "failed to verify Delta commit {target_version}: {error}"
387                    )));
388                }
389            }
390
391            let scan_table = table.clone();
392            let cdf_batches = delta_io::read_cdf_batches(
393                scan_table,
394                target_version,
395                target_version,
396                MAX_CDF_COMMIT_ROWS,
397                MAX_CDF_COMMIT_BYTES,
398            )
399            .await?;
400
401            let expected_schema =
402                self.schema
403                    .as_ref()
404                    .ok_or_else(|| ConnectorError::InvalidState {
405                        expected: "source schema initialized".into(),
406                        actual: "source schema missing".into(),
407                    })?;
408            let mut batches = Vec::with_capacity(cdf_batches.len());
409            for batch in cdf_batches {
410                let mapped = delta_io::map_cdf_to_changelog(&batch)?;
411                if !cdf_output_matches(expected_schema, &mapped) {
412                    return Err(ConnectorError::SchemaMismatch(format!(
413                        "Delta CDF schema evolved at version {target_version}"
414                    )));
415                }
416                batches.push(mapped);
417            }
418
419            // Buffer all batches. Do NOT advance current_version yet —
420            // it is only safe to checkpoint this version after the
421            // buffer is fully drained. Store it as inflight_version.
422            for batch in batches {
423                if batch.num_rows() == 0 {
424                    continue;
425                }
426                self.pending_batches.push_back(batch);
427            }
428
429            if self.pending_batches.is_empty() {
430                // Version fully consumed with no data rows (metadata-only).
431                self.current_version = target_version;
432            } else {
433                // Version fully consumed, batches buffered. Advance after drain.
434                self.inflight_version = Some(target_version);
435            }
436
437            if let Some(batch) = self.pending_batches.pop_front() {
438                self.records_read += batch.num_rows() as u64;
439
440                // Single-batch version: buffer is already empty, advance now.
441                if self.pending_batches.is_empty() {
442                    if let Some(v) = self.inflight_version.take() {
443                        self.current_version = v;
444                    }
445                }
446
447                return Ok(Some(SourceBatch::new(batch)));
448            }
449        }
450
451        Ok(None)
452    }
453
454    fn schema(&self) -> SchemaRef {
455        self.schema
456            .clone()
457            .unwrap_or_else(|| Arc::new(arrow_schema::Schema::empty()))
458    }
459
460    fn checkpoint(&self) -> SourceCheckpoint {
461        let mut cp = SourceCheckpoint::new();
462        cp.set_offset("delta_version", self.current_version.to_string());
463        cp
464    }
465
466    async fn close(&mut self) -> Result<(), ConnectorError> {
467        info!("closing Delta Lake source connector");
468
469        #[cfg(feature = "delta-lake")]
470        {
471            self.table = None;
472        }
473
474        self.pending_batches.clear();
475        self.state = ConnectorState::Closed;
476
477        info!(
478            current_version = self.current_version,
479            records_read = self.records_read,
480            "Delta Lake source connector closed"
481        );
482
483        Ok(())
484    }
485}
486
487impl std::fmt::Debug for DeltaSource {
488    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
489        f.debug_struct("DeltaSource")
490            .field("state", &self.state)
491            .field("table_path", &"<configured>")
492            .field("current_version", &self.current_version)
493            .field("pending_batches", &self.pending_batches.len())
494            .field("records_read", &self.records_read)
495            .finish_non_exhaustive()
496    }
497}
498
499#[cfg(test)]
500mod tests;