laminar_connectors/lakehouse/iceberg_reference/
mod.rs1use 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
31pub 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 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 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;