use crate::Result;
use crate::error::CoreError;
use crate::metadata::meta_field::MetaField;
use crate::util::arrow::{create_row_converter, get_column_arrays};
use arrow_array::RecordBatch;
use arrow_row::{RowConverter, Rows};
use arrow_schema::SchemaRef;
pub fn create_record_key_converter(schema: SchemaRef) -> Result<RowConverter> {
create_row_converter(schema, [MetaField::RecordKey.as_ref()])
}
pub fn create_event_time_ordering_converter(
schema: SchemaRef,
ordering_field: &str,
) -> Result<RowConverter> {
create_row_converter(schema, [ordering_field])
}
pub fn create_commit_time_ordering_converter(schema: SchemaRef) -> Result<RowConverter> {
create_row_converter(schema, [MetaField::CommitTime.as_ref()])
}
pub fn extract_record_keys(converter: &RowConverter, batch: &RecordBatch) -> Result<Rows> {
let columns = get_column_arrays(batch, [MetaField::RecordKey.as_ref()])?;
converter
.convert_columns(&columns)
.map_err(CoreError::ArrowError)
}
pub fn extract_event_time_ordering_values(
converter: &RowConverter,
batch: &RecordBatch,
ordering_field: &str,
) -> Result<Rows> {
let ordering_columns = get_column_arrays(batch, [ordering_field])?;
converter
.convert_columns(&ordering_columns)
.map_err(CoreError::ArrowError)
}
pub fn extract_commit_time_ordering_values(
converter: &RowConverter,
batch: &RecordBatch,
) -> Result<Rows> {
let ordering_columns = get_column_arrays(batch, [MetaField::CommitTime.as_ref()])?;
converter
.convert_columns(&ordering_columns)
.map_err(CoreError::ArrowError)
}