Skip to main content

rbt/materializer/
stream.rs

1//! Streaming materialize: pull DF `RecordBatch` streams batch-by-batch, write, drop.
2//!
3//! Peak retained memory ≈ in-flight batch + Parquet row-group encoder + unique tracker.
4//! Never holds a full `Vec<RecordBatch>` for the model result.
5
6use crate::core::dag::OutputFormat;
7use crate::core::project::{
8    MaterializeConfig, DEFAULT_MAX_ROW_GROUP_BYTES, DEFAULT_MAX_ROW_GROUP_ROWS,
9};
10use crate::testing::{Assertion, StreamingAssertionRunner, ValidationResult};
11use anyhow::{bail, Context, Result};
12use arrow::datatypes::SchemaRef;
13use arrow::record_batch::RecordBatch;
14use datafusion::physical_plan::SendableRecordBatchStream;
15use futures::StreamExt;
16use parquet::arrow::ArrowWriter;
17use parquet::basic::Compression;
18use parquet::file::properties::WriterProperties;
19use std::fs::{self, File};
20use std::io::{BufWriter, Write};
21use std::path::{Path, PathBuf};
22
23/// Result of a successful stream materialize.
24#[derive(Debug, Clone)]
25pub struct StreamWriteStats {
26    pub rows: usize,
27    pub batches: usize,
28    pub path: PathBuf,
29    pub bytes_written: u64,
30    pub validation: ValidationResult,
31}
32
33/// Options for stream / collect writers (derived from [`MaterializeConfig`]).
34#[derive(Debug, Clone)]
35pub struct MaterializeWriteOptions {
36    pub max_row_group_rows: usize,
37    pub max_row_group_bytes: usize,
38    /// Abort on first assertion failure (default true for fail_on_error models).
39    pub fail_fast_assertions: bool,
40}
41
42impl Default for MaterializeWriteOptions {
43    fn default() -> Self {
44        Self {
45            max_row_group_rows: DEFAULT_MAX_ROW_GROUP_ROWS,
46            max_row_group_bytes: DEFAULT_MAX_ROW_GROUP_BYTES,
47            fail_fast_assertions: true,
48        }
49    }
50}
51
52impl MaterializeWriteOptions {
53    pub fn from_config(cfg: &MaterializeConfig, fail_fast_assertions: bool) -> Self {
54        Self {
55            max_row_group_rows: cfg.max_row_group_rows.max(1),
56            max_row_group_bytes: cfg.max_row_group_bytes.max(1),
57            fail_fast_assertions,
58        }
59    }
60}
61
62fn parquet_props(opts: &MaterializeWriteOptions) -> WriterProperties {
63    WriterProperties::builder()
64        .set_max_row_group_row_count(Some(opts.max_row_group_rows))
65        .set_compression(Compression::SNAPPY)
66        .build()
67}
68
69/// Staging path for atomic publish: `dir/.name.ext.rbt-partial`.
70pub fn partial_path_for(dest: &Path) -> PathBuf {
71    let parent = dest.parent().unwrap_or_else(|| Path::new("."));
72    let name = dest
73        .file_name()
74        .map(|s| s.to_string_lossy().into_owned())
75        .unwrap_or_else(|| "output".into());
76    parent.join(format!(".{name}.rbt-partial"))
77}
78
79fn remove_if_exists(path: &Path) {
80    if path.exists() {
81        let _ = if path.is_dir() {
82            fs::remove_dir_all(path)
83        } else {
84            fs::remove_file(path)
85        };
86    }
87}
88
89/// Atomically replace `dest` with `partial` (same filesystem). Cleans partial on failure.
90pub fn atomic_publish(partial: &Path, dest: &Path) -> Result<()> {
91    if let Some(parent) = dest.parent() {
92        fs::create_dir_all(parent)
93            .with_context(|| format!("E_RBT_MATERIALIZE_IO: mkdir {}", parent.display()))?;
94    }
95    // Replace existing destination (full refresh).
96    if dest.exists() {
97        if dest.is_dir() {
98            fs::remove_dir_all(dest).with_context(|| {
99                format!(
100                    "E_RBT_MATERIALIZE_IO: remove existing dir {}",
101                    dest.display()
102                )
103            })?;
104        } else {
105            fs::remove_file(dest).with_context(|| {
106                format!(
107                    "E_RBT_MATERIALIZE_IO: remove existing file {}",
108                    dest.display()
109                )
110            })?;
111        }
112    }
113    fs::rename(partial, dest).with_context(|| {
114        format!(
115            "E_RBT_MATERIALIZE_ATOMIC: rename {} → {} failed. \
116             Partial file left for inspection if rename partially failed.",
117            partial.display(),
118            dest.display()
119        )
120    })?;
121    Ok(())
122}
123
124/// Stream a DataFusion result into `destination_path` for the given format.
125///
126/// On any error, partial artifacts are deleted (previous successful dest is left intact
127/// until a successful atomic replace).
128pub async fn materialize_stream(
129    mut stream: SendableRecordBatchStream,
130    format: &OutputFormat,
131    destination_path: &Path,
132    opts: &MaterializeWriteOptions,
133    assertions: &[Assertion],
134) -> Result<StreamWriteStats> {
135    match format {
136        OutputFormat::Parquet | OutputFormat::ZeroCopyClone => {
137            write_parquet_stream(&mut stream, destination_path, opts, assertions).await
138        }
139        OutputFormat::Jsonl => {
140            write_line_stream(&mut stream, destination_path, opts, assertions, LineFormat::Jsonl)
141                .await
142        }
143        OutputFormat::Csv => {
144            write_line_stream(&mut stream, destination_path, opts, assertions, LineFormat::Csv)
145                .await
146        }
147        OutputFormat::Iceberg => {
148            write_iceberg_stream(&mut stream, destination_path, opts, assertions).await
149        }
150        OutputFormat::ParquetAndIceberg => {
151            // Dual-write: stream once into parquet, then re-read path for iceberg layout
152            // would double IO. For dual-write we buffer is bad — write parquet stream,
153            // then copy data file into iceberg layout + metadata (metadata only needs schema+rows).
154            let parquet_path =
155                if destination_path.extension().and_then(|e| e.to_str()) == Some("parquet") {
156                    destination_path.to_path_buf()
157                } else {
158                    destination_path.with_extension("parquet")
159                };
160            let stats =
161                write_parquet_stream(&mut stream, &parquet_path, opts, assertions).await?;
162            // Build iceberg sidecar from written parquet (schema + row count) without re-materializing batches.
163            write_iceberg_sidecar_from_parquet(&parquet_path, stats.rows, &stats.path)?;
164            Ok(stats)
165        }
166    }
167}
168
169/// Stream write Parquet with atomic publish + optional streaming assertions.
170pub async fn write_parquet_stream(
171    stream: &mut SendableRecordBatchStream,
172    destination_path: &Path,
173    opts: &MaterializeWriteOptions,
174    assertions: &[Assertion],
175) -> Result<StreamWriteStats> {
176    let schema = stream.schema();
177    let partial = partial_path_for(destination_path);
178    remove_if_exists(&partial);
179    if let Some(parent) = partial.parent() {
180        fs::create_dir_all(parent)?;
181    }
182
183    let mut runner = StreamingAssertionRunner::new(assertions, opts.fail_fast_assertions);
184    let props = parquet_props(opts);
185    let file = File::create(&partial).with_context(|| {
186        format!(
187            "E_RBT_MATERIALIZE_IO: create partial parquet {}",
188            partial.display()
189        )
190    })?;
191    // Large buffer reduces syscalls on multi-million-row writes.
192    let buf = BufWriter::with_capacity(8 * 1024 * 1024, file);
193    let mut writer = ArrowWriter::try_new(buf, schema.clone(), Some(props)).with_context(|| {
194        format!(
195            "E_RBT_MATERIALIZE_PARQUET: ArrowWriter::try_new for {}",
196            partial.display()
197        )
198    })?;
199
200    let mut rows = 0usize;
201    let mut batches = 0usize;
202    let result = async {
203        while let Some(item) = stream.next().await {
204            let batch = item.map_err(|e| {
205                anyhow::anyhow!("E_RBT_MATERIALIZE_STREAM: DataFusion stream error: {e}")
206            })?;
207            if batch.num_rows() == 0 && batch.num_columns() == 0 {
208                continue;
209            }
210            if !runner.is_empty() {
211                runner.observe_batch(&batch).map_err(|e| {
212                    anyhow::anyhow!("E_RBT_MATERIALIZE_ASSERT: {e}")
213                })?;
214            }
215            writer.write(&batch).with_context(|| {
216                format!(
217                    "E_RBT_MATERIALIZE_PARQUET: write batch #{batches} to {}",
218                    partial.display()
219                )
220            })?;
221            rows += batch.num_rows();
222            batches += 1;
223            // Soft flush when in-progress row group grows large.
224            let in_progress = writer.in_progress_size();
225            if in_progress >= opts.max_row_group_bytes {
226                writer.flush().with_context(|| {
227                    format!(
228                        "E_RBT_MATERIALIZE_PARQUET: flush row group at {in_progress} bytes"
229                    )
230                })?;
231            }
232            // batch dropped here
233        }
234        Ok::<(), anyhow::Error>(())
235    }
236    .await;
237
238    if let Err(e) = result {
239        let _ = writer.close();
240        remove_if_exists(&partial);
241        return Err(e);
242    }
243
244    writer.close().with_context(|| {
245        format!(
246            "E_RBT_MATERIALIZE_PARQUET: close writer {}",
247            partial.display()
248        )
249    })?;
250
251    let validation = runner.finish();
252    if validation.failed_assertions > 0 {
253        remove_if_exists(&partial);
254        bail!(
255            "E_RBT_MATERIALIZE_ASSERT: {} assertion(s) failed: {}",
256            validation.failed_assertions,
257            validation.errors.join("; ")
258        );
259    }
260
261    atomic_publish(&partial, destination_path)?;
262    let bytes_written = fs::metadata(destination_path).map(|m| m.len()).unwrap_or(0);
263
264    Ok(StreamWriteStats {
265        rows,
266        batches,
267        path: destination_path.to_path_buf(),
268        bytes_written,
269        validation,
270    })
271}
272
273#[derive(Clone, Copy)]
274enum LineFormat {
275    Jsonl,
276    Csv,
277}
278
279async fn write_line_stream(
280    stream: &mut SendableRecordBatchStream,
281    destination_path: &Path,
282    opts: &MaterializeWriteOptions,
283    assertions: &[Assertion],
284    line_fmt: LineFormat,
285) -> Result<StreamWriteStats> {
286    let partial = partial_path_for(destination_path);
287    remove_if_exists(&partial);
288    if let Some(parent) = partial.parent() {
289        fs::create_dir_all(parent)?;
290    }
291    let file = File::create(&partial).with_context(|| {
292        format!(
293            "E_RBT_MATERIALIZE_IO: create partial {}",
294            partial.display()
295        )
296    })?;
297    let mut runner = StreamingAssertionRunner::new(assertions, opts.fail_fast_assertions);
298    let mut rows = 0usize;
299    let mut batches = 0usize;
300
301    let write_result = async {
302        match line_fmt {
303            LineFormat::Jsonl => {
304                let mut writer = arrow::json::LineDelimitedWriter::new(file);
305                while let Some(item) = stream.next().await {
306                    let batch = item.map_err(|e| {
307                        anyhow::anyhow!("E_RBT_MATERIALIZE_STREAM: {e}")
308                    })?;
309                    if !runner.is_empty() {
310                        runner.observe_batch(&batch)?;
311                    }
312                    writer.write(&batch)?;
313                    rows += batch.num_rows();
314                    batches += 1;
315                }
316                writer.finish()?;
317            }
318            LineFormat::Csv => {
319                let mut writer = arrow::csv::Writer::new(file);
320                while let Some(item) = stream.next().await {
321                    let batch = item.map_err(|e| {
322                        anyhow::anyhow!("E_RBT_MATERIALIZE_STREAM: {e}")
323                    })?;
324                    if !runner.is_empty() {
325                        runner.observe_batch(&batch)?;
326                    }
327                    writer.write(&batch)?;
328                    rows += batch.num_rows();
329                    batches += 1;
330                }
331            }
332        }
333        Ok::<(), anyhow::Error>(())
334    }
335    .await;
336
337    if let Err(e) = write_result {
338        remove_if_exists(&partial);
339        return Err(e);
340    }
341
342    let validation = runner.finish();
343    if validation.failed_assertions > 0 {
344        remove_if_exists(&partial);
345        bail!(
346            "E_RBT_MATERIALIZE_ASSERT: {} assertion(s) failed: {}",
347            validation.failed_assertions,
348            validation.errors.join("; ")
349        );
350    }
351
352    atomic_publish(&partial, destination_path)?;
353    let bytes_written = fs::metadata(destination_path).map(|m| m.len()).unwrap_or(0);
354    Ok(StreamWriteStats {
355        rows,
356        batches,
357        path: destination_path.to_path_buf(),
358        bytes_written,
359        validation,
360    })
361}
362
363async fn write_iceberg_stream(
364    stream: &mut SendableRecordBatchStream,
365    table_root: &Path,
366    opts: &MaterializeWriteOptions,
367    assertions: &[Assertion],
368) -> Result<StreamWriteStats> {
369    // Data file is full-refresh; metadata versions are retained for a local snapshot log
370    // (not multi-writer OCC / REST catalog — honest FS Iceberg-style history).
371    let prior = read_iceberg_version_hint(table_root);
372    let next_version = prior.map(|v| v + 1).unwrap_or(1);
373    let mut meta_log = prior_metadata_log(table_root, prior);
374
375    let staging = table_root.with_extension("rbt-partial-table");
376    remove_if_exists(&staging);
377    let data_dir = staging.join("data");
378    let meta_dir = staging.join("metadata");
379    fs::create_dir_all(&data_dir)?;
380    fs::create_dir_all(&meta_dir)?;
381
382    // Preserve prior metadata JSON files into staging for history.
383    if let Some(old_meta) = table_root.join("metadata").exists().then(|| table_root.join("metadata"))
384    {
385        if let Ok(entries) = fs::read_dir(&old_meta) {
386            for e in entries.flatten() {
387                let p = e.path();
388                if p.extension().and_then(|x| x.to_str()) == Some("json") {
389                    if let Some(name) = p.file_name() {
390                        let _ = fs::copy(&p, meta_dir.join(name));
391                    }
392                }
393            }
394        }
395    }
396
397    let data_path = data_dir.join("part-00000.parquet");
398    let schema = stream.schema();
399    let stats = write_parquet_stream(stream, &data_path, opts, assertions).await?;
400
401    write_iceberg_metadata(
402        &staging,
403        &schema,
404        stats.rows,
405        "part-00000.parquet",
406        next_version,
407        &mut meta_log,
408    )?;
409
410    if table_root.exists() {
411        fs::remove_dir_all(table_root).with_context(|| {
412            format!(
413                "E_RBT_MATERIALIZE_IO: clear iceberg table {}",
414                table_root.display()
415            )
416        })?;
417    }
418    if let Some(parent) = table_root.parent() {
419        fs::create_dir_all(parent)?;
420    }
421    fs::rename(&staging, table_root).with_context(|| {
422        format!(
423            "E_RBT_MATERIALIZE_ATOMIC: rename iceberg staging {} → {}",
424            staging.display(),
425            table_root.display()
426        )
427    })?;
428
429    tracing::info!(
430        "Iceberg FS table written (stream): {} ({} rows, metadata v{}, data/part-00000.parquet)",
431        table_root.display(),
432        stats.rows,
433        next_version
434    );
435
436    Ok(StreamWriteStats {
437        rows: stats.rows,
438        batches: stats.batches,
439        path: table_root.to_path_buf(),
440        bytes_written: stats.bytes_written,
441        validation: stats.validation,
442    })
443}
444
445fn read_iceberg_version_hint(table_root: &Path) -> Option<u64> {
446    let hint = table_root.join("metadata/version-hint.text");
447    let s = fs::read_to_string(hint).ok()?;
448    s.trim().parse().ok()
449}
450
451fn prior_metadata_log(table_root: &Path, prior: Option<u64>) -> Vec<serde_json::Value> {
452    use serde_json::json;
453    let mut log = Vec::new();
454    if let Some(v) = prior {
455        let meta_path = table_root.join(format!("metadata/v{v}.metadata.json"));
456        if meta_path.exists() {
457            let now_ms = std::time::SystemTime::now()
458                .duration_since(std::time::UNIX_EPOCH)
459                .map(|d| d.as_millis() as u64)
460                .unwrap_or(0);
461            log.push(json!({
462                "timestamp-ms": now_ms,
463                "metadata-file": format!("v{v}.metadata.json"),
464            }));
465        }
466    }
467    log
468}
469
470fn write_iceberg_sidecar_from_parquet(
471    parquet_path: &Path,
472    row_count: usize,
473    _stats_path: &Path,
474) -> Result<()> {
475    let table_root = super::sibling_iceberg_dir(parquet_path);
476    let prior = read_iceberg_version_hint(&table_root);
477    let next = prior.map(|v| v + 1).unwrap_or(1);
478    let mut log = prior_metadata_log(&table_root, prior);
479    // Preserve prior metadata JSON into a temp list of copies.
480    let mut prior_meta_files: Vec<(String, Vec<u8>)> = Vec::new();
481    let old_meta = table_root.join("metadata");
482    if old_meta.is_dir() {
483        if let Ok(entries) = fs::read_dir(&old_meta) {
484            for e in entries.flatten() {
485                let p = e.path();
486                if p.extension().and_then(|x| x.to_str()) == Some("json") {
487                    if let (Some(name), Ok(bytes)) = (
488                        p.file_name().map(|n| n.to_string_lossy().into_owned()),
489                        fs::read(&p),
490                    ) {
491                        prior_meta_files.push((name, bytes));
492                    }
493                }
494            }
495        }
496    }
497
498    let file = File::open(parquet_path)
499        .with_context(|| format!("open {} for iceberg sidecar", parquet_path.display()))?;
500    let builder = parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder::try_new(file)
501        .with_context(|| format!("parquet reader {}", parquet_path.display()))?;
502    let schema = builder.schema().clone();
503
504    if table_root.exists() {
505        fs::remove_dir_all(&table_root)?;
506    }
507    let data_dir = table_root.join("data");
508    let meta_dir = table_root.join("metadata");
509    fs::create_dir_all(&data_dir)?;
510    fs::create_dir_all(&meta_dir)?;
511    for (name, bytes) in prior_meta_files {
512        let _ = fs::write(meta_dir.join(name), bytes);
513    }
514    let data_name = "part-00000.parquet";
515    fs::copy(parquet_path, data_dir.join(data_name))?;
516    write_iceberg_metadata(
517        &table_root,
518        &schema,
519        row_count,
520        data_name,
521        next,
522        &mut log,
523    )?;
524    Ok(())
525}
526
527fn write_iceberg_metadata(
528    table_root: &Path,
529    schema: &SchemaRef,
530    total_rows: usize,
531    data_file_name: &str,
532    version: u64,
533    metadata_log: &mut Vec<serde_json::Value>,
534) -> Result<()> {
535    use serde_json::json;
536    use std::time::{SystemTime, UNIX_EPOCH};
537
538    let meta_dir = table_root.join("metadata");
539    fs::create_dir_all(&meta_dir)?;
540
541    let mut fields = Vec::new();
542    for (i, f) in schema.fields().iter().enumerate() {
543        fields.push(json!({
544            "id": i + 1,
545            "name": f.name(),
546            "required": !f.is_nullable(),
547            "type": arrow_type_to_iceberg_json(f.data_type()),
548        }));
549    }
550
551    let now_ms = SystemTime::now()
552        .duration_since(UNIX_EPOCH)
553        .map(|d| d.as_millis() as u64)
554        .unwrap_or(0);
555    let snapshot_id = now_ms.wrapping_add(version);
556    let location = table_root
557        .canonicalize()
558        .unwrap_or_else(|_| table_root.to_path_buf());
559    let location_uri = format!("file://{}", location.display());
560
561    let metadata = json!({
562        "format-version": 2,
563        "table-uuid": format!("{:032x}", snapshot_id),
564        "location": location_uri,
565        "last-sequence-number": version,
566        "last-updated-ms": now_ms,
567        "last-column-id": fields.len(),
568        "current-schema-id": 0,
569        "schemas": [{
570            "type": "struct",
571            "schema-id": 0,
572            "fields": fields,
573        }],
574        "default-spec-id": 0,
575        "partition-specs": [{ "spec-id": 0, "fields": [] }],
576        "last-partition-id": 0,
577        "default-sort-order-id": 0,
578        "sort-orders": [{ "order-id": 0, "fields": [] }],
579        "properties": {
580            "rbt.writer": "rbt",
581            "rbt.layout": "filesystem-iceberg-v1",
582            "write.format.default": "parquet",
583            "rbt.materialize": "stream",
584            "rbt.metadata-version": version.to_string()
585        },
586        "current-snapshot-id": snapshot_id,
587        "snapshots": [{
588            "snapshot-id": snapshot_id,
589            "sequence-number": version,
590            "timestamp-ms": now_ms,
591            "summary": {
592                "operation": "overwrite",
593                "rbt.added-records": total_rows.to_string(),
594                "rbt.added-data-files": "1",
595                "rbt.data-file": format!("data/{data_file_name}")
596            },
597            "schema-id": 0
598        }],
599        "snapshot-log": [{
600            "timestamp-ms": now_ms,
601            "snapshot-id": snapshot_id
602        }],
603        "metadata-log": metadata_log,
604        "rbt": {
605            "note": "Filesystem Iceberg-style table (full-refresh data, versioned metadata). Not REST/Glue OCC.",
606            "data_files": [format!("data/{data_file_name}")],
607            "row_count": total_rows,
608            "metadata_version": version
609        }
610    });
611
612    let meta_name = format!("v{version}.metadata.json");
613    let meta_path = meta_dir.join(&meta_name);
614    let mut meta_file = File::create(&meta_path)?;
615    writeln!(meta_file, "{}", serde_json::to_string_pretty(&metadata)?)?;
616    let mut hint = File::create(meta_dir.join("version-hint.text"))?;
617    writeln!(hint, "{version}")?;
618    fs::copy(&meta_path, meta_dir.join("metadata.json"))?;
619    Ok(())
620}
621
622fn arrow_type_to_iceberg_json(dt: &arrow::datatypes::DataType) -> serde_json::Value {
623    use arrow::datatypes::DataType;
624    use serde_json::json;
625    match dt {
626        DataType::Boolean => json!("boolean"),
627        DataType::Int32 => json!("int"),
628        DataType::Int64 => json!("long"),
629        DataType::Float32 => json!("float"),
630        DataType::Float64 => json!("double"),
631        DataType::Utf8 | DataType::LargeUtf8 | DataType::Utf8View => json!("string"),
632        DataType::Binary | DataType::LargeBinary => json!("binary"),
633        DataType::Date32 | DataType::Date64 => json!("date"),
634        DataType::Timestamp(_, _) => json!("timestamptz"),
635        other => json!(format!("string /* arrow:{:?} */", other)),
636    }
637}
638
639/// Collect-mode helper: write batches with same Parquet props / atomic publish as stream.
640pub fn write_parquet_batches_atomic(
641    batches: &[RecordBatch],
642    path: &Path,
643    opts: &MaterializeWriteOptions,
644) -> Result<usize> {
645    if batches.is_empty() {
646        return Ok(0);
647    }
648    let schema = batches[0].schema();
649    let partial = partial_path_for(path);
650    remove_if_exists(&partial);
651    if let Some(parent) = partial.parent() {
652        fs::create_dir_all(parent)?;
653    }
654    let file = File::create(&partial)?;
655    let buf = BufWriter::with_capacity(8 * 1024 * 1024, file);
656    let props = parquet_props(opts);
657    let mut writer = ArrowWriter::try_new(buf, schema, Some(props))?;
658    let mut rows = 0usize;
659    for batch in batches {
660        writer.write(batch)?;
661        rows += batch.num_rows();
662        if writer.in_progress_size() >= opts.max_row_group_bytes {
663            writer.flush()?;
664        }
665    }
666    writer.close()?;
667    atomic_publish(&partial, path)?;
668    Ok(rows)
669}
670
671/// Load small Parquet file into memory for optional MemTable ref() after stream write.
672pub fn load_parquet_batches(path: &Path) -> Result<Vec<RecordBatch>> {
673    let file = File::open(path)
674        .with_context(|| format!("E_RBT_REF_LOAD: open {} for MemTable", path.display()))?;
675    let builder = parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder::try_new(file)
676        .with_context(|| format!("E_RBT_REF_LOAD: parquet builder {}", path.display()))?;
677    let reader = builder
678        .build()
679        .with_context(|| format!("E_RBT_REF_LOAD: parquet reader {}", path.display()))?;
680    let mut out = Vec::new();
681    for item in reader {
682        out.push(item.with_context(|| {
683            format!("E_RBT_REF_LOAD: read batch from {}", path.display())
684        })?);
685    }
686    Ok(out)
687}
688
689/// Empty schema-only Parquet (0 rows) so ref() registration has a file.
690pub fn write_empty_parquet(schema: SchemaRef, path: &Path, opts: &MaterializeWriteOptions) -> Result<()> {
691    let partial = partial_path_for(path);
692    remove_if_exists(&partial);
693    if let Some(parent) = partial.parent() {
694        fs::create_dir_all(parent)?;
695    }
696    let file = File::create(&partial)?;
697    let props = parquet_props(opts);
698    let writer = ArrowWriter::try_new(file, schema, Some(props))?;
699    writer.close()?;
700    atomic_publish(&partial, path)?;
701    Ok(())
702}
703
704#[cfg(test)]
705mod tests {
706    use super::*;
707    use crate::testing::Assertion;
708    use arrow::datatypes::{DataType, Field, Schema};
709    use datafusion::prelude::SessionContext;
710    use std::sync::Arc;
711
712    fn sample_schema() -> SchemaRef {
713        Arc::new(Schema::new(vec![
714            Field::new("id", DataType::Int64, false),
715            Field::new("name", DataType::Utf8, true),
716        ]))
717    }
718
719    #[tokio::test]
720    async fn stream_parquet_many_batches_row_count() -> Result<()> {
721        let temp = tempfile::tempdir()?;
722        let dest = temp.path().join("out.parquet");
723        let ctx = SessionContext::new();
724        // Produce multiple small batches via UNION ALL chain
725        let df = ctx
726            .sql(
727                "SELECT * FROM (VALUES (1, 'a'), (2, 'b'), (3, 'c'), (4, 'd'), (5, 'e')) \
728                 AS t(id, name)",
729            )
730            .await?;
731        let stream = df.execute_stream().await?;
732        let opts = MaterializeWriteOptions {
733            max_row_group_rows: 2,
734            max_row_group_bytes: 1024,
735            fail_fast_assertions: true,
736        };
737        let assertions = vec![Assertion::UniqueKey {
738            columns: vec!["id".into()],
739        }];
740        let mut stream = stream;
741        let stats = write_parquet_stream(&mut stream, &dest, &opts, &assertions).await?;
742        assert_eq!(stats.rows, 5);
743        assert!(dest.exists());
744        assert!(!partial_path_for(&dest).exists());
745        let loaded = load_parquet_batches(&dest)?;
746        let n: usize = loaded.iter().map(|b| b.num_rows()).sum();
747        assert_eq!(n, 5);
748        Ok(())
749    }
750
751    #[tokio::test]
752    async fn iceberg_stream_versions_metadata() -> Result<()> {
753        let temp = tempfile::tempdir()?;
754        let root = temp.path().join("tbl");
755        let ctx = SessionContext::new();
756        let opts = MaterializeWriteOptions::default();
757
758        let df1 = ctx.sql("SELECT 1 AS id").await?;
759        let mut s1 = df1.execute_stream().await?;
760        write_iceberg_stream(&mut s1, &root, &opts, &[]).await?;
761        assert!(root.join("metadata/v1.metadata.json").exists());
762        assert_eq!(
763            fs::read_to_string(root.join("metadata/version-hint.text"))?.trim(),
764            "1"
765        );
766
767        let df2 = ctx.sql("SELECT 2 AS id").await?;
768        let mut s2 = df2.execute_stream().await?;
769        write_iceberg_stream(&mut s2, &root, &opts, &[]).await?;
770        assert!(root.join("metadata/v2.metadata.json").exists());
771        // prior v1 preserved
772        assert!(root.join("metadata/v1.metadata.json").exists());
773        assert_eq!(
774            fs::read_to_string(root.join("metadata/version-hint.text"))?.trim(),
775            "2"
776        );
777        Ok(())
778    }
779
780    #[tokio::test]
781    async fn stream_unique_failure_removes_partial() -> Result<()> {
782        let temp = tempfile::tempdir()?;
783        let dest = temp.path().join("dup.parquet");
784        let ctx = SessionContext::new();
785        let df = ctx
786            .sql("SELECT * FROM (VALUES (1), (1)) AS t(id)")
787            .await?;
788        let stream = df.execute_stream().await?;
789        let opts = MaterializeWriteOptions::default();
790        let assertions = vec![Assertion::UniqueKey {
791            columns: vec!["id".into()],
792        }];
793        let mut stream = stream;
794        let err = write_parquet_stream(&mut stream, &dest, &opts, &assertions)
795            .await
796            .unwrap_err()
797            .to_string();
798        assert!(
799            err.contains("E_RBT_MATERIALIZE_ASSERT") || err.contains("Duplicate"),
800            "got: {err}"
801        );
802        assert!(!dest.exists(), "failed assert must not publish dest");
803        assert!(
804            !partial_path_for(&dest).exists(),
805            "partial must be cleaned on assert fail"
806        );
807        Ok(())
808    }
809
810    #[test]
811    fn atomic_publish_replaces_existing() -> Result<()> {
812        let temp = tempfile::tempdir()?;
813        let dest = temp.path().join("f.parquet");
814        fs::write(&dest, b"old")?;
815        let partial = partial_path_for(&dest);
816        fs::write(&partial, b"new-data")?;
817        atomic_publish(&partial, &dest)?;
818        assert_eq!(fs::read(&dest)?, b"new-data");
819        assert!(!partial.exists());
820        Ok(())
821    }
822
823    #[test]
824    fn write_empty_parquet_ok() -> Result<()> {
825        let temp = tempfile::tempdir()?;
826        let dest = temp.path().join("empty.parquet");
827        write_empty_parquet(sample_schema(), &dest, &MaterializeWriteOptions::default())?;
828        assert!(dest.exists());
829        let batches = load_parquet_batches(&dest)?;
830        let n: usize = batches.iter().map(|b| b.num_rows()).sum();
831        assert_eq!(n, 0);
832        Ok(())
833    }
834
835}