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