1use std::collections::HashSet;
18
19use arrow::datatypes::SchemaRef;
20use arrow::record_batch::RecordBatch;
21use async_trait::async_trait;
22use deltalake::DeltaTable;
23use deltalake::kernel::StructType;
24use deltalake::kernel::engine::arrow_conversion::TryIntoKernel;
25use deltalake::operations::create::CreateBuilder;
26use deltalake::writer::{DeltaWriter, RecordBatchWriter};
27use faucet_common_delta::convert::infer_arrow_schema;
28use faucet_core::{FaucetError, WriteMode};
29use serde_json::Value;
30use tokio::sync::Mutex;
31
32use crate::config::DeltaSinkConfig;
33
34pub struct DeltaSink {
36 config: DeltaSinkConfig,
37 state: Mutex<SinkState>,
38}
39
40struct SinkState {
43 table: Option<DeltaTable>,
46 writer: Option<RecordBatchWriter>,
48 schema: Option<SchemaRef>,
52 warned_fields: HashSet<String>,
55 pending: bool,
58}
59
60impl SinkState {
61 fn new() -> Self {
62 Self {
63 table: None,
64 writer: None,
65 schema: None,
66 warned_fields: HashSet::new(),
67 pending: false,
68 }
69 }
70}
71
72impl DeltaSink {
73 pub async fn new(config: DeltaSinkConfig) -> Result<Self, FaucetError> {
76 config
77 .validate()
78 .map_err(|e| FaucetError::Config(format!("invalid delta sink config: {e}")))?;
79 config.connection.register_handlers();
82 Ok(Self {
83 config,
84 state: Mutex::new(SinkState::new()),
85 })
86 }
87
88 async fn ensure_open(
91 &self,
92 state: &mut SinkState,
93 records: &[Value],
94 ) -> Result<(), FaucetError> {
95 if state.schema.is_none() {
96 let schema = infer_arrow_schema(records, self.config.effective_sample_size())?;
97 state.schema = Some(schema);
98 }
99 let schema = state.schema.clone().expect("schema set above");
100 self.ensure_table_writer(state, &schema).await
101 }
102
103 async fn ensure_table_writer(
109 &self,
110 state: &mut SinkState,
111 schema: &SchemaRef,
112 ) -> Result<(), FaucetError> {
113 if state.table.is_none() {
114 let table = self.open_or_create(schema).await?;
115 state.table = Some(table);
116 }
117 if state.writer.is_none() {
118 let table = state.table.as_ref().expect("table set above");
119 let writer = RecordBatchWriter::for_table(table).map_err(|e| {
120 FaucetError::Sink(format!("delta: could not build record-batch writer: {e}"))
121 })?;
122 state.writer = Some(writer);
123 }
124 Ok(())
125 }
126
127 async fn open_or_create(&self, schema: &SchemaRef) -> Result<DeltaTable, FaucetError> {
129 if let Some(table) = self.config.connection.open_optional().await? {
130 return Ok(table);
131 }
132 if !self.config.create_if_not_missing {
133 return Err(FaucetError::Sink(format!(
134 "delta: table '{}' does not exist and create_if_not_missing is false",
135 self.config.connection.redacted_uri()
136 )));
137 }
138 self.create_table(schema).await
139 }
140
141 async fn create_table(&self, schema: &SchemaRef) -> Result<DeltaTable, FaucetError> {
143 for col in &self.config.partition_by {
145 if schema.field_with_name(col).is_err() {
146 return Err(FaucetError::Sink(format!(
147 "delta: partition column '{col}' not present in the inferred record schema"
148 )));
149 }
150 }
151
152 let delta_schema: StructType = schema.as_ref().try_into_kernel().map_err(|e| {
153 FaucetError::Sink(format!(
154 "delta: could not convert Arrow schema to Delta: {e}"
155 ))
156 })?;
157
158 let mut builder = CreateBuilder::new()
159 .with_location(self.config.connection.location_string()?)
160 .with_storage_options(self.config.connection.merged_storage_options())
161 .with_columns(delta_schema.fields().cloned());
162 if !self.config.partition_by.is_empty() {
163 builder = builder.with_partition_columns(self.config.partition_by.clone());
164 }
165
166 builder.await.map_err(|e| {
167 FaucetError::Sink(format!(
168 "delta: could not create table '{}': {e}",
169 self.config.connection.redacted_uri()
170 ))
171 })
172 }
173
174 fn encode_batch(
177 &self,
178 warned_fields: &mut HashSet<String>,
179 schema: SchemaRef,
180 records: &[Value],
181 ) -> Result<RecordBatch, FaucetError> {
182 warn_on_unknown_fields(warned_fields, &schema, records);
183
184 let mut decoder = arrow_json::ReaderBuilder::new(schema.clone())
185 .build_decoder()
186 .map_err(|e| FaucetError::Sink(format!("delta: could not build json decoder: {e}")))?;
187 decoder.serialize(records).map_err(|e| {
188 FaucetError::Sink(format!("delta: record does not match table schema: {e}"))
189 })?;
190 decoder
191 .flush()
192 .map_err(|e| FaucetError::Sink(format!("delta: json decode error: {e}")))?
193 .ok_or_else(|| FaucetError::Sink("delta: json decoder produced no batch".to_string()))
194 }
195
196 async fn write_chunk(
198 &self,
199 state: &mut SinkState,
200 records: &[Value],
201 ) -> Result<usize, FaucetError> {
202 if records.is_empty() {
203 return Ok(0);
204 }
205 self.ensure_open(state, records).await?;
206 let schema = state.schema.clone().expect("schema set");
207 let batch = self.encode_batch(&mut state.warned_fields, schema, records)?;
208 let rows = batch.num_rows();
209 let writer = state.writer.as_mut().expect("writer set");
210 writer
211 .write(batch)
212 .await
213 .map_err(|e| FaucetError::Sink(format!("delta: write failed: {e}")))?;
214 state.pending = true;
215 Ok(rows)
216 }
217}
218
219#[async_trait]
220impl faucet_core::Sink for DeltaSink {
221 fn config_schema(&self) -> Value {
222 serde_json::to_value(faucet_core::schema_for!(DeltaSinkConfig))
223 .expect("schema serialization")
224 }
225
226 fn connector_name(&self) -> &'static str {
227 "delta"
228 }
229
230 fn dataset_uri(&self) -> String {
231 self.config.connection.redacted_uri()
232 }
233
234 fn supported_write_modes(&self) -> &'static [WriteMode] {
235 &[WriteMode::Append]
238 }
239
240 async fn check(
241 &self,
242 ctx: &faucet_core::check::CheckContext,
243 ) -> Result<faucet_core::check::CheckReport, FaucetError> {
244 use faucet_core::check::{CheckReport, Probe};
245 let started = std::time::Instant::now();
246 let probe =
250 match tokio::time::timeout(ctx.timeout, self.config.connection.open_optional()).await {
251 Ok(Ok(_)) => Probe::pass("table", started.elapsed()),
252 Ok(Err(e)) => Probe::fail_hint(
253 "table",
254 started.elapsed(),
255 format!("delta sink probe failed: {e}"),
256 "Verify table_uri, credentials, and object-store reachability.",
257 ),
258 Err(_) => Probe::fail_hint(
259 "table",
260 started.elapsed(),
261 format!("delta sink probe timed out after {:?}", ctx.timeout),
262 "Check object-store network reachability.",
263 ),
264 };
265 Ok(CheckReport::single(probe))
266 }
267
268 async fn write_batch(&self, records: &[Value]) -> Result<usize, FaucetError> {
269 if records.is_empty() {
270 return Ok(0);
271 }
272 let mut state = self.state.lock().await;
273 let bs = self.config.batch_size;
274 let mut total = 0;
275 if bs == 0 || records.len() <= bs {
276 total += self.write_chunk(&mut state, records).await?;
277 } else {
278 for chunk in records.chunks(bs) {
279 total += self.write_chunk(&mut state, chunk).await?;
280 }
281 }
282 Ok(total)
283 }
284
285 #[cfg(feature = "arrow")]
289 fn supports_columnar(&self) -> bool {
290 true
291 }
292
293 #[cfg(feature = "arrow")]
299 async fn write_batch_columnar(&self, batch: &RecordBatch) -> Result<usize, FaucetError> {
300 if batch.num_rows() == 0 {
301 return Ok(0);
302 }
303 let mut state = self.state.lock().await;
304 if state.schema.is_none() {
305 state.schema = Some(batch.schema());
306 }
307 let schema = state.schema.clone().expect("schema set above");
308 self.ensure_table_writer(&mut state, &schema).await?;
309 let rows = batch.num_rows();
310 let writer = state.writer.as_mut().expect("writer set");
311 writer
312 .write(batch.clone())
313 .await
314 .map_err(|e| FaucetError::Sink(format!("delta: columnar write failed: {e}")))?;
315 state.pending = true;
316 Ok(rows)
317 }
318
319 async fn flush(&self) -> Result<(), FaucetError> {
320 let mut state = self.state.lock().await;
321 if !state.pending {
322 return Ok(());
323 }
324 let mut writer = match state.writer.take() {
328 Some(w) => w,
329 None => return Ok(()),
330 };
331 let mut table = state
332 .table
333 .take()
334 .ok_or_else(|| FaucetError::Sink("delta: flush without an open table".to_string()))?;
335 let version = writer
336 .flush_and_commit(&mut table)
337 .await
338 .map_err(|e| FaucetError::Sink(format!("delta: commit failed: {e}")))?;
339 tracing::debug!(version, uri = %self.config.connection.redacted_uri(), "delta commit");
340 state.table = Some(table);
341 state.pending = false;
342 Ok(())
343 }
344}
345
346fn warn_on_unknown_fields(
349 warned_fields: &mut HashSet<String>,
350 schema: &SchemaRef,
351 records: &[Value],
352) {
353 for rec in records {
354 if let Value::Object(map) = rec {
355 for key in map.keys() {
356 if schema.field_with_name(key).is_err() && warned_fields.insert(key.clone()) {
357 tracing::warn!(
358 field = %key,
359 "delta sink: dropping field not present in the table schema"
360 );
361 }
362 }
363 }
364 }
365}
366
367#[cfg(test)]
368mod tests {
369 use super::*;
370 use faucet_core::Sink;
371
372 #[tokio::test]
373 async fn trait_metadata_methods() {
374 let sink = DeltaSink::new(DeltaSinkConfig::new("file:///tmp/delta_meta"))
375 .await
376 .unwrap();
377 assert_eq!(sink.connector_name(), "delta");
378 assert_eq!(sink.dataset_uri(), "file:///tmp/delta_meta");
379 assert_eq!(sink.supported_write_modes(), &[WriteMode::Append]);
380 assert!(sink.config_schema().is_object());
381 }
382
383 #[tokio::test]
384 async fn create_table_rejects_missing_partition_column() {
385 let dir = tempfile::tempdir().unwrap();
386 let uri = dir.path().join("p").to_string_lossy().into_owned();
387 let mut cfg = DeltaSinkConfig::new(&uri);
388 cfg.partition_by = vec!["nope".into()];
389 let sink = DeltaSink::new(cfg).await.unwrap();
390 let err = sink
391 .write_batch(&[serde_json::json!({"id": 1})])
392 .await
393 .unwrap_err();
394 assert!(err.to_string().contains("partition column"), "{err}");
395 }
396}