Skip to main content

laminar_connectors/lakehouse/iceberg_reference/
mod.rs

1//! Iceberg startup snapshot source for reference tables.
2
3use arrow_array::RecordBatch;
4use arrow_schema::SchemaRef;
5use async_trait::async_trait;
6use futures_util::StreamExt;
7use iceberg::scan::ArrowRecordBatchStream;
8use iceberg::spec::SnapshotRef;
9use tracing::info;
10
11use crate::config::ConnectorConfig;
12use crate::error::ConnectorError;
13use crate::reference::ReferenceTableSource;
14
15use super::iceberg_config::IcebergReadMode;
16use super::iceberg_config::IcebergSourceConfig;
17use super::iceberg_scan::{
18    connector_scan_error, plan_files, preflight_snapshot, ManifestReadLimits,
19};
20use super::snapshot_schema::{conform_snapshot_batch, validate_snapshot_schema};
21
22#[derive(Debug, Clone, Copy, PartialEq, Eq)]
23enum Phase {
24    Ready,
25    Draining,
26    Done,
27    Failed,
28    Closed,
29}
30
31/// A finite snapshot of one Iceberg table.
32pub struct IcebergReferenceTableSource {
33    config: IcebergSourceConfig,
34    declared_schema: SchemaRef,
35    phase: Phase,
36    snapshot_stream: Option<ArrowRecordBatchStream>,
37    snapshot_id: Option<i64>,
38    emitted_rows: u64,
39}
40
41impl IcebergReferenceTableSource {
42    /// Creates a source from parsed configuration and the declared table schema.
43    ///
44    /// # Errors
45    ///
46    /// Returns an error when an explicit projection conflicts with the declared schema.
47    pub fn new(
48        config: IcebergSourceConfig,
49        declared_schema: SchemaRef,
50    ) -> Result<Self, ConnectorError> {
51        config.validate_read_limits()?;
52        if config.read_mode != IcebergReadMode::Snapshot {
53            return Err(ConnectorError::ConfigurationError(
54                "Iceberg reference tables require read.mode=snapshot".into(),
55            ));
56        }
57        let declared_columns = declared_schema
58            .fields()
59            .iter()
60            .map(|field| field.name().clone())
61            .collect::<Vec<_>>();
62        if !config.select_columns.is_empty() && config.select_columns != declared_columns {
63            return Err(ConnectorError::ConfigurationError(
64                "Iceberg select.columns must exactly match the declared reference-table columns"
65                    .into(),
66            ));
67        }
68        let config = {
69            let mut config = config;
70            config.select_columns = declared_columns;
71            config
72        };
73        Ok(Self {
74            config,
75            declared_schema,
76            phase: Phase::Ready,
77            snapshot_stream: None,
78            snapshot_id: None,
79            emitted_rows: 0,
80        })
81    }
82
83    /// Creates a source from SQL connector properties and the declared table schema.
84    ///
85    /// # Errors
86    ///
87    /// Returns an error when the Iceberg configuration or projection is invalid.
88    pub fn from_connector_config(
89        config: &ConnectorConfig,
90        declared_schema: SchemaRef,
91    ) -> Result<Self, ConnectorError> {
92        Self::new(IcebergSourceConfig::from_config(config)?, declared_schema)
93    }
94
95    async fn load_initial_snapshot(&mut self) -> Result<(), ConnectorError> {
96        let catalog =
97            super::iceberg_io::build_catalog(&self.config.catalog, &self.config.storage).await?;
98        let table = super::iceberg_io::load_table_with_timeout(
99            catalog.as_ref(),
100            &self.config.catalog.namespace,
101            &self.config.catalog.table_name,
102            self.config.catalog.request_timeout,
103        )
104        .await?;
105
106        let snapshot = selected_snapshot(&table, &self.config)?;
107        let snapshot_schema = match &snapshot {
108            Some(snapshot) => snapshot
109                .schema(table.metadata())
110                .map_err(|error| connector_scan_error("resolve Iceberg snapshot schema", &error))?,
111            None => table.current_schema_ref(),
112        };
113        let predicate = super::iceberg_scan::parse_and_bind_filter(
114            self.config.filter.as_deref(),
115            snapshot_schema.clone(),
116        )?;
117        let physical_schema = iceberg::arrow::schema_to_arrow_schema(&snapshot_schema)
118            .map_err(|error| connector_scan_error("convert Iceberg snapshot schema", &error))?;
119        let projected_fields = self
120            .declared_schema
121            .fields()
122            .iter()
123            .map(|declared| {
124                physical_schema
125                    .index_of(declared.name())
126                    .map(|index| physical_schema.field(index).clone())
127                    .map_err(|_| {
128                        ConnectorError::ReadError(format!(
129                            "Iceberg snapshot is missing declared column '{}'",
130                            declared.name()
131                        ))
132                    })
133            })
134            .collect::<Result<Vec<_>, _>>()?;
135        let projected_schema = arrow_schema::Schema::new(projected_fields);
136        validate_snapshot_schema(&projected_schema, self.declared_schema.as_ref())?;
137        let Some(snapshot) = snapshot else {
138            return Ok(());
139        };
140
141        let mut builder = table
142            .scan()
143            .snapshot_id(snapshot.snapshot_id())
144            .with_batch_size(Some(8_192))
145            .with_concurrency_limit(self.config.scan_concurrency)
146            .select(self.config.select_columns.iter().map(String::as_str));
147        if let Some(predicate) = predicate {
148            builder = builder.with_filter(predicate);
149        }
150        let scan = builder
151            .build()
152            .map_err(|error| connector_scan_error("build Iceberg reference scan", &error))?;
153        let deadline = super::iceberg_io::checked_deadline(
154            self.config.storage.request_timeout,
155            "storage.request_timeout",
156        )?;
157        preflight_snapshot(
158            &table,
159            &snapshot,
160            ManifestReadLimits::from_source(&self.config),
161            deadline,
162        )
163        .await?;
164        let tasks = plan_files(&scan, self.config.max_planned_files, deadline).await?;
165        let reader = table
166            .reader_builder()
167            .with_batch_size(8_192)
168            .with_data_file_concurrency_limit(self.config.scan_concurrency)
169            .build()
170            .read(tasks)
171            .map_err(|error| connector_scan_error("create Iceberg reference reader", &error))?;
172        self.snapshot_id = Some(snapshot.snapshot_id());
173        self.snapshot_stream = Some(reader.stream());
174        Ok(())
175    }
176
177    async fn next_batch(&mut self) -> Result<Option<RecordBatch>, ConnectorError> {
178        loop {
179            let Some(stream) = self.snapshot_stream.as_mut() else {
180                return Ok(None);
181            };
182            let next = tokio::time::timeout(self.config.storage.request_timeout, stream.next())
183                .await
184                .map_err(|_| {
185                    ConnectorError::ReadError(
186                        "[LDB-ICEBERG-STORAGE-TIMEOUT] reference snapshot read made no progress"
187                            .into(),
188                    )
189                })?;
190            let Some(result) = next else {
191                self.snapshot_stream = None;
192                return Ok(None);
193            };
194            let batch = result.map_err(|error| {
195                connector_scan_error("Iceberg reference snapshot read failed", &error)
196            })?;
197            if batch.num_rows() == 0 {
198                continue;
199            }
200            let batch = conform_snapshot_batch(&batch, &self.declared_schema)?;
201            self.emitted_rows = self
202                .emitted_rows
203                .saturating_add(u64::try_from(batch.num_rows()).unwrap_or(u64::MAX));
204            return Ok(Some(batch));
205        }
206    }
207}
208
209fn selected_snapshot(
210    table: &iceberg::table::Table,
211    config: &IcebergSourceConfig,
212) -> Result<Option<SnapshotRef>, ConnectorError> {
213    if let Some(snapshot_id) = config.snapshot_id {
214        return table
215            .metadata()
216            .snapshot_by_id(snapshot_id)
217            .cloned()
218            .map(Some)
219            .ok_or_else(|| {
220                ConnectorError::ReadError(format!(
221                    "[LDB-ICEBERG-SNAPSHOT-MISSING] snapshot {snapshot_id} does not exist"
222                ))
223            });
224    }
225    if let Some(snapshot) = table.metadata().snapshot_for_ref(&config.table_ref) {
226        return Ok(Some(snapshot.clone()));
227    }
228    if config.table_ref == "main" && table.metadata().current_snapshot().is_none() {
229        return Ok(None);
230    }
231    Err(ConnectorError::ReadError(format!(
232        "[LDB-ICEBERG-REF-MISSING] table ref '{}' does not exist",
233        config.table_ref
234    )))
235}
236
237#[async_trait]
238impl ReferenceTableSource for IcebergReferenceTableSource {
239    async fn poll_snapshot(&mut self) -> Result<Option<RecordBatch>, ConnectorError> {
240        match self.phase {
241            Phase::Closed => {
242                return Err(ConnectorError::InvalidState {
243                    expected: "open reference snapshot source".into(),
244                    actual: "closed".into(),
245                });
246            }
247            Phase::Done => return Ok(None),
248            Phase::Failed => {
249                return Err(ConnectorError::InvalidState {
250                    expected: "readable reference snapshot source".into(),
251                    actual: "failed".into(),
252                });
253            }
254            Phase::Ready => {
255                if let Err(error) = self.load_initial_snapshot().await {
256                    self.phase = Phase::Failed;
257                    return Err(error);
258                }
259                self.phase = Phase::Draining;
260            }
261            Phase::Draining => {}
262        }
263
264        match self.next_batch().await {
265            Ok(Some(batch)) => Ok(Some(batch)),
266            Ok(None) => {
267                self.phase = Phase::Done;
268                info!(
269                    snapshot = ?self.snapshot_id,
270                    rows = self.emitted_rows,
271                    "Iceberg reference snapshot completed"
272                );
273                Ok(None)
274            }
275            Err(error) => {
276                self.snapshot_stream = None;
277                self.phase = Phase::Failed;
278                Err(error)
279            }
280        }
281    }
282
283    async fn close(&mut self) -> Result<(), ConnectorError> {
284        self.phase = Phase::Closed;
285        self.snapshot_stream = None;
286        Ok(())
287    }
288}
289
290#[cfg(test)]
291mod tests;