1use super::TargetLoader;
75use crate::types::target::TargetColumnSpec;
76use anyhow::{Context, Result, bail};
77use std::collections::HashMap;
78use std::process::{Command, Output};
79const MAX_CLUSTER_COLUMNS: usize = 4;
83
84const DEFAULT_MAX_PARTITIONS_PER_JOB: usize = 4000;
86
87#[derive(Debug, Clone)]
89pub struct BigQueryLoader {
90 pub project: String,
91 pub dataset: String,
92 pub partition_by: Option<String>,
95 pub cluster_by: Vec<String>,
97 pub run_id: Option<String>,
102 pub max_partitions_per_job: usize,
108}
109
110impl BigQueryLoader {
111 pub fn new(project: impl Into<String>, dataset: impl Into<String>) -> Self {
112 Self {
113 project: project.into(),
114 dataset: dataset.into(),
115 partition_by: None,
116 cluster_by: Vec::new(),
117 run_id: None,
118 max_partitions_per_job: DEFAULT_MAX_PARTITIONS_PER_JOB,
119 }
120 }
121
122 pub fn partition_by(mut self, expr: impl Into<String>) -> Self {
123 self.partition_by = Some(expr.into());
124 self
125 }
126
127 pub fn run_id(mut self, id: impl Into<String>) -> Self {
129 self.run_id = Some(id.into());
130 self
131 }
132
133 pub fn cluster_by(mut self, columns: Vec<String>) -> Self {
134 self.cluster_by = columns;
135 self
136 }
137
138 fn run_bq(&self, args: &[String]) -> Result<Output> {
143 let out = Command::new("bq")
144 .arg(format!("--project_id={}", self.project))
145 .args(args)
146 .output()
147 .context("failed to run `bq` — is the Google Cloud SDK installed and on PATH?")?;
148 if !out.status.success() {
149 let detail = [clean_bq_output(&out.stdout), clean_bq_output(&out.stderr)]
150 .into_iter()
151 .filter(|s| !s.is_empty())
152 .collect::<Vec<_>>()
153 .join(" | ");
154 bail!(
155 "bq {} failed: {detail}",
156 args.first()
157 .map(String::as_str)
158 .unwrap_or("<no-subcommand>"),
159 );
160 }
161 Ok(out)
162 }
163
164 fn label_flags(&self, op: &str, table: &str) -> Vec<String> {
166 build_label_flags(op, table, self.run_id.as_deref())
167 }
168
169 fn run_sql(&self, sql: &str, op: &str, table: &str) -> Result<Output> {
172 self.run_bq(&query_args(sql, &self.label_flags(op, table)))
173 .map_err(augment_partition_limit)
174 }
175
176 fn count_rows(&self, fqtn: &str, table: &str) -> Result<u64> {
177 let out = self.run_bq(&count_args(fqtn, &self.label_flags("count", table)))?;
179 parse_count_csv(&String::from_utf8_lossy(&out.stdout))
180 }
181
182 fn plan_load_batches(&self, uris: &[String]) -> Vec<Vec<String>> {
188 match self.partition_by.as_deref() {
189 Some(col) if is_bare_column(col) => {
190 plan_hive_batches(uris, col, self.max_partitions_per_job)
191 .unwrap_or_else(|_| vec![uris.to_vec()])
192 }
193 _ => vec![uris.to_vec()],
194 }
195 }
196}
197
198impl TargetLoader for BigQueryLoader {
199 fn fqtn(&self, table: &str) -> String {
200 format!("{}.{}.{}", self.project, self.dataset, table)
201 }
202
203 fn materialize(&self, table: &str, specs: &[TargetColumnSpec], uris: &[String]) -> Result<u64> {
204 if self.cluster_by.len() > MAX_CLUSTER_COLUMNS {
205 bail!(
206 "BigQuery allows at most {MAX_CLUSTER_COLUMNS} clustering columns, got {}",
207 self.cluster_by.len()
208 );
209 }
210 for c in &self.cluster_by {
217 if !super::is_safe_load_ident(c) {
218 bail!(
219 "BigQuery load: clustering column `{}` is not a plain SQL identifier \
220 ([A-Za-z_][A-Za-z0-9_]*) — it splices into CLUSTER BY. Rename it.",
221 c.escape_default()
222 );
223 }
224 }
225 let target = self.fqtn(table);
226 let schema = build_schema(specs);
227
228 for (i, batch) in self.plan_load_batches(uris).iter().enumerate() {
235 let sql = build_load_data_sql(
236 &target,
237 i == 0, &schema,
239 &self.partition_by,
240 &self.cluster_by,
241 batch,
242 );
243 self.run_sql(&sql, "load", table)?;
244 }
245 self.count_rows(&target, table)
249 }
250
251 fn append_changelog(
252 &self,
253 table: &str,
254 specs: &[TargetColumnSpec],
255 uris: &[String],
256 pk: &[String],
257 ) -> Result<u64> {
258 use crate::load::cdc::Warehouse;
259 let mut full = crate::load::cdc::meta_column_specs(Warehouse::BigQuery);
262 full.extend(
263 specs
264 .iter()
265 .filter(|s| !is_meta_column(&s.column_name))
266 .cloned(),
267 );
268 let schema = build_schema(&full);
269
270 let changes = format!("{table}__changes");
271 let changes_fqtn = self.fqtn(&changes);
272
273 let create = build_create_changes_sql(&changes_fqtn, &schema, pk);
276 self.run_sql(&create, "create", &changes)?;
277
278 if let Some(alter) = build_alter_add_columns_sql(&changes_fqtn, &full) {
284 self.run_sql(&alter, "alter", &changes)?;
285 }
286
287 let before = self.count_rows(&changes_fqtn, &changes)?;
290 let load = build_load_data_sql(&changes_fqtn, false, &schema, &None, &[], uris);
291 self.run_sql(&load, "load", &changes)?;
292 let after = self.count_rows(&changes_fqtn, &changes)?;
293 Ok(after.saturating_sub(before))
294 }
295
296 fn warehouse(&self) -> crate::load::cdc::Warehouse {
297 crate::load::cdc::Warehouse::BigQuery
298 }
299
300 fn create_view(&self, table: &str, view_sql: &str) -> Result<()> {
301 self.run_sql(view_sql, "view", table)?;
302 Ok(())
303 }
304}
305
306fn is_meta_column(name: &str) -> bool {
310 crate::load::cdc::is_meta_column(name)
311}
312
313fn build_create_changes_sql(fqtn: &str, schema: &str, pk: &[String]) -> String {
317 let cluster_cols = pk
318 .iter()
319 .take(MAX_CLUSTER_COLUMNS)
320 .cloned()
321 .collect::<Vec<_>>()
322 .join(", ");
323 format!("CREATE TABLE IF NOT EXISTS `{fqtn}` (\n{schema}\n)\nCLUSTER BY {cluster_cols};")
324}
325
326fn build_alter_add_columns_sql(fqtn: &str, specs: &[TargetColumnSpec]) -> Option<String> {
344 if specs.is_empty() {
345 return None;
346 }
347 let adds = specs
348 .iter()
349 .map(|s| {
350 format!(
351 "ADD COLUMN IF NOT EXISTS `{}` {}",
352 s.column_name, s.target_type
353 )
354 })
355 .collect::<Vec<_>>()
356 .join(",\n ");
357 Some(format!("ALTER TABLE `{fqtn}`\n {adds};"))
358}
359
360fn is_bare_column(c: &str) -> bool {
363 !c.is_empty() && c.chars().all(|ch| ch.is_ascii_alphanumeric() || ch == '_')
364}
365
366fn hive_partition_value(uri: &str, column: &str) -> Option<String> {
369 let needle = format!("{column}=");
370 uri.split('/')
371 .find_map(|seg| seg.strip_prefix(&needle).map(str::to_string))
372}
373
374fn plan_hive_batches(uris: &[String], column: &str, max: usize) -> Result<Vec<Vec<String>>> {
378 let pairs: Vec<(&String, String)> = uris
379 .iter()
380 .map(|u| {
381 hive_partition_value(u, column)
382 .map(|v| (u, v))
383 .ok_or_else(|| anyhow::anyhow!("uri has no `{column}=` Hive segment: {u}"))
384 })
385 .collect::<Result<_>>()?;
386
387 let mut values: Vec<&str> = pairs.iter().map(|(_, v)| v.as_str()).collect();
388 values.sort_unstable();
389 values.dedup();
390 if values.len() <= max {
391 return Ok(vec![uris.to_vec()]);
392 }
393
394 let batch_of: HashMap<&str, usize> = values
396 .iter()
397 .enumerate()
398 .map(|(i, v)| (*v, i / max))
399 .collect();
400 let mut batches: Vec<Vec<String>> = vec![Vec::new(); values.len().div_ceil(max)];
401 for (u, v) in &pairs {
402 batches[batch_of[v.as_str()]].push((*u).clone());
403 }
404 Ok(batches)
405}
406
407fn table_shape_clauses(partition_by: &Option<String>, cluster_by: &[String]) -> String {
410 let mut s = String::new();
411 if let Some(expr) = partition_by {
412 s.push_str(&format!("\nPARTITION BY {expr}"));
413 }
414 if !cluster_by.is_empty() {
415 s.push_str(&format!("\nCLUSTER BY {}", cluster_by.join(", ")));
416 }
417 s
418}
419
420fn from_files(uris: &[String]) -> String {
427 let list = uris
428 .iter()
429 .map(|u| format!(" '{u}'"))
430 .collect::<Vec<_>>()
431 .join(",\n");
432 format!(
433 "FROM FILES (\n format = 'PARQUET',\n enable_list_inference = true,\n uris = [\n{list}\n ]\n)"
434 )
435}
436
437fn build_schema(specs: &[TargetColumnSpec]) -> String {
442 specs
443 .iter()
444 .map(|s| format!(" {} {}", s.column_name, s.target_type))
445 .collect::<Vec<_>>()
446 .join(",\n")
447}
448
449fn build_load_data_sql(
452 fqtn: &str,
453 overwrite: bool,
454 schema: &str,
455 partition_by: &Option<String>,
456 cluster_by: &[String],
457 uris: &[String],
458) -> String {
459 let kw = if overwrite { "OVERWRITE" } else { "INTO" };
460 let clauses = table_shape_clauses(partition_by, cluster_by);
461 format!(
462 "LOAD DATA {kw} `{fqtn}` (\n{schema}\n){clauses}\n{};",
463 from_files(uris)
464 )
465}
466
467fn query_args(sql: &str, labels: &[String]) -> Vec<String> {
468 let mut a = vec![
470 "query".into(),
471 "--use_legacy_sql=false".into(),
472 "--format=none".into(),
473 ];
474 a.extend_from_slice(labels);
475 a.push(sql.into());
476 a
477}
478
479fn count_args(fqtn: &str, labels: &[String]) -> Vec<String> {
480 let mut a = vec![
481 "query".into(),
482 "--use_legacy_sql=false".into(),
483 "--format=csv".into(),
484 ];
485 a.extend_from_slice(labels);
486 a.push(format!("SELECT COUNT(*) AS n FROM `{fqtn}`"));
487 a
488}
489
490fn build_label_flags(op: &str, table: &str, run_id: Option<&str>) -> Vec<String> {
496 let mut labels: Vec<(String, String)> = vec![
497 ("managed_by".into(), "rivet".into()),
498 ("rivet_op".into(), sanitize_label(op)),
499 ("rivet_table".into(), sanitize_label(table)),
500 ];
501 if let Some(id) = run_id {
502 labels.push(("rivet_run".into(), sanitize_label(id)));
503 }
504 labels
505 .into_iter()
506 .flat_map(|(k, v)| ["--label".to_string(), format!("{k}:{v}")])
507 .collect()
508}
509
510fn sanitize_label(s: &str) -> String {
513 let mut out: String = s
514 .chars()
515 .map(|c| {
516 let c = c.to_ascii_lowercase();
517 if c.is_ascii_alphanumeric() || c == '_' || c == '-' {
518 c
519 } else {
520 '_'
521 }
522 })
523 .collect();
524 out.truncate(63);
525 if out.is_empty() {
526 "unnamed".clone_into(&mut out);
527 }
528 out
529}
530
531fn parse_count_csv(stdout: &str) -> Result<u64> {
534 stdout
535 .lines()
536 .rev()
537 .find_map(|l| l.trim().parse::<u64>().ok())
538 .context("could not parse a row count from bq output")
539}
540
541fn clean_bq_output(bytes: &[u8]) -> String {
545 String::from_utf8_lossy(bytes)
546 .replace('\r', "\n")
547 .lines()
548 .map(str::trim)
549 .filter(|l| !l.is_empty() && !l.starts_with("Waiting on") && !l.contains("Current status:"))
550 .collect::<Vec<_>>()
551 .join(" ")
552}
553
554fn augment_partition_limit(e: anyhow::Error) -> anyhow::Error {
556 let s = e.to_string().to_lowercase();
557 if s.contains("partition")
558 && (s.contains("4000") || s.contains("quota") || s.contains("exceed"))
559 {
560 return e.context(
561 "BigQuery caps a single load/query job at 4,000 modified partitions — split the \
562 Parquet URIs into batches whose partition span is <= 4,000 (e.g. load by date range)",
563 );
564 }
565 e
566}
567#[cfg(test)]
568mod tests {
569 use super::*;
570 use crate::types::target::TargetStatus;
571
572 fn spec(name: &str, cast: Option<&str>, status: TargetStatus) -> TargetColumnSpec {
573 TargetColumnSpec {
574 column_name: name.into(),
575 target_type: "X".into(),
576 autoload_type: "Y".into(),
577 status,
578 note: None,
579 cast_sql: cast.map(String::from),
580 }
581 }
582
583 fn uris() -> Vec<String> {
584 vec!["gs://b/a.parquet".into(), "gs://b/b.parquet".into()]
585 }
586
587 fn typed(name: &str, target_type: &str) -> TargetColumnSpec {
588 TargetColumnSpec {
589 column_name: name.into(),
590 target_type: target_type.into(),
591 autoload_type: "BYTES".into(),
592 status: TargetStatus::Ok,
593 note: None,
594 cast_sql: None,
595 }
596 }
597
598 #[test]
599 fn schema_declares_each_columns_native_target_type() {
600 let s = build_schema(&[
601 typed("id", "INT64"),
602 typed("json_col", "JSON"),
603 typed("dt_col", "DATETIME"),
604 ]);
605 assert!(s.contains("id INT64"));
606 assert!(s.contains("json_col JSON"));
607 assert!(s.contains("dt_col DATETIME"));
608 }
609
610 #[test]
611 fn load_data_declares_native_schema_and_is_a_free_batch_load() {
612 let schema = build_schema(&[typed("id", "INT64"), typed("json_col", "JSON")]);
613 let sql = build_load_data_sql("p.d.orders", true, &schema, &None, &[], &uris());
614 assert!(sql.starts_with("LOAD DATA OVERWRITE `p.d.orders` ("));
615 assert!(sql.contains("json_col JSON"));
617 assert!(sql.contains("format = 'PARQUET'"));
618 assert!(sql.contains("'gs://b/a.parquet'"));
619 assert!(!sql.contains("PARTITION BY"));
620 }
621
622 #[test]
623 fn load_data_append_uses_into() {
624 let schema = build_schema(&[typed("id", "INT64")]);
625 let sql = build_load_data_sql("p.d.orders", false, &schema, &None, &[], &uris());
626 assert!(sql.starts_with("LOAD DATA INTO `p.d.orders`"));
627 }
628
629 #[test]
630 fn load_data_emits_partition_and_cluster_when_configured() {
631 let schema = build_schema(&[typed("id", "INT64")]);
632 let sql = build_load_data_sql(
633 "p.d.orders",
634 true,
635 &schema,
636 &Some("DATE(created_at)".into()),
637 &["customer_id".into(), "region".into()],
638 &uris(),
639 );
640 assert!(sql.contains("PARTITION BY DATE(created_at)"));
641 assert!(sql.contains("CLUSTER BY customer_id, region"));
642 }
643
644 #[test]
645 fn create_changes_clusters_on_pk_capped_at_four_columns() {
646 let schema = build_schema(&[typed("__op", "STRING"), typed("id", "INT64")]);
647 let sql = build_create_changes_sql("p.d.orders__changes", &schema, &["id".into()]);
648 assert!(sql.starts_with("CREATE TABLE IF NOT EXISTS `p.d.orders__changes` ("));
649 assert!(sql.contains("CLUSTER BY id"));
650 let wide: Vec<String> = ["a", "b", "c", "d", "e"]
652 .iter()
653 .map(|s| s.to_string())
654 .collect();
655 let sql2 = build_create_changes_sql("t", &schema, &wide);
656 assert!(sql2.contains("CLUSTER BY a, b, c, d"));
657 assert!(!sql2.contains(", e"));
658 }
659
660 #[test]
661 fn is_meta_column_matches_only_the_three_cdc_columns() {
662 assert!(is_meta_column("__op") && is_meta_column("__pos") && is_meta_column("__seq"));
663 assert!(!is_meta_column("id") && !is_meta_column("__op_code"));
664 }
665
666 #[test]
667 fn count_csv_skips_header() {
668 assert_eq!(parse_count_csv("n\n42\n").unwrap(), 42);
669 assert_eq!(parse_count_csv("n\n0\n").unwrap(), 0);
670 assert!(parse_count_csv("n\n").is_err());
671 }
672
673 #[test]
674 fn clean_bq_output_drops_standalone_status_and_waiting_lines() {
675 let raw = b"Waiting on bqjob_x\nCurrent status: RUNNING\nError: boom\n";
679 assert_eq!(clean_bq_output(raw), "Error: boom");
680 }
681
682 #[test]
683 fn augment_partition_limit_fires_only_on_partition_plus_signal() {
684 let aug = |m: &str| augment_partition_limit(anyhow::anyhow!("{m}")).to_string();
685 assert!(aug("too many partitions, allowed 4000").contains("split the"));
687 assert!(aug("partition quota reached").contains("split the"));
688 assert!(aug("partition count will exceed the limit").contains("split the"));
689 assert!(!aug("partition pruning is disabled").contains("split the"));
691 assert!(!aug("row quota 4000 reached").contains("split the"));
692 }
693
694 #[test]
695 fn partition_limit_error_is_augmented() {
696 let raw = anyhow::anyhow!("Too many partitions: cannot modify more than 4000 partitions");
697 let msg = augment_partition_limit(raw).to_string();
698 assert!(
699 msg.contains("split the"),
700 "expected the actionable hint: {msg}"
701 );
702 }
703
704 #[test]
705 fn job_labels_tag_managed_by_op_and_table() {
706 let flags = build_label_flags("recover", "Orders", Some("Run-7"));
707 let kv: Vec<&String> = flags.iter().skip(1).step_by(2).collect();
708 assert!(kv.iter().any(|s| *s == "managed_by:rivet"));
709 assert!(kv.iter().any(|s| *s == "rivet_op:recover"));
710 assert!(kv.iter().any(|s| *s == "rivet_table:orders")); assert!(kv.iter().any(|s| *s == "rivet_run:run-7")); assert!(flags.iter().step_by(2).all(|s| s == "--label"));
714 }
715
716 #[test]
717 fn no_run_id_omits_the_rivet_run_label() {
718 let flags = build_label_flags("load", "orders", None);
719 let kv: Vec<&String> = flags.iter().skip(1).step_by(2).collect();
720 assert!(kv.iter().any(|s| *s == "rivet_table:orders"));
721 assert!(!kv.iter().any(|s| s.starts_with("rivet_run:")));
722 }
723
724 #[test]
725 fn fqtn_qualifies_project_dataset_table() {
726 let l = BigQueryLoader::new("proj", "ds");
727 assert_eq!(l.fqtn("orders"), "proj.ds.orders");
728 }
729
730 #[test]
731 fn sanitize_label_coerces_to_bq_charset() {
732 assert_eq!(sanitize_label("My.Table!"), "my_table_");
733 assert_eq!(sanitize_label(""), "unnamed");
734 assert_eq!(sanitize_label("ok-name_1"), "ok-name_1");
735 assert_eq!(sanitize_label(&"x".repeat(80)).len(), 63);
736 }
737
738 #[test]
739 fn clean_bq_output_keeps_real_error_drops_spinner() {
740 let stdout = b"Error in query string: Too many partitions produced by query, \
744 allowed 4000, query produces at least 4200 partitions";
745 let cleaned = clean_bq_output(stdout);
746 assert!(cleaned.contains("Too many partitions") && cleaned.contains("4000"));
747 let augmented = augment_partition_limit(anyhow::anyhow!("{cleaned}")).to_string();
749 assert!(augmented.contains("split the"), "{augmented}");
750 let stderr = "Waiting on bqjob_x ... (0s) Current status: RUNNING\r\
752 Waiting on bqjob_x ... (0s) Current status: DONE";
753 assert!(clean_bq_output(stderr.as_bytes()).is_empty());
754 }
755
756 #[test]
761 fn schema_reconciliation_adds_columns_and_never_replaces() {
762 let specs = [
763 spec("id", None, TargetStatus::Ok),
764 spec("_rivet_row_hash", None, TargetStatus::Ok),
765 ];
766 let sql = build_alter_add_columns_sql("p.d.t__changes", &specs).unwrap();
767 assert!(sql.starts_with("ALTER TABLE `p.d.t__changes`"), "{sql}");
768 assert_eq!(sql.matches("ADD COLUMN IF NOT EXISTS").count(), 2, "{sql}");
771 assert!(sql.contains("`_rivet_row_hash` X"), "{sql}");
772 for forbidden in ["REPLACE", "DROP", "CREATE", "TRUNCATE", "OVERWRITE"] {
773 assert!(
774 !sql.contains(forbidden),
775 "reconciliation must be additive only, found {forbidden}: {sql}"
776 );
777 }
778 }
779
780 #[test]
783 fn schema_reconciliation_emits_nothing_for_an_empty_spec_list() {
784 assert!(build_alter_add_columns_sql("p.d.t", &[]).is_none());
785 }
786
787 #[test]
791 fn changelog_sql_is_create_if_not_exists_plus_append_only() {
792 let create = build_create_changes_sql("p.d.t__changes", " `id` INT64", &["id".into()]);
793 assert!(create.starts_with("CREATE TABLE IF NOT EXISTS"), "{create}");
794 let load =
795 build_load_data_sql("p.d.t__changes", false, " `id` INT64", &None, &[], &uris());
796 assert!(load.starts_with("LOAD DATA INTO"), "{load}");
797 assert!(!load.contains("OVERWRITE"), "{load}");
798 }
799
800 #[test]
801 fn materialize_refuses_too_many_cluster_columns() {
802 let l = BigQueryLoader::new("p", "d").cluster_by(vec![
806 "a".into(),
807 "b".into(),
808 "c".into(),
809 "d".into(),
810 "e".into(),
811 ]);
812 let err = l
813 .materialize("t", &[spec("id", None, TargetStatus::Ok)], &uris())
814 .unwrap_err()
815 .to_string();
816 assert!(err.contains("clustering"), "{err}");
817 }
818
819 #[test]
820 fn materialize_refuses_a_non_identifier_cluster_column() {
821 let l = BigQueryLoader::new("p", "d").cluster_by(vec!["id) FROM secrets; --".into()]);
826 let err = l
827 .materialize("t", &[spec("id", None, TargetStatus::Ok)], &uris())
828 .unwrap_err()
829 .to_string();
830 assert!(
831 err.contains("not a plain SQL identifier") && err.contains("CLUSTER BY"),
832 "{err}"
833 );
834 }
835
836 #[test]
837 fn hive_partition_value_parses_col_segment() {
838 assert_eq!(
839 hive_partition_value("gs://b/t/d=2023-01-01/part-0.parquet", "d").as_deref(),
840 Some("2023-01-01")
841 );
842 assert_eq!(
843 hive_partition_value("gs://b/t/created_at=2023-01-01/p.parquet", "created_at")
844 .as_deref(),
845 Some("2023-01-01")
846 );
847 assert!(hive_partition_value("gs://b/t/part-0.parquet", "d").is_none());
848 }
849
850 #[test]
851 fn is_bare_column_rejects_expressions() {
852 assert!(is_bare_column("d"));
853 assert!(is_bare_column("created_at"));
854 assert!(!is_bare_column("DATE(d)"));
855 assert!(!is_bare_column("DATE_TRUNC(d, MONTH)"));
856 assert!(!is_bare_column(""));
857 }
858
859 #[test]
860 fn hive_batches_split_by_distinct_partition_cap() {
861 let uris: Vec<String> = [
863 "gs://b/t/d=2023-01-01/a.parquet",
864 "gs://b/t/d=2023-01-01/b.parquet",
865 "gs://b/t/d=2023-01-02/a.parquet",
866 "gs://b/t/d=2023-01-03/a.parquet",
867 "gs://b/t/d=2023-01-04/a.parquet",
868 "gs://b/t/d=2023-01-05/a.parquet",
869 ]
870 .iter()
871 .map(|s| s.to_string())
872 .collect();
873 let batches = plan_hive_batches(&uris, "d", 2).unwrap();
874 assert_eq!(batches.len(), 3);
875 for b in &batches {
876 let mut days: Vec<_> = b
877 .iter()
878 .map(|u| hive_partition_value(u, "d").unwrap())
879 .collect();
880 days.sort();
881 days.dedup();
882 assert!(
883 days.len() <= 2,
884 "batch touches {} distinct days",
885 days.len()
886 );
887 }
888 assert_eq!(batches.iter().map(Vec::len).sum::<usize>(), uris.len());
890 }
891
892 #[test]
893 fn hive_batches_single_when_under_cap() {
894 let uris = vec![
895 "gs://b/t/d=2023-01-01/a.parquet".to_string(),
896 "gs://b/t/d=2023-01-02/a.parquet".to_string(),
897 ];
898 assert_eq!(plan_hive_batches(&uris, "d", 4000).unwrap().len(), 1);
899 }
900
901 #[test]
902 fn hive_batches_error_when_uri_lacks_segment() {
903 let uris = vec!["gs://b/t/no-hive/a.parquet".to_string()];
904 assert!(plan_hive_batches(&uris, "d", 2).is_err());
905 }
906
907 #[test]
914 #[ignore = "live: needs bq CLI + ADC + a GCS Parquet fixture"]
915 fn bigquery_live_load_round_trips() {
916 let Ok(project) = std::env::var("BIGQUERY_TEST_PROJECT") else {
920 eprintln!("skipping bigquery_live_load_round_trips: BIGQUERY_TEST_PROJECT unset");
921 return;
922 };
923 let dataset =
924 std::env::var("RIVET_BQ_TEST_DATASET").unwrap_or_else(|_| "rivet_test".to_string());
925 let uri = std::env::var("RIVET_BQ_TEST_PARQUET_URI").expect(
926 "set RIVET_BQ_TEST_PARQUET_URI to a GCS Parquet object matching the specs below",
927 );
928
929 let specs = vec![spec("id", None, TargetStatus::Ok)];
931
932 let loader = BigQueryLoader::new(project, dataset);
933 let report =
936 crate::load::run_load(&loader, "rivet_bq_live_test", &specs, &[uri], None, None)
937 .expect("live load should succeed");
938 assert!(
939 report.rows_loaded > 0,
940 "expected rows, got {}",
941 report.rows_loaded
942 );
943 }
944
945 #[test]
959 #[ignore = "live: needs bq CLI + ADC + a CDC change-log Parquet fixture"]
960 fn bigquery_live_cdc_view_dedups_at_least_once() {
961 let Ok(project) = std::env::var("BIGQUERY_TEST_PROJECT") else {
963 eprintln!(
964 "skipping bigquery_live_cdc_view_dedups_at_least_once: BIGQUERY_TEST_PROJECT unset"
965 );
966 return;
967 };
968 let dataset =
969 std::env::var("RIVET_BQ_TEST_DATASET").unwrap_or_else(|_| "rivet_test".to_string());
970 let uri = std::env::var("RIVET_BQ_CDC_PARQUET_URI")
971 .expect("set RIVET_BQ_CDC_PARQUET_URI to a CDC change-log Parquet object");
972 let pk = std::env::var("RIVET_BQ_CDC_PK").unwrap_or_else(|_| "id".to_string());
973 let data_cols =
976 std::env::var("RIVET_BQ_CDC_DATA_COLS").unwrap_or_else(|_| "id:INT64".to_string());
977 let specs: Vec<TargetColumnSpec> = data_cols
978 .split(',')
979 .map(|c| {
980 let (name, ty) = c.split_once(':').expect("data col must be name:TYPE");
981 typed(name, ty)
982 })
983 .collect();
984 let expected_state: u64 = std::env::var("RIVET_BQ_CDC_EXPECTED_STATE")
985 .ok()
986 .and_then(|s| s.parse().ok())
987 .unwrap_or(0);
988
989 let table = "rivet_bq_live_cdc_test";
990 let pk_cols: Vec<String> = pk.split(',').map(str::to_string).collect();
991 let loader = BigQueryLoader::new(&project, &dataset);
992
993 crate::load::run_load_cdc(
996 &loader,
997 table,
998 &specs,
999 std::slice::from_ref(&uri),
1000 &pk_cols,
1001 crate::load::cdc::SourceEngine::MySql,
1002 None,
1003 None,
1004 )
1005 .expect("first CDC append + view build should succeed");
1006 let second = crate::load::run_load_cdc(
1007 &loader,
1008 table,
1009 &specs,
1010 &[uri],
1011 &pk_cols,
1012 crate::load::cdc::SourceEngine::MySql,
1013 None,
1014 None,
1015 )
1016 .expect("second CDC append (at-least-once) should succeed");
1017 assert!(second.rows_appended > 0, "second append added rows");
1018
1019 let state_rows = loader
1022 .count_rows(&second.view, table)
1023 .expect("counting the dedup view should succeed");
1024 if expected_state > 0 {
1025 assert_eq!(
1026 state_rows, expected_state,
1027 "the view must collapse duplicates to {expected_state} distinct-PK rows \
1028 (incl tombstones), got {state_rows}"
1029 );
1030 }
1031 }
1032}