laminar_connectors/lakehouse/delta_source/
mod.rs1use 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
36pub struct DeltaSource {
48 config: DeltaSourceConfig,
50 state: ConnectorState,
52 schema: Option<SchemaRef>,
54 current_version: i64,
57 #[cfg(feature = "delta-lake")]
61 inflight_version: Option<i64>,
62 #[cfg(feature = "delta-lake")]
66 known_latest_version: i64,
67 pending_batches: VecDeque<RecordBatch>,
69 records_read: u64,
71 #[cfg(feature = "delta-lake")]
73 table: Option<DeltaTable>,
74 #[cfg(feature = "delta-lake")]
76 resolved_table_path: String,
77 #[cfg(feature = "delta-lake")]
79 stable_storage_options: std::collections::HashMap<String, String>,
80 #[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 #[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 #[must_use]
137 pub fn state(&self) -> ConnectorState {
138 self.state
139 }
140
141 #[must_use]
143 pub fn current_version(&self) -> i64 {
144 self.current_version
145 }
146
147 #[must_use]
149 pub fn config(&self) -> &DeltaSourceConfig {
150 &self.config
151 }
152
153 #[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 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 #[cfg(feature = "delta-lake")]
308 {
309 use super::delta_io;
310
311 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 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); }
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 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 self.current_version = target_version;
432 } else {
433 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 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;