1use std::cmp::Ordering;
16use std::sync::Arc;
17use std::time::{Duration, Instant};
18
19use datatypes::arrow::array::{
20 Array, ArrayRef, BinaryArray, DictionaryArray, Int64Array, StringArray, UInt32Array,
21 UInt64Array,
22};
23use datatypes::arrow::datatypes::{DataType, Field, Schema, SchemaRef, UInt32Type};
24use datatypes::arrow::record_batch::RecordBatch;
25use datatypes::prelude::ConcreteDataType;
26use datatypes::timestamp::timestamp_array_to_primitive;
27use datatypes::value::Value;
28use mito_codec::row_converter::{CompositeValues, PrimaryKeyCodec, build_primary_key_codec};
29use object_store::ObjectStore;
30use parquet::file::metadata::KeyValue;
31use snafu::{OptionExt, ResultExt, ensure};
32use store_api::codec::PrimaryKeyEncoding;
33use store_api::metadata::RegionMetadataRef;
34use store_api::storage::ColumnId;
35use store_api::storage::consts::ReservedColumnId;
36
37use crate::error::{
38 DecodeSnafu, InvalidMetaSnafu, InvalidRecordBatchSnafu, NewRecordBatchSnafu, Result,
39};
40use crate::series_index::{
41 MAX_TS_COLUMN, MIN_TS_COLUMN, ROW_COUNT_COLUMN, TABLE_ID_COLUMN, TSID_COLUMN,
42};
43use crate::sst::parquet::DEFAULT_ROW_GROUP_SIZE;
44use crate::sst::parquet::flat_format::{primary_key_column_index, time_index_column_index};
45use crate::sst::parquet::index_writer::ParquetIndexWriter;
46
47const WRITE_BATCH_SIZE: usize = 1024;
48
49#[derive(Debug, Clone)]
51pub struct SeriesIndexWriterOptions {
52 pub row_group_size: usize,
54}
55
56impl Default for SeriesIndexWriterOptions {
57 fn default() -> Self {
58 Self {
59 row_group_size: DEFAULT_ROW_GROUP_SIZE,
60 }
61 }
62}
63
64#[derive(Debug, Clone, Default)]
66pub struct SeriesIndexWriterMetrics {
67 pub input_batches: usize,
69 pub input_rows: usize,
71 pub num_series: usize,
73 pub output_bytes: u64,
75 pub open_elapsed: Duration,
77 pub aggregate_elapsed: Duration,
79 pub write_elapsed: Duration,
81 pub finish_elapsed: Duration,
83 pub cleanup_elapsed: Duration,
85 pub aborted: bool,
87}
88
89impl SeriesIndexWriterMetrics {
90 pub fn total_elapsed(&self) -> Duration {
92 self.open_elapsed
93 + self.aggregate_elapsed
94 + self.write_elapsed
95 + self.finish_elapsed
96 + self.cleanup_elapsed
97 }
98}
99
100#[derive(Debug)]
101struct SeriesIndexRow {
102 min_ts: i64,
103 max_ts: i64,
104 row_count: u64,
105 table_id: u32,
106 tsid: u64,
107 tags: Vec<Option<String>>,
108}
109
110pub struct SeriesIndexWriter {
112 codec: Arc<dyn PrimaryKeyCodec>,
113 tag_columns: Vec<(ColumnId, String)>,
114 schema: SchemaRef,
115 writer: ParquetIndexWriter,
116 current_primary_key: Option<Vec<u8>>,
117 current_row: Option<SeriesIndexRow>,
118 buffered_rows: Vec<SeriesIndexRow>,
119 metrics: SeriesIndexWriterMetrics,
120 failed: bool,
121}
122
123impl SeriesIndexWriter {
124 pub async fn try_new(
129 metadata: RegionMetadataRef,
130 object_store: ObjectStore,
131 path: &str,
132 options: SeriesIndexWriterOptions,
133 key_value_metadata: Option<Vec<KeyValue>>,
134 ) -> Result<Self> {
135 let open_start = Instant::now();
136 ensure!(
137 options.row_group_size > 0,
138 InvalidMetaSnafu {
139 reason: "series index row group size must be greater than zero",
140 }
141 );
142 let schema = series_index_schema(&metadata)?;
143 let tag_columns = tag_columns(&metadata);
144 let writer = ParquetIndexWriter::try_new(
145 "series index",
146 object_store,
147 path,
148 &schema,
149 options.row_group_size,
150 key_value_metadata,
151 )
152 .await?;
153 let codec = build_primary_key_codec(&metadata);
154
155 Ok(Self {
156 codec,
157 tag_columns,
158 schema,
159 writer,
160 current_primary_key: None,
161 current_row: None,
162 buffered_rows: Vec::with_capacity(WRITE_BATCH_SIZE),
163 metrics: SeriesIndexWriterMetrics {
164 open_elapsed: open_start.elapsed(),
165 ..Default::default()
166 },
167 failed: false,
168 })
169 }
170
171 pub fn metrics(&self) -> &SeriesIndexWriterMetrics {
173 &self.metrics
174 }
175
176 pub async fn write(&mut self, batch: &RecordBatch) -> Result<()> {
180 ensure!(
181 !self.failed,
182 InvalidRecordBatchSnafu {
183 reason: "cannot write to a failed series index writer",
184 }
185 );
186
187 self.metrics.input_batches += 1;
188 self.metrics.input_rows += batch.num_rows();
189 let aggregate_start = Instant::now();
190 let write_before = self.metrics.write_elapsed;
191 let result = self.write_inner(batch).await;
192 let write_cost = self.metrics.write_elapsed.saturating_sub(write_before);
193 self.metrics.aggregate_elapsed += aggregate_start.elapsed().saturating_sub(write_cost);
194 if result.is_err() {
195 self.failed = true;
196 }
197 result
198 }
199
200 pub async fn finish(mut self) -> Result<SeriesIndexWriterMetrics> {
202 if self.failed {
203 let error = InvalidRecordBatchSnafu {
204 reason: "cannot finish a failed series index writer",
205 }
206 .build();
207 self.cleanup().await;
208 return Err(error);
209 }
210
211 let result = self.finish_inner().await;
212 if result.is_err() {
213 self.cleanup().await;
214 }
215 result.map(|_| self.metrics)
216 }
217
218 pub async fn abort(mut self) -> Result<SeriesIndexWriterMetrics> {
220 self.metrics.aborted = true;
221 self.cleanup().await;
222 Ok(self.metrics)
223 }
224
225 async fn write_inner(&mut self, batch: &RecordBatch) -> Result<()> {
226 if batch.num_rows() == 0 {
227 return Ok(());
228 }
229 ensure!(
230 batch.num_columns() >= 4,
231 InvalidRecordBatchSnafu {
232 reason: format!(
233 "series index input has too few columns: {}",
234 batch.num_columns()
235 ),
236 }
237 );
238
239 let pk_idx = primary_key_column_index(batch.num_columns());
240 let ts_idx = time_index_column_index(batch.num_columns());
241 let primary_keys = batch.column(pk_idx);
242 let timestamps = timestamp_values(batch.column(ts_idx))?;
243 ensure!(
244 primary_keys.len() == batch.num_rows() && timestamps.len() == batch.num_rows(),
245 InvalidRecordBatchSnafu {
246 reason: "primary-key or timestamp array length does not match the batch",
247 }
248 );
249
250 if let Some(array) = primary_keys.as_any().downcast_ref::<BinaryArray>() {
251 ensure!(
252 array.null_count() == 0,
253 InvalidRecordBatchSnafu {
254 reason: "series index input contains null primary keys",
255 }
256 );
257 self.write_binary_primary_keys(array, timestamps.values())
258 .await
259 } else if let Some(array) = primary_keys
260 .as_any()
261 .downcast_ref::<DictionaryArray<UInt32Type>>()
262 {
263 ensure!(
264 array.null_count() == 0,
265 InvalidRecordBatchSnafu {
266 reason: "series index input contains null primary keys",
267 }
268 );
269 self.write_dictionary_primary_keys(array, timestamps.values())
270 .await
271 } else {
272 InvalidRecordBatchSnafu {
273 reason: format!(
274 "series index requires Binary or Dictionary(UInt32, Binary) primary keys, got {:?}",
275 primary_keys.data_type()
276 ),
277 }
278 .fail()
279 }
280 }
281
282 async fn write_binary_primary_keys(
283 &mut self,
284 primary_keys: &BinaryArray,
285 timestamps: &[i64],
286 ) -> Result<()> {
287 let mut start = 0;
288 while start < primary_keys.len() {
289 let primary_key = primary_keys.value(start);
290 let mut end = start + 1;
291 while end < primary_keys.len() && primary_keys.value(end) == primary_key {
292 end += 1;
293 }
294
295 self.update_primary_key(
296 primary_key,
297 timestamps[start],
298 timestamps[end - 1],
299 (end - start) as u64,
300 )
301 .await?;
302 start = end;
303 }
304 Ok(())
305 }
306
307 async fn write_dictionary_primary_keys(
308 &mut self,
309 primary_keys: &DictionaryArray<UInt32Type>,
310 timestamps: &[i64],
311 ) -> Result<()> {
312 let values = primary_keys
313 .values()
314 .as_any()
315 .downcast_ref::<BinaryArray>()
316 .context(InvalidRecordBatchSnafu {
317 reason: "primary-key dictionary values are not binary",
318 })?;
319 let keys = primary_keys.keys().values();
320 let mut start = 0;
321 while start < keys.len() {
322 let key = keys[start];
323 let mut end = start + 1;
324 while end < keys.len() && keys[end] == key {
325 end += 1;
326 }
327
328 self.update_primary_key(
329 values.value(key as usize),
330 timestamps[start],
331 timestamps[end - 1],
332 (end - start) as u64,
333 )
334 .await?;
335 start = end;
336 }
337 Ok(())
338 }
339
340 async fn update_primary_key(
341 &mut self,
342 primary_key: &[u8],
343 min_ts: i64,
344 max_ts: i64,
345 row_count: u64,
346 ) -> Result<()> {
347 if let Some(current) = self.current_primary_key.as_deref() {
348 match primary_key.cmp(current) {
349 Ordering::Less => {
350 return InvalidRecordBatchSnafu {
351 reason: "series index input is not sorted by primary key",
352 }
353 .fail();
354 }
355 Ordering::Equal => {
356 let row = self.current_row.as_mut().context(InvalidRecordBatchSnafu {
357 reason: "series index aggregation state is incomplete",
358 })?;
359 row.min_ts = row.min_ts.min(min_ts);
360 row.max_ts = row.max_ts.max(max_ts);
361 row.row_count += row_count;
362 return Ok(());
363 }
364 Ordering::Greater => self.finish_current_row().await?,
365 }
366 }
367
368 let row = decode_primary_key(
369 self.codec.as_ref(),
370 primary_key,
371 min_ts,
372 max_ts,
373 row_count,
374 &self.tag_columns,
375 )?;
376 self.current_primary_key = Some(primary_key.to_vec());
377 self.current_row = Some(row);
378 Ok(())
379 }
380
381 async fn finish_current_row(&mut self) -> Result<()> {
382 self.current_primary_key = None;
383 if let Some(row) = self.current_row.take() {
384 self.buffered_rows.push(row);
385 self.metrics.num_series += 1;
386 }
387 if self.buffered_rows.len() >= WRITE_BATCH_SIZE {
388 self.flush_rows().await?;
389 }
390 Ok(())
391 }
392
393 async fn flush_rows(&mut self) -> Result<()> {
394 if self.buffered_rows.is_empty() {
395 return Ok(());
396 }
397 let batch = rows_to_batch(&self.schema, &self.buffered_rows)?;
398 let start = Instant::now();
399 let result = self.writer.write(&batch).await;
400 self.metrics.write_elapsed += start.elapsed();
401 result?;
402 self.buffered_rows.clear();
403 Ok(())
404 }
405
406 async fn finish_inner(&mut self) -> Result<()> {
407 let aggregate_start = Instant::now();
408 let write_before = self.metrics.write_elapsed;
409 self.finish_current_row().await?;
410 self.flush_rows().await?;
411 let write_cost = self.metrics.write_elapsed.saturating_sub(write_before);
412 self.metrics.aggregate_elapsed += aggregate_start.elapsed().saturating_sub(write_cost);
413
414 let finish_start = Instant::now();
415 self.metrics.output_bytes = self.writer.finish().await?;
416 self.metrics.finish_elapsed += finish_start.elapsed();
417 Ok(())
418 }
419
420 async fn cleanup(&mut self) {
421 let start = Instant::now();
422 self.writer.abort().await;
423 self.current_primary_key = None;
424 self.current_row = None;
425 self.buffered_rows.clear();
426
427 self.metrics.output_bytes = 0;
428 self.metrics.cleanup_elapsed += start.elapsed();
429 }
430}
431
432pub fn series_index_schema(metadata: &RegionMetadataRef) -> Result<SchemaRef> {
434 validate_metadata(metadata)?;
435 let mut fields = vec![
436 Field::new(MIN_TS_COLUMN, DataType::Int64, false),
437 Field::new(MAX_TS_COLUMN, DataType::Int64, false),
438 Field::new(ROW_COUNT_COLUMN, DataType::UInt64, false),
439 Field::new(TABLE_ID_COLUMN, DataType::UInt32, false),
440 Field::new(TSID_COLUMN, DataType::UInt64, false),
441 ];
442 fields.extend(
443 tag_columns(metadata)
444 .into_iter()
445 .map(|(_, name)| Field::new(name, DataType::Utf8, true)),
446 );
447 Ok(Arc::new(Schema::new(fields)))
448}
449
450fn validate_metadata(metadata: &RegionMetadataRef) -> Result<()> {
451 for column in &metadata.column_metadatas {
452 ensure!(
453 !matches!(
454 column.column_schema.name.as_str(),
455 MIN_TS_COLUMN | MAX_TS_COLUMN | ROW_COUNT_COLUMN
456 ),
457 InvalidMetaSnafu {
458 reason: format!(
459 "series index internal column name {} is already in use",
460 column.column_schema.name
461 ),
462 }
463 );
464 }
465 ensure!(
466 metadata.primary_key_encoding == PrimaryKeyEncoding::Sparse,
467 InvalidMetaSnafu {
468 reason: "series index only supports sparse primary-key encoding",
469 }
470 );
471 ensure!(
472 metadata
473 .primary_key
474 .starts_with(&[ReservedColumnId::table_id(), ReservedColumnId::tsid()]),
475 InvalidMetaSnafu {
476 reason: "series index requires (__table_id, __tsid) as the primary-key prefix",
477 }
478 );
479 let table_id = metadata
480 .column_by_id(ReservedColumnId::table_id())
481 .context(InvalidMetaSnafu {
482 reason: "series index metadata is missing __table_id",
483 })?;
484 let tsid = metadata
485 .column_by_id(ReservedColumnId::tsid())
486 .context(InvalidMetaSnafu {
487 reason: "series index metadata is missing __tsid",
488 })?;
489 ensure!(
490 table_id.column_schema.data_type == ConcreteDataType::uint32_datatype()
491 && tsid.column_schema.data_type == ConcreteDataType::uint64_datatype(),
492 InvalidMetaSnafu {
493 reason: "series index requires UInt32 __table_id and UInt64 __tsid",
494 }
495 );
496 for column in metadata.primary_key_columns() {
497 if is_reserved_column(column.column_id) {
498 continue;
499 }
500 ensure!(
501 column.column_schema.data_type == ConcreteDataType::string_datatype(),
502 InvalidMetaSnafu {
503 reason: format!(
504 "series index requires string tag column {}, got {}",
505 column.column_schema.name, column.column_schema.data_type
506 ),
507 }
508 );
509 }
510 Ok(())
511}
512
513fn tag_columns(metadata: &RegionMetadataRef) -> Vec<(ColumnId, String)> {
514 metadata
515 .primary_key_columns()
516 .filter(|column| !is_reserved_column(column.column_id))
517 .map(|column| (column.column_id, column.column_schema.name.clone()))
518 .collect()
519}
520
521fn is_reserved_column(column_id: ColumnId) -> bool {
522 column_id == ReservedColumnId::table_id() || column_id == ReservedColumnId::tsid()
523}
524
525fn timestamp_values(array: &ArrayRef) -> Result<Int64Array> {
532 let timestamps = if let Some(array) = array.as_any().downcast_ref::<Int64Array>() {
533 array.clone()
534 } else {
535 timestamp_array_to_primitive(array)
536 .map(|(array, _)| array)
537 .with_context(|| InvalidRecordBatchSnafu {
538 reason: format!(
539 "series index requires an Int64 or timestamp time index, got {:?}",
540 array.data_type()
541 ),
542 })?
543 };
544 ensure!(
545 timestamps.null_count() == 0,
546 InvalidRecordBatchSnafu {
547 reason: "series index input contains null timestamps",
548 }
549 );
550 Ok(timestamps)
551}
552
553fn decode_primary_key(
555 codec: &dyn PrimaryKeyCodec,
556 primary_key: &[u8],
557 min_ts: i64,
558 max_ts: i64,
559 row_count: u64,
560 tag_columns: &[(ColumnId, String)],
561) -> Result<SeriesIndexRow> {
562 let CompositeValues::Sparse(values) = codec.decode(primary_key).context(DecodeSnafu)? else {
563 return InvalidRecordBatchSnafu {
564 reason: "decoded primary key is not sparse",
565 }
566 .fail();
567 };
568 let table_id = match values.get(&ReservedColumnId::table_id()) {
569 Some(Value::UInt32(value)) => *value,
570 value => {
571 return InvalidRecordBatchSnafu {
572 reason: format!("missing or invalid sparse __table_id: {value:?}"),
573 }
574 .fail();
575 }
576 };
577 let tsid = match values.get(&ReservedColumnId::tsid()) {
578 Some(Value::UInt64(value)) => *value,
579 value => {
580 return InvalidRecordBatchSnafu {
581 reason: format!("missing or invalid sparse __tsid: {value:?}"),
582 }
583 .fail();
584 }
585 };
586 let tags = tag_columns
587 .iter()
588 .map(|(column_id, _)| match values.get(column_id) {
589 None | Some(Value::Null) => Ok(None),
590 Some(Value::String(value)) => Ok(Some(value.as_utf8().to_string())),
591 value => InvalidRecordBatchSnafu {
592 reason: format!(
593 "invalid sparse string tag value for column {column_id}: {value:?}"
594 ),
595 }
596 .fail(),
597 })
598 .collect::<Result<Vec<_>>>()?;
599 Ok(SeriesIndexRow {
600 min_ts,
601 max_ts,
602 row_count,
603 table_id,
604 tsid,
605 tags,
606 })
607}
608
609fn rows_to_batch(schema: &SchemaRef, rows: &[SeriesIndexRow]) -> Result<RecordBatch> {
610 let mut arrays: Vec<ArrayRef> = vec![
611 Arc::new(Int64Array::from_iter_values(
612 rows.iter().map(|row| row.min_ts),
613 )),
614 Arc::new(Int64Array::from_iter_values(
615 rows.iter().map(|row| row.max_ts),
616 )),
617 Arc::new(UInt64Array::from_iter_values(
618 rows.iter().map(|row| row.row_count),
619 )),
620 Arc::new(UInt32Array::from_iter_values(
621 rows.iter().map(|row| row.table_id),
622 )),
623 Arc::new(UInt64Array::from_iter_values(
624 rows.iter().map(|row| row.tsid),
625 )),
626 ];
627 for tag_idx in 0..schema.fields().len() - 5 {
628 arrays.push(Arc::new(StringArray::from_iter(
629 rows.iter().map(|row| row.tags[tag_idx].as_deref()),
630 )));
631 }
632 RecordBatch::try_new(schema.clone(), arrays).context(NewRecordBatchSnafu)
633}
634
635#[cfg(test)]
636mod tests {
637 use api::v1::SemanticType;
638 use bytes::Bytes;
639 use datatypes::arrow::array::{
640 BinaryDictionaryBuilder, TimestampMicrosecondArray, TimestampMillisecondArray,
641 TimestampNanosecondArray, TimestampSecondArray, UInt8Array,
642 };
643 use datatypes::arrow::datatypes::{TimeUnit, UInt32Type};
644 use datatypes::schema::ColumnSchema;
645 use mito_codec::row_converter::{PrimaryKeyCodec, SparsePrimaryKeyCodec};
646 use object_store::ErrorKind;
647 use object_store::services::Memory;
648 use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
649 use store_api::codec::PrimaryKeyEncoding;
650 use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
651
652 use super::*;
653 use crate::test_util::sst_util::{new_sparse_primary_key, sst_region_metadata_with_encoding};
654
655 fn object_store() -> ObjectStore {
656 ObjectStore::new(Memory::default()).unwrap()
657 }
658
659 fn flat_schema(primary_key_type: DataType) -> SchemaRef {
660 Arc::new(Schema::new(vec![
661 Field::new(
662 "ts",
663 DataType::Timestamp(TimeUnit::Millisecond, None),
664 false,
665 ),
666 Field::new("__primary_key", primary_key_type, false),
667 Field::new("__sequence", DataType::UInt64, false),
668 Field::new("__op_type", DataType::UInt8, false),
669 ]))
670 }
671
672 fn binary_batch(primary_keys: &[&[u8]], timestamps: &[i64]) -> RecordBatch {
673 RecordBatch::try_new(
674 flat_schema(DataType::Binary),
675 vec![
676 Arc::new(TimestampMillisecondArray::from(timestamps.to_vec())),
677 Arc::new(BinaryArray::from_iter_values(primary_keys.iter().copied())),
678 Arc::new(UInt64Array::from(vec![1; timestamps.len()])),
679 Arc::new(UInt8Array::from(vec![0; timestamps.len()])),
680 ],
681 )
682 .unwrap()
683 }
684
685 fn dictionary_batch(primary_keys: &[&[u8]], timestamps: &[i64]) -> RecordBatch {
686 let mut builder = BinaryDictionaryBuilder::<UInt32Type>::new();
687 for primary_key in primary_keys {
688 builder.append(*primary_key).unwrap();
689 }
690 RecordBatch::try_new(
691 flat_schema(DataType::Dictionary(
692 Box::new(DataType::UInt32),
693 Box::new(DataType::Binary),
694 )),
695 vec![
696 Arc::new(TimestampMillisecondArray::from(timestamps.to_vec())),
697 Arc::new(builder.finish()),
698 Arc::new(UInt64Array::from(vec![1; timestamps.len()])),
699 Arc::new(UInt8Array::from(vec![0; timestamps.len()])),
700 ],
701 )
702 .unwrap()
703 }
704
705 async fn read_index(store: &ObjectStore, path: &str) -> (u64, usize, Vec<RecordBatch>) {
706 let bytes = store.read(path).await.unwrap().to_bytes();
707 let output_bytes = bytes.len() as u64;
708 let builder = ParquetRecordBatchReaderBuilder::try_new(bytes).unwrap();
709 let row_groups = builder.metadata().num_row_groups();
710 let batches = builder
711 .build()
712 .unwrap()
713 .collect::<std::result::Result<Vec<_>, _>>()
714 .unwrap();
715 (output_bytes, row_groups, batches)
716 }
717
718 #[test]
719 fn test_series_index_schema() {
720 let metadata = Arc::new(sst_region_metadata_with_encoding(
721 PrimaryKeyEncoding::Sparse,
722 ));
723 let schema = series_index_schema(&metadata).unwrap();
724 assert_eq!(
725 schema
726 .fields()
727 .iter()
728 .map(|field| field.name().as_str())
729 .collect::<Vec<_>>(),
730 [
731 "__series_min_ts",
732 "__series_max_ts",
733 "__series_row_count",
734 "__table_id",
735 "__tsid",
736 "tag_0",
737 "tag_1",
738 ]
739 );
740 assert!(!schema.field(4).is_nullable());
741 assert!(schema.field(5).is_nullable());
742
743 let dense = Arc::new(sst_region_metadata_with_encoding(PrimaryKeyEncoding::Dense));
744 assert!(series_index_schema(&dense).is_err());
745 }
746
747 #[tokio::test]
748 async fn test_reject_series_index_internal_column_names_before_opening_writer() {
749 for (index, name) in [MIN_TS_COLUMN, MAX_TS_COLUMN, ROW_COUNT_COLUMN]
750 .into_iter()
751 .enumerate()
752 {
753 let mut builder = RegionMetadataBuilder::from_existing(
754 sst_region_metadata_with_encoding(PrimaryKeyEncoding::Sparse),
755 );
756 builder.push_column_metadata(ColumnMetadata {
757 column_schema: ColumnSchema::new(name, ConcreteDataType::string_datatype(), true),
758 semantic_type: SemanticType::Field,
759 column_id: 100 + index as u32,
760 });
761 let metadata = Arc::new(builder.build().unwrap());
762 let store = object_store();
763 let path = format!("collision-{index}.parquet");
764 let error = SeriesIndexWriter::try_new(
765 metadata,
766 store.clone(),
767 &path,
768 SeriesIndexWriterOptions::default(),
769 None,
770 )
771 .await
772 .err()
773 .unwrap();
774
775 assert!(error.to_string().contains(name), "{error}");
776 assert_eq!(
777 store.stat(&path).await.unwrap_err().kind(),
778 ErrorKind::NotFound
779 );
780 }
781 }
782
783 #[test]
784 fn test_timestamp_values() {
785 let arrays: Vec<ArrayRef> = vec![
786 Arc::new(Int64Array::from(vec![1, 2])),
787 Arc::new(TimestampSecondArray::from(vec![1, 2])),
788 Arc::new(TimestampMillisecondArray::from(vec![1, 2])),
789 Arc::new(TimestampMicrosecondArray::from(vec![1, 2])),
790 Arc::new(TimestampNanosecondArray::from(vec![1, 2])),
791 ];
792 for array in arrays {
793 assert_eq!(timestamp_values(&array).unwrap().values().as_ref(), &[1, 2]);
794 }
795
796 let nulls: ArrayRef = Arc::new(TimestampMillisecondArray::from(vec![Some(1), None]));
797 let error = timestamp_values(&nulls).unwrap_err();
798 assert!(error.to_string().contains("null timestamps"), "{error}");
799
800 let unsupported: ArrayRef = Arc::new(UInt8Array::from(vec![1, 2]));
801 let error = timestamp_values(&unsupported).unwrap_err();
802 assert!(
803 error.to_string().contains("requires an Int64 or timestamp"),
804 "{error}"
805 );
806 }
807
808 #[tokio::test]
809 async fn test_write_batches_and_metrics() {
810 let metadata = Arc::new(sst_region_metadata_with_encoding(
811 PrimaryKeyEncoding::Sparse,
812 ));
813 let primary_key_1 = new_sparse_primary_key(&["a", "x"], &metadata, 1, 10);
814 let primary_key_2 = new_sparse_primary_key(&["b", "y"], &metadata, 1, 20);
815 let store = object_store();
816 let mut writer = SeriesIndexWriter::try_new(
817 metadata,
818 store.clone(),
819 "series.parquet",
820 SeriesIndexWriterOptions { row_group_size: 2 },
821 None,
822 )
823 .await
824 .unwrap();
825
826 writer
827 .write(&dictionary_batch(
828 &[
829 primary_key_1.as_slice(),
830 primary_key_1.as_slice(),
831 primary_key_1.as_slice(),
832 primary_key_1.as_slice(),
833 ],
834 &[70, 80, 90, 100],
835 ))
836 .await
837 .unwrap();
838 writer
839 .write(&binary_batch(
840 &[
841 primary_key_1.as_slice(),
842 primary_key_1.as_slice(),
843 primary_key_1.as_slice(),
844 primary_key_2.as_slice(),
845 primary_key_2.as_slice(),
846 primary_key_2.as_slice(),
847 primary_key_2.as_slice(),
848 ],
849 &[110, 120, 130, 200, 210, 220, 230],
850 ))
851 .await
852 .unwrap();
853 assert_eq!(writer.metrics().input_batches, 2);
854 assert_eq!(writer.metrics().input_rows, 11);
855
856 let metrics = writer.finish().await.unwrap();
857 assert_eq!(metrics.input_batches, 2);
858 assert_eq!(metrics.input_rows, 11);
859 assert_eq!(metrics.num_series, 2);
860 assert!(!metrics.aborted);
861
862 let (output_bytes, row_groups, batches) = read_index(&store, "series.parquet").await;
863 assert_eq!(metrics.output_bytes, output_bytes);
864 assert_eq!(row_groups, 1);
865 assert_eq!(batches.len(), 1);
866 let batch = &batches[0];
867 assert_eq!(batch.num_rows(), 2);
868 assert_eq!(
869 batch
870 .column(0)
871 .as_any()
872 .downcast_ref::<Int64Array>()
873 .unwrap(),
874 &Int64Array::from(vec![70, 200])
875 );
876 assert_eq!(
877 batch
878 .column(1)
879 .as_any()
880 .downcast_ref::<Int64Array>()
881 .unwrap(),
882 &Int64Array::from(vec![130, 230])
883 );
884 assert_eq!(
885 batch
886 .column(2)
887 .as_any()
888 .downcast_ref::<UInt64Array>()
889 .unwrap(),
890 &UInt64Array::from(vec![7, 4])
891 );
892 assert_eq!(
893 batch
894 .column(3)
895 .as_any()
896 .downcast_ref::<UInt32Array>()
897 .unwrap(),
898 &UInt32Array::from(vec![1, 1])
899 );
900 assert_eq!(
901 batch
902 .column(4)
903 .as_any()
904 .downcast_ref::<UInt64Array>()
905 .unwrap(),
906 &UInt64Array::from(vec![10, 20])
907 );
908 assert_eq!(
909 batch
910 .column(5)
911 .as_any()
912 .downcast_ref::<StringArray>()
913 .unwrap(),
914 &StringArray::from(vec![Some("a"), Some("b")])
915 );
916 }
917
918 #[tokio::test]
919 async fn test_row_group_size_and_empty_file() {
920 let metadata = Arc::new(sst_region_metadata_with_encoding(
921 PrimaryKeyEncoding::Sparse,
922 ));
923 let store = object_store();
924 let mut writer = SeriesIndexWriter::try_new(
925 metadata.clone(),
926 store.clone(),
927 "groups.parquet",
928 SeriesIndexWriterOptions { row_group_size: 2 },
929 None,
930 )
931 .await
932 .unwrap();
933 let keys = (0..5)
934 .map(|tsid| new_sparse_primary_key(&["a", "x"], &metadata, 1, tsid))
935 .collect::<Vec<_>>();
936 let key_refs = keys.iter().map(Vec::as_slice).collect::<Vec<_>>();
937 writer
938 .write(&binary_batch(&key_refs, &[1, 2, 3, 4, 5]))
939 .await
940 .unwrap();
941 writer.finish().await.unwrap();
942 let (_, row_groups, _) = read_index(&store, "groups.parquet").await;
943 assert_eq!(row_groups, 3);
944
945 let empty = SeriesIndexWriter::try_new(
946 metadata,
947 store.clone(),
948 "empty.parquet",
949 SeriesIndexWriterOptions::default(),
950 None,
951 )
952 .await
953 .unwrap()
954 .finish()
955 .await
956 .unwrap();
957 assert_eq!(empty.input_rows, 0);
958 assert_eq!(empty.num_series, 0);
959 let (output_bytes, _, batches) = read_index(&store, "empty.parquet").await;
960 assert_eq!(empty.output_bytes, output_bytes);
961 assert!(batches.is_empty());
962 }
963
964 #[tokio::test]
965 async fn test_abort_and_out_of_order_input() {
966 let metadata = Arc::new(sst_region_metadata_with_encoding(
967 PrimaryKeyEncoding::Sparse,
968 ));
969 let primary_key_1 = new_sparse_primary_key(&["a", "x"], &metadata, 1, 10);
970 let primary_key_2 = new_sparse_primary_key(&["b", "y"], &metadata, 1, 20);
971 let store = object_store();
972 let mut writer = SeriesIndexWriter::try_new(
973 metadata.clone(),
974 store.clone(),
975 "abort.parquet",
976 SeriesIndexWriterOptions { row_group_size: 1 },
977 None,
978 )
979 .await
980 .unwrap();
981 let error = writer
982 .write(&binary_batch(
983 &[primary_key_2.as_slice(), primary_key_1.as_slice()],
984 &[1, 2],
985 ))
986 .await
987 .unwrap_err();
988 assert!(error.to_string().contains("not sorted"), "{error}");
989 store
990 .write("abort.parquet", Bytes::from_static(b"existing"))
991 .await
992 .unwrap();
993 let metrics = writer.abort().await.unwrap();
994 assert!(metrics.aborted);
995 assert_eq!(metrics.output_bytes, 0);
996 assert_eq!(
997 store.read("abort.parquet").await.unwrap().to_bytes(),
998 Bytes::from_static(b"existing")
999 );
1000
1001 let mut writer = SeriesIndexWriter::try_new(
1002 metadata,
1003 store.clone(),
1004 "dictionary-abort.parquet",
1005 SeriesIndexWriterOptions { row_group_size: 1 },
1006 None,
1007 )
1008 .await
1009 .unwrap();
1010 let error = writer
1011 .write(&dictionary_batch(
1012 &[
1013 primary_key_2.as_slice(),
1014 primary_key_2.as_slice(),
1015 primary_key_1.as_slice(),
1016 ],
1017 &[1, 2, 3],
1018 ))
1019 .await
1020 .unwrap_err();
1021 assert!(error.to_string().contains("not sorted"), "{error}");
1022 writer.abort().await.unwrap();
1023 assert_eq!(
1024 store
1025 .stat("dictionary-abort.parquet")
1026 .await
1027 .unwrap_err()
1028 .kind(),
1029 ErrorKind::NotFound
1030 );
1031 }
1032
1033 #[tokio::test]
1034 async fn test_nullable_tag_and_invalid_options() {
1035 let metadata = Arc::new(sst_region_metadata_with_encoding(
1036 PrimaryKeyEncoding::Sparse,
1037 ));
1038 let codec = SparsePrimaryKeyCodec::new(&metadata);
1039 let mut primary_key = Vec::new();
1040 codec
1041 .encode_value_refs(
1042 &[
1043 (
1044 ReservedColumnId::table_id(),
1045 datatypes::value::ValueRef::UInt32(1),
1046 ),
1047 (
1048 ReservedColumnId::tsid(),
1049 datatypes::value::ValueRef::UInt64(10),
1050 ),
1051 (0, datatypes::value::ValueRef::String("a")),
1052 ],
1053 &mut primary_key,
1054 )
1055 .unwrap();
1056 let store = object_store();
1057 let mut writer = SeriesIndexWriter::try_new(
1058 metadata.clone(),
1059 store.clone(),
1060 "nullable.parquet",
1061 SeriesIndexWriterOptions::default(),
1062 None,
1063 )
1064 .await
1065 .unwrap();
1066 writer
1067 .write(&binary_batch(&[primary_key.as_slice()], &[1]))
1068 .await
1069 .unwrap();
1070 writer.finish().await.unwrap();
1071 let (_, _, batches) = read_index(&store, "nullable.parquet").await;
1072 let tag_1 = batches[0]
1073 .column(6)
1074 .as_any()
1075 .downcast_ref::<StringArray>()
1076 .unwrap();
1077 assert!(tag_1.is_null(0));
1078
1079 assert!(
1080 SeriesIndexWriter::try_new(
1081 metadata,
1082 store,
1083 "invalid.parquet",
1084 SeriesIndexWriterOptions { row_group_size: 0 },
1085 None,
1086 )
1087 .await
1088 .is_err()
1089 );
1090 }
1091}