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 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/// Options for writing a series index.
52#[derive(Debug, Clone)]
53pub struct SeriesIndexWriterOptions {
54    /// Maximum number of rows in a Parquet row group.
55    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/// Metrics collected by a [`SeriesIndexWriter`].
67#[derive(Debug, Clone, Default)]
68pub struct SeriesIndexWriterMetrics {
69    /// Number of input record batches passed to the writer.
70    pub input_batches: usize,
71    /// Number of logical rows passed to the writer.
72    pub input_rows: usize,
73    /// Number of unique time series (primary keys) aggregated by the writer.
74    pub num_series: usize,
75    /// Size of the committed index file. This remains zero for an aborted writer.
76    pub output_bytes: u64,
77    /// Time spent opening the object-store and Parquet writers.
78    pub open_elapsed: Duration,
79    /// Time spent validating, aggregating, and decoding primary keys.
80    pub aggregate_elapsed: Duration,
81    /// Time spent encoding and writing Parquet batches.
82    pub write_elapsed: Duration,
83    /// Time spent closing a completed file.
84    pub finish_elapsed: Duration,
85    /// Time spent removing an incomplete output.
86    pub cleanup_elapsed: Duration,
87    /// Whether this writer was explicitly aborted.
88    pub aborted: bool,
89}
90
91impl SeriesIndexWriterMetrics {
92    /// Returns time spent by writer-owned work.
93    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
112/// Incrementally aggregates sorted flat record batches into series-index Parquet files.
113pub struct SeriesIndexWriter {
114    tag_columns: Vec<(ColumnId, String)>,
115    schema: SchemaRef,
116    /// Unit of the min/max ts Timestamp columns, i.e. the region's time
117    /// index unit at writer creation.
118    time_unit: TimeUnit,
119    writer: ParquetIndexWriter,
120    current_primary_key: Option<Vec<u8>>,
121    current_row: Option<SeriesIndexRow>,
122    buffered_rows: Vec<SeriesIndexRow>,
123    /// Scratch offsets shared when extracting tags from sparse primary keys.
124    pk_offsets: SparseOffsetsCache,
125    /// Reusable buffer for extracting tag values.
126    tag_buf: Vec<u8>,
127    metrics: SeriesIndexWriterMetrics,
128    failed: bool,
129}
130
131impl SeriesIndexWriter {
132    /// Creates a writer for `path` in `object_store`.
133    ///
134    /// The file name in `path` must be unique in the object store so aborting
135    /// the writer only removes temporary files that belong to this writer.
136    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    /// Returns the metrics collected so far.
182    pub fn metrics(&self) -> &SeriesIndexWriterMetrics {
183        &self.metrics
184    }
185
186    /// Aggregates one sorted flat record batch.
187    ///
188    /// After this method returns an error, callers must call [`Self::abort`].
189    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    /// Finishes and commits the index file.
211    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    /// Aborts the writer and removes incomplete output files.
229    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
443/// Returns the Arrow schema of a series index.
444pub fn series_index_schema(metadata: &RegionMetadataRef) -> Result<SchemaRef> {
445    validate_metadata(metadata)?;
446    // Native Timestamp columns carry the time index unit in their datatype,
447    // so each file is interpreted in the unit it was written with even after
448    // the region's time index unit has been widened.
449    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
465/// Returns the unit of the region's timestamp time index; the writer stamps
466/// it into index files so a searcher interprets each file in the unit it was
467/// written with (the region's time index unit may have been widened since).
468fn 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
555/// Extracts the time index column's raw values in the writer's `unit`.
556/// A timestamp array carrying a different unit is an invariant violation:
557/// the alter path flushes memtables before widening the region's time index
558/// unit, so a writer never receives batches in the region's previous unit.
559/// Reject the mismatch rather than reinterpreting or rescaling the values.
560fn 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
586/// Decodes one sparse primary key into a series-index row by extracting only
587/// the reserved ids and the tag columns the index needs, instead of decoding
588/// every label into owned values.
589fn 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
619/// Builds a timestamp array in `unit` from raw i64 values.
620fn 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
629/// Builds a batch in the index `schema` from aggregated rows. The rows carry
630/// raw timestamps in `unit` (guaranteed by `timestamp_values`), which
631/// `series_index_schema` stamps into the min/max ts columns.
632fn 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        // The min/max ts columns are native timestamps in the time index unit.
766        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        // Values already in the writer's unit pass through unchanged.
818        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        // A timestamp array carrying a different unit is rejected instead of
824        // being silently reinterpreted or rescaled.
825        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        // Plain Int64 columns carry no unit to validate against and are
833        // rejected.
834        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}