Skip to main content

mito2/series_index/
writer.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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/// Options for writing a series index.
50#[derive(Debug, Clone)]
51pub struct SeriesIndexWriterOptions {
52    /// Maximum number of rows in a Parquet row group.
53    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/// Metrics collected by a [`SeriesIndexWriter`].
65#[derive(Debug, Clone, Default)]
66pub struct SeriesIndexWriterMetrics {
67    /// Number of input record batches passed to the writer.
68    pub input_batches: usize,
69    /// Number of logical rows passed to the writer.
70    pub input_rows: usize,
71    /// Number of unique time series (primary keys) aggregated by the writer.
72    pub num_series: usize,
73    /// Size of the committed index file. This remains zero for an aborted writer.
74    pub output_bytes: u64,
75    /// Time spent opening the object-store and Parquet writers.
76    pub open_elapsed: Duration,
77    /// Time spent validating, aggregating, and decoding primary keys.
78    pub aggregate_elapsed: Duration,
79    /// Time spent encoding and writing Parquet batches.
80    pub write_elapsed: Duration,
81    /// Time spent closing a completed file.
82    pub finish_elapsed: Duration,
83    /// Time spent removing an incomplete output.
84    pub cleanup_elapsed: Duration,
85    /// Whether this writer was explicitly aborted.
86    pub aborted: bool,
87}
88
89impl SeriesIndexWriterMetrics {
90    /// Returns time spent by writer-owned work.
91    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
110/// Incrementally aggregates sorted flat record batches into series-index Parquet files.
111pub 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    /// Creates a writer for `path` in `object_store`.
125    ///
126    /// The file name in `path` must be unique in the object store so aborting
127    /// the writer only removes temporary files that belong to this writer.
128    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    /// Returns the metrics collected so far.
172    pub fn metrics(&self) -> &SeriesIndexWriterMetrics {
173        &self.metrics
174    }
175
176    /// Aggregates one sorted flat record batch.
177    ///
178    /// After this method returns an error, callers must call [`Self::abort`].
179    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    /// Finishes and commits the index file.
201    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    /// Aborts the writer and removes incomplete output files.
219    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
432/// Returns the Arrow schema of a series index.
433pub 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
525/// Extracts raw i64 timestamp values for the `__series_min_ts`/`__series_max_ts`
526/// columns. NOTE: the unit is dropped; the values must always be interpreted in
527/// the unit of the region metadata the index is written with. Once this index
528/// is wired into scans, files written before/after a time index unit widening
529/// would carry mixed units — the index schema or the searcher must record the
530/// unit per file, or the index must be rebuilt on such an alter.
531fn 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
553// TODO(yingwen): Bench and optimize the performance if this is costly.
554fn 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}