1use 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#[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#[derive(Debug, Clone)]
35pub struct MaterializeWriteOptions {
36 pub max_row_group_rows: usize,
37 pub max_row_group_bytes: usize,
38 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
69pub 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
89pub 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 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
124pub 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 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 write_iceberg_sidecar_from_parquet(&parquet_path, stats.rows, &stats.path)?;
164 Ok(stats)
165 }
166 }
167}
168
169pub 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 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 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 }
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 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 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 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
639pub 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
671pub 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
689pub 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 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 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}