mirror of
https://github.com/GreptimeTeam/greptimedb.git
synced 2026-10-02 18:15:36 +00:00
feat: allow customized time index unit for metric engine table (#9236)
* feat: allow customized time index unit for metric engine table * test: provide query tests * refactor: revert unnecessary change * refactor: share timestamp unit conversions in api helper Address review feedback on the time index unit changeset: - Add shared timestamp_unit/timestamp_datatype helpers to api::helper (the only crate that sees both proto ColumnDataType and TimeUnit due to layering; common-time and datatypes have no greptime-proto dep). This removes the ColumnDataType -> TimeUnit match duplicated between operator's insert path and the OTLP logs path. - Collapse the two TimeUnit <-> ValueData matches in convert_timestamp_value_data by reusing api::helper::to_grpc_value for the construction side. - Note that convert_rows_time_unit rewrites the schema before the values, so an overflow mid-batch leaves the request half-converted; harmless because the error aborts the whole insert request. Signed-off-by: Ning Sun <sunning@greptime.com> * fix: align time units per destination table and floor remote-read timestamps Address review feedback on PR #9236: - Align each metric insert request to the unit of the table it actually targets: an existing logical table keeps its own unit (it may be bound to a different physical table than the one selected by the request), and only new tables use the selected physical table's unit. The previous blanket conversion rewrote valid millisecond samples to the selected physical table's unit and the engine rejected them. Regression test: writing an existing millisecond logical table and a new table in one request that selects a microsecond physical table. - Remote read now floors narrowing timestamp conversions towards negative infinity (div_euclid), consistent with Timestamp::convert_to on the ingestion path; arrow's cast truncates towards zero and returned -1ms for a stored -1001us. Widening (second -> millisecond) keeps the exact arrow cast. Regression test: a negative, non-aligned timestamp round-trips as -2ms. Signed-off-by: Ning Sun <sunning@greptime.com> * perf: fold time unit alignment into existing table lookups Address review feedback on PR #9236: - The per-destination unit alignment no longer runs its own pass of table lookups: create_or_alter_tables_on_demand gains an align_time_index_unit parameter (metric engine path only) and converts each request inside the lookups it already performs — existing tables to their own unit, new tables to the selected physical table's. Default ingest paths now issue zero additional catalog lookups compared to main; the separate alignment pass remains only in the opt-in logical batcher pre-gate, next to the eligibility check that already looks up the same tables. - convert_rows_time_unit indexes the time index position directly (validate_column_count_match guarantees row widths) instead of Optional get_mut; the gate-side alignment validates widths itself. Signed-off-by: Ning Sun <sunning@greptime.com> * perf: resolve the batcher time index guard once per write target All batches of one remote write request share the same write target (catalog, schema, physical table), so the batcher time index guard now resolves each distinct target once instead of once per batch. Signed-off-by: Ning Sun <sunning@greptime.com> * fix: address review comments * fix: address review issue --------- Signed-off-by: Ning Sun <sunning@greptime.com>
This commit is contained in:
@@ -87,6 +87,28 @@ impl ColumnDataTypeWrapper {
|
||||
}
|
||||
}
|
||||
|
||||
/// Returns the time unit if `datatype` is a timestamp type.
|
||||
pub fn timestamp_unit(datatype: ColumnDataType) -> Option<TimeUnit> {
|
||||
match datatype {
|
||||
ColumnDataType::TimestampSecond => Some(TimeUnit::Second),
|
||||
ColumnDataType::TimestampMillisecond => Some(TimeUnit::Millisecond),
|
||||
ColumnDataType::TimestampMicrosecond => Some(TimeUnit::Microsecond),
|
||||
ColumnDataType::TimestampNanosecond => Some(TimeUnit::Nanosecond),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Returns the timestamp [ColumnDataType] for the given time unit.
|
||||
/// This is the inverse of [timestamp_unit].
|
||||
pub fn timestamp_datatype(unit: TimeUnit) -> ColumnDataType {
|
||||
match unit {
|
||||
TimeUnit::Second => ColumnDataType::TimestampSecond,
|
||||
TimeUnit::Millisecond => ColumnDataType::TimestampMillisecond,
|
||||
TimeUnit::Microsecond => ColumnDataType::TimestampMicrosecond,
|
||||
TimeUnit::Nanosecond => ColumnDataType::TimestampNanosecond,
|
||||
}
|
||||
}
|
||||
|
||||
impl From<ColumnDataTypeWrapper> for ConcreteDataType {
|
||||
fn from(datatype_wrapper: ColumnDataTypeWrapper) -> Self {
|
||||
match datatype_wrapper.datatype {
|
||||
@@ -1215,6 +1237,21 @@ mod tests {
|
||||
use super::*;
|
||||
use crate::v1::Column;
|
||||
|
||||
#[test]
|
||||
fn test_timestamp_unit_roundtrip() {
|
||||
for unit in [
|
||||
TimeUnit::Second,
|
||||
TimeUnit::Millisecond,
|
||||
TimeUnit::Microsecond,
|
||||
TimeUnit::Nanosecond,
|
||||
] {
|
||||
assert_eq!(timestamp_unit(timestamp_datatype(unit)), Some(unit));
|
||||
}
|
||||
// Non-timestamp types have no time unit.
|
||||
assert_eq!(timestamp_unit(ColumnDataType::String), None);
|
||||
assert_eq!(timestamp_unit(ColumnDataType::Datetime), None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_values_with_capacity() {
|
||||
let values = values_with_capacity(ColumnDataType::Int8, 2);
|
||||
|
||||
@@ -124,7 +124,7 @@ impl OpenTelemetryProtocolHandler for Instance {
|
||||
metric_ctx.resource_info = self.otlp_resource_info;
|
||||
|
||||
let otlp::metrics::MetricsConversion {
|
||||
requests,
|
||||
mut requests,
|
||||
rows,
|
||||
semantic_index,
|
||||
resource_info,
|
||||
@@ -173,15 +173,30 @@ impl OpenTelemetryProtocolHandler for Instance {
|
||||
let batcher = self.logical_batcher().filter(|_| {
|
||||
ctx.logical_batching_enabled() && !metric_ctx.is_legacy && metric_ctx.with_metric_engine
|
||||
});
|
||||
let batcher = if batcher.is_some()
|
||||
&& self
|
||||
let batcher = if batcher.is_some() {
|
||||
// Align the requests' time index unit with the physical table's
|
||||
// before the bulk eligibility check: the OTLP encoder keeps
|
||||
// nanosecond precision on the metric engine path, while the bulk
|
||||
// path assumes the physical table's unit (today millisecond).
|
||||
// Without this, nanosecond requests would silently skip the
|
||||
// batcher, and a non-millisecond physical table must not enter
|
||||
// the bulk path either (the eligibility check rejects its unit).
|
||||
self.inserter
|
||||
.align_metric_row_inserts_time_unit(&ctx, &physical_table, &mut requests)
|
||||
.await
|
||||
.map_err(BoxedError::new)
|
||||
.context(error::ExecuteGrpcQuerySnafu)?;
|
||||
if self
|
||||
.inserter
|
||||
.can_batch_metric_rows(&requests, &ctx, &physical_table)
|
||||
.await
|
||||
.map_err(BoxedError::new)
|
||||
.context(error::ExecuteGrpcQuerySnafu)?
|
||||
{
|
||||
batcher
|
||||
{
|
||||
batcher
|
||||
} else {
|
||||
None
|
||||
}
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
@@ -791,7 +791,9 @@ mod tests {
|
||||
use common_function::utils::partition_expr_version;
|
||||
use common_query::prelude::{greptime_native_histogram, greptime_timestamp, greptime_value};
|
||||
use common_recordbatch::RecordBatches;
|
||||
use datatypes::arrow::array::{Float64Array, TimestampMillisecondArray};
|
||||
use datatypes::arrow::array::{
|
||||
Float64Array, TimestampMicrosecondArray, TimestampMillisecondArray,
|
||||
};
|
||||
use datatypes::prelude::ConcreteDataType;
|
||||
use datatypes::schema::{ColumnDefaultConstraint, ColumnSchema};
|
||||
use datatypes::value::Value as PartitionValue;
|
||||
@@ -1116,6 +1118,91 @@ mod tests {
|
||||
expr.as_json_str().unwrap()
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_put_and_scan_microsecond_physical_region() {
|
||||
let env = TestEnv::new().await;
|
||||
let engine = env.metric();
|
||||
let physical_region_id = env.default_physical_region_id();
|
||||
let logical_region_id = env.default_logical_region_id();
|
||||
env.create_physical_region_with_ts_type(
|
||||
physical_region_id,
|
||||
&TestEnv::default_table_dir(),
|
||||
vec![],
|
||||
ConcreteDataType::timestamp_microsecond_datatype(),
|
||||
)
|
||||
.await;
|
||||
|
||||
// A logical region with a matching microsecond time index is accepted.
|
||||
let region_create_request = test_util::create_logical_region_request_with_ts_type(
|
||||
&["job"],
|
||||
physical_region_id,
|
||||
&table_dir("test", logical_region_id.table_id()),
|
||||
ConcreteDataType::timestamp_microsecond_datatype(),
|
||||
);
|
||||
engine
|
||||
.handle_request(
|
||||
logical_region_id,
|
||||
RegionRequest::Create(region_create_request),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
// Writing microsecond rows works.
|
||||
let affected_rows = engine
|
||||
.handle_request(
|
||||
logical_region_id,
|
||||
RegionRequest::Put(RegionPutRequest {
|
||||
rows: Rows {
|
||||
schema: test_util::row_schema_with_tags_and_ts_datatype(
|
||||
&["job"],
|
||||
ColumnDataType::TimestampMicrosecond,
|
||||
),
|
||||
rows: test_util::build_rows_with_ts_datatype(
|
||||
1,
|
||||
2,
|
||||
ColumnDataType::TimestampMicrosecond,
|
||||
),
|
||||
},
|
||||
hint: None,
|
||||
partition_expr_version: None,
|
||||
skip_wal: false,
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(affected_rows.affected_rows, 2);
|
||||
|
||||
// The scan returns timestamps in the physical region's unit.
|
||||
let stream = engine
|
||||
.scan_to_stream(logical_region_id, ScanRequest::default())
|
||||
.await
|
||||
.unwrap();
|
||||
let batches = RecordBatches::try_collect(stream).await.unwrap();
|
||||
let mut rows = Vec::new();
|
||||
for batch in batches.iter() {
|
||||
let batch = batch.df_record_batch();
|
||||
let timestamps = batch
|
||||
.column(batch.schema().index_of(greptime_timestamp()).unwrap())
|
||||
.as_any()
|
||||
.downcast_ref::<TimestampMicrosecondArray>()
|
||||
.unwrap();
|
||||
let values = batch
|
||||
.column(batch.schema().index_of(greptime_value()).unwrap())
|
||||
.as_any()
|
||||
.downcast_ref::<Float64Array>()
|
||||
.unwrap();
|
||||
rows.extend(
|
||||
timestamps
|
||||
.values()
|
||||
.iter()
|
||||
.copied()
|
||||
.zip(values.values().iter().copied()),
|
||||
);
|
||||
}
|
||||
rows.sort_unstable_by_key(|(timestamp, _)| *timestamp);
|
||||
assert_eq!(rows, vec![(0, 0.0), (1, 1.0)]);
|
||||
}
|
||||
|
||||
async fn create_logical_region_with_tags(
|
||||
env: &TestEnv,
|
||||
physical_region_id: RegionId,
|
||||
|
||||
@@ -158,6 +158,24 @@ impl TestEnv {
|
||||
physical_region_id: RegionId,
|
||||
table_dir: &str,
|
||||
options: Vec<(String, String)>,
|
||||
) {
|
||||
self.create_physical_region_with_ts_type(
|
||||
physical_region_id,
|
||||
table_dir,
|
||||
options,
|
||||
ConcreteDataType::timestamp_millisecond_datatype(),
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
/// Create regions in [MetricEngine] with specific `physical_region_id`
|
||||
/// and time index column type.
|
||||
pub async fn create_physical_region_with_ts_type(
|
||||
&self,
|
||||
physical_region_id: RegionId,
|
||||
table_dir: &str,
|
||||
options: Vec<(String, String)>,
|
||||
ts_type: ConcreteDataType,
|
||||
) {
|
||||
let region_create_request = RegionCreateRequest {
|
||||
engine: METRIC_ENGINE_NAME.to_string(),
|
||||
@@ -165,11 +183,7 @@ impl TestEnv {
|
||||
ColumnMetadata {
|
||||
column_id: 0,
|
||||
semantic_type: SemanticType::Timestamp,
|
||||
column_schema: ColumnSchema::new(
|
||||
greptime_timestamp(),
|
||||
ConcreteDataType::timestamp_millisecond_datatype(),
|
||||
false,
|
||||
),
|
||||
column_schema: ColumnSchema::new(greptime_timestamp(), ts_type, false),
|
||||
},
|
||||
ColumnMetadata {
|
||||
column_id: 1,
|
||||
@@ -330,16 +344,28 @@ pub fn create_logical_region_request(
|
||||
tags: &[&str],
|
||||
physical_region_id: RegionId,
|
||||
table_dir: &str,
|
||||
) -> RegionCreateRequest {
|
||||
create_logical_region_request_with_ts_type(
|
||||
tags,
|
||||
physical_region_id,
|
||||
table_dir,
|
||||
ConcreteDataType::timestamp_millisecond_datatype(),
|
||||
)
|
||||
}
|
||||
|
||||
/// Generate a [RegionCreateRequest] for logical region with the given time
|
||||
/// index column type.
|
||||
pub fn create_logical_region_request_with_ts_type(
|
||||
tags: &[&str],
|
||||
physical_region_id: RegionId,
|
||||
table_dir: &str,
|
||||
ts_type: ConcreteDataType,
|
||||
) -> RegionCreateRequest {
|
||||
let mut column_metadatas = vec![
|
||||
ColumnMetadata {
|
||||
column_id: 0,
|
||||
semantic_type: SemanticType::Timestamp,
|
||||
column_schema: ColumnSchema::new(
|
||||
greptime_timestamp(),
|
||||
ConcreteDataType::timestamp_millisecond_datatype(),
|
||||
false,
|
||||
),
|
||||
column_schema: ColumnSchema::new(greptime_timestamp(), ts_type, false),
|
||||
},
|
||||
ColumnMetadata {
|
||||
column_id: 1,
|
||||
@@ -407,10 +433,20 @@ pub fn alter_logical_region_request(tags: &[&str]) -> RegionAlterRequest {
|
||||
///
|
||||
/// The result will also contains default timestamp and value column at beginning.
|
||||
pub fn row_schema_with_tags(tags: &[&str]) -> Vec<PbColumnSchema> {
|
||||
row_schema_with_tags_and_ts_datatype(tags, ColumnDataType::TimestampMillisecond)
|
||||
}
|
||||
|
||||
/// Generate a row schema with given tag columns and time index datatype.
|
||||
///
|
||||
/// The result will also contains default timestamp and value column at beginning.
|
||||
pub fn row_schema_with_tags_and_ts_datatype(
|
||||
tags: &[&str],
|
||||
ts_datatype: ColumnDataType,
|
||||
) -> Vec<PbColumnSchema> {
|
||||
let mut schema = vec![
|
||||
PbColumnSchema {
|
||||
column_name: greptime_timestamp().to_string(),
|
||||
datatype: ColumnDataType::TimestampMillisecond as i32,
|
||||
datatype: ts_datatype as i32,
|
||||
semantic_type: SemanticType::Timestamp as _,
|
||||
datatype_extension: None,
|
||||
options: None,
|
||||
@@ -440,11 +476,33 @@ pub fn row_schema_with_tags(tags: &[&str]) -> Vec<PbColumnSchema> {
|
||||
/// The schema is generated by [row_schema_with_tags]. `num_tags` doesn't need to be precise,
|
||||
/// it's used to determine the column id for new columns.
|
||||
pub fn build_rows(num_tags: usize, num_rows: usize) -> Vec<Row> {
|
||||
build_rows_with_ts_datatype(num_tags, num_rows, ColumnDataType::TimestampMillisecond)
|
||||
}
|
||||
|
||||
/// Build [Row]s whose time index values use the given timestamp datatype.
|
||||
///
|
||||
/// The schema is generated by [row_schema_with_tags_and_ts_datatype].
|
||||
/// `num_tags` doesn't need to be precise, it's used to determine the column id
|
||||
/// for new columns.
|
||||
pub fn build_rows_with_ts_datatype(
|
||||
num_tags: usize,
|
||||
num_rows: usize,
|
||||
ts_datatype: ColumnDataType,
|
||||
) -> Vec<Row> {
|
||||
let ts_value = |i: usize| -> ValueData {
|
||||
match ts_datatype {
|
||||
ColumnDataType::TimestampSecond => ValueData::TimestampSecondValue(i as _),
|
||||
ColumnDataType::TimestampMillisecond => ValueData::TimestampMillisecondValue(i as _),
|
||||
ColumnDataType::TimestampMicrosecond => ValueData::TimestampMicrosecondValue(i as _),
|
||||
ColumnDataType::TimestampNanosecond => ValueData::TimestampNanosecondValue(i as _),
|
||||
_ => unreachable!("not a timestamp datatype"),
|
||||
}
|
||||
};
|
||||
let mut rows = vec![];
|
||||
for i in 0..num_rows {
|
||||
let mut values = vec![
|
||||
Value {
|
||||
value_data: Some(ValueData::TimestampMillisecondValue(i as _)),
|
||||
value_data: Some(ts_value(i)),
|
||||
},
|
||||
Value {
|
||||
value_data: Some(ValueData::F64Value(i as f64)),
|
||||
|
||||
+469
-9
@@ -22,6 +22,7 @@ use api::v1::region::{
|
||||
InsertRequest as RegionInsertRequest, InsertRequests as RegionInsertRequests,
|
||||
RegionRequestHeader,
|
||||
};
|
||||
use api::v1::value::ValueData;
|
||||
use api::v1::{
|
||||
AlterTableExpr, ColumnDataType, ColumnSchema, CreateTableExpr, InsertRequests,
|
||||
RowInsertRequest, RowInsertRequests, Rows, SemanticType,
|
||||
@@ -46,6 +47,8 @@ use common_query::native_histogram::{is_native_histogram_value_type, native_hist
|
||||
use common_query::prelude::{greptime_timestamp, greptime_value};
|
||||
use common_telemetry::tracing_context::TracingContext;
|
||||
use common_telemetry::{debug, error, warn};
|
||||
use common_time::Timestamp;
|
||||
use common_time::timestamp::TimeUnit;
|
||||
use datatypes::schema::SkippingIndexOptions;
|
||||
use futures_util::future;
|
||||
use meter_core::data::MeterRecord;
|
||||
@@ -429,6 +432,7 @@ impl Inserter {
|
||||
statement_executor,
|
||||
accommodate_existing_schema,
|
||||
is_single_value,
|
||||
None,
|
||||
)
|
||||
.await?;
|
||||
|
||||
@@ -565,10 +569,21 @@ impl Inserter {
|
||||
validate_column_count_match(&requests)?;
|
||||
|
||||
// check and create physical table
|
||||
self.create_physical_table_on_demand(&ctx, physical_table.clone(), statement_executor)
|
||||
let physical_table_ref = self
|
||||
.create_physical_table_on_demand(&ctx, physical_table.clone(), statement_executor)
|
||||
.await?;
|
||||
|
||||
// check and create logical tables
|
||||
// check and create logical tables; `create_or_alter_tables_on_demand`
|
||||
// aligns each request's time index unit with the unit of the table it
|
||||
// targets, inside its existing table lookups: existing tables keep
|
||||
// their own unit (which matches the physical table they are bound
|
||||
// to), and new tables use the selected physical table's unit (from
|
||||
// `physical_table_ref`). Ingestion endpoints encode timestamps in a
|
||||
// fixed unit (prometheus remote write always uses millisecond; OTLP
|
||||
// keeps nanosecond precision on the metric engine path), and the
|
||||
// metric engine requires each logical table's requests to match its
|
||||
// time index unit. Narrowing conversions truncate the sub-unit part
|
||||
// (floor), following `Timestamp::convert_to`.
|
||||
let CreateAlterTableResult {
|
||||
instant_table_ids,
|
||||
table_infos,
|
||||
@@ -580,6 +595,7 @@ impl Inserter {
|
||||
statement_executor,
|
||||
true,
|
||||
true,
|
||||
table_time_index_unit(&physical_table_ref),
|
||||
)
|
||||
.await?;
|
||||
let name_to_info = table_infos
|
||||
@@ -913,6 +929,7 @@ impl Inserter {
|
||||
statement_executor,
|
||||
false,
|
||||
false,
|
||||
None,
|
||||
)
|
||||
.await?;
|
||||
Ok(())
|
||||
@@ -930,6 +947,13 @@ impl Inserter {
|
||||
/// custom schema, and then inserts data with endpoints that have default schema setting, like prometheus
|
||||
/// remote write. This will modify the `RowInsertRequests` in place.
|
||||
/// `is_single_value` indicates whether the default schema only contains single value column so we can accommodate it.
|
||||
///
|
||||
/// `align_time_index_unit` is the selected physical metric table's time
|
||||
/// index unit; passing `Some` (metric engine path only) rewrites each
|
||||
/// request's time index column, inside this function's existing table
|
||||
/// lookups (no extra catalog access): existing destination tables are
|
||||
/// converted to their own unit, and new tables to the given unit.
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
async fn create_or_alter_tables_on_demand(
|
||||
&self,
|
||||
requests: &mut RowInsertRequests,
|
||||
@@ -938,6 +962,7 @@ impl Inserter {
|
||||
statement_executor: &StatementExecutor,
|
||||
accommodate_existing_schema: bool,
|
||||
is_single_value: bool,
|
||||
align_time_index_unit: Option<TimeUnit>,
|
||||
) -> Result<CreateAlterTableResult> {
|
||||
let _timer = crate::metrics::CREATE_ALTER_ON_DEMAND
|
||||
.with_label_values(&[auto_create_table_type.as_str()])
|
||||
@@ -958,7 +983,7 @@ impl Inserter {
|
||||
&& !has_auto_create_exempt_table
|
||||
{
|
||||
let mut instant_table_ids = HashSet::new();
|
||||
for req in &requests.inserts {
|
||||
for req in &mut requests.inserts {
|
||||
let table = match self.get_table(catalog, &schema, &req.table_name).await? {
|
||||
Some(table) => table,
|
||||
// System-defined table: created canonically by the system,
|
||||
@@ -978,6 +1003,15 @@ impl Inserter {
|
||||
.fail();
|
||||
}
|
||||
};
|
||||
// Metric path: an existing destination table keeps its own
|
||||
// time index unit (it may be bound to another physical
|
||||
// table than the one selected by this request).
|
||||
if align_time_index_unit.is_some()
|
||||
&& let Some(rows) = req.rows.as_mut()
|
||||
&& let Some(target_unit) = table_time_index_unit(&table)
|
||||
{
|
||||
convert_rows_time_unit(rows, target_unit)?;
|
||||
}
|
||||
let table_info = table.table_info();
|
||||
if matches!(auto_create_table_type, AutoCreateTableType::Trace { .. }) {
|
||||
validate_trace_table_model(&table_info, ctx)?;
|
||||
@@ -1013,6 +1047,15 @@ impl Inserter {
|
||||
if table_info.is_ttl_instant_table() {
|
||||
instant_table_ids.insert(table_info.table_id());
|
||||
}
|
||||
// Metric path: an existing destination table keeps its
|
||||
// own time index unit (it may be bound to another
|
||||
// physical table than the one selected by this request).
|
||||
if align_time_index_unit.is_some()
|
||||
&& let Some(rows) = req.rows.as_mut()
|
||||
&& let Some(target_unit) = table_time_index_unit(&table)
|
||||
{
|
||||
convert_rows_time_unit(rows, target_unit)?;
|
||||
}
|
||||
if auto_create_allowed
|
||||
&& let Some(alter_expr) = self.get_alter_table_expr_on_demand(
|
||||
req,
|
||||
@@ -1058,6 +1101,14 @@ impl Inserter {
|
||||
.fail();
|
||||
}
|
||||
None => {
|
||||
// Metric path: a new table uses the selected physical
|
||||
// table's unit; convert before the create expression is
|
||||
// derived from the request schema.
|
||||
if let Some(physical_unit) = align_time_index_unit
|
||||
&& let Some(rows) = req.rows.as_mut()
|
||||
{
|
||||
convert_rows_time_unit(rows, physical_unit)?;
|
||||
}
|
||||
let semantic_index = per_table_semantics
|
||||
.get_or_insert_with(|| parse_per_table_semantic_index(ctx))
|
||||
.as_ref();
|
||||
@@ -1256,17 +1307,16 @@ impl Inserter {
|
||||
ctx: &QueryContextRef,
|
||||
physical_table: String,
|
||||
statement_executor: &StatementExecutor,
|
||||
) -> Result<()> {
|
||||
) -> Result<TableRef> {
|
||||
let catalog_name = ctx.current_catalog();
|
||||
let schema_name = ctx.current_schema();
|
||||
|
||||
// check if exist
|
||||
if self
|
||||
if let Some(table) = self
|
||||
.get_table(catalog_name, &schema_name, &physical_table)
|
||||
.await?
|
||||
.is_some()
|
||||
{
|
||||
return Ok(());
|
||||
return Ok(table);
|
||||
}
|
||||
|
||||
// Gate here too, otherwise a disabled switch would still leak the physical table.
|
||||
@@ -1318,7 +1368,7 @@ impl Inserter {
|
||||
.await;
|
||||
|
||||
match res {
|
||||
Ok(_) => Ok(()),
|
||||
Ok(table) => Ok(table),
|
||||
Err(err) => {
|
||||
error!(err; "Failed to create table {table_reference}");
|
||||
Err(err)
|
||||
@@ -1326,6 +1376,75 @@ impl Inserter {
|
||||
}
|
||||
}
|
||||
|
||||
/// Aligns each request's time index unit with the unit of the table the
|
||||
/// request targets, so that write-path gates (e.g. the logical batcher
|
||||
/// eligibility check) observe the destination table's unit instead of the
|
||||
/// ingestion endpoint's encoding unit (Prometheus remote write always
|
||||
/// uses millisecond; OTLP keeps nanosecond precision on the metric engine
|
||||
/// path). Existing destination tables keep their own unit — they may be
|
||||
/// bound to a different physical table than `physical_table` — while new
|
||||
/// tables use `physical_table`'s unit. A missing physical table is
|
||||
/// treated as the millisecond unit of the auto-created default; the
|
||||
/// actual creation (if any) happens later in
|
||||
/// [`Inserter::handle_metric_row_inserts`].
|
||||
///
|
||||
/// This is only needed on paths that choose a write path before
|
||||
/// `handle_metric_row_inserts` (whose own table lookups perform the
|
||||
/// alignment again, as a no-op after this).
|
||||
pub async fn align_metric_row_inserts_time_unit(
|
||||
&self,
|
||||
ctx: &QueryContextRef,
|
||||
physical_table: &str,
|
||||
requests: &mut RowInsertRequests,
|
||||
) -> Result<()> {
|
||||
// The unit conversion indexes rows by the time index position, which
|
||||
// requires well-formed requests.
|
||||
validate_column_count_match(requests)?;
|
||||
let physical_unit = match self
|
||||
.get_table(ctx.current_catalog(), &ctx.current_schema(), physical_table)
|
||||
.await?
|
||||
{
|
||||
Some(table) => table_time_index_unit(&table),
|
||||
// A missing physical table is auto-created with the millisecond
|
||||
// unit later.
|
||||
None => Some(TimeUnit::Millisecond),
|
||||
};
|
||||
self.align_metric_rows_per_destination(ctx, physical_unit, requests)
|
||||
.await
|
||||
}
|
||||
|
||||
/// Converts each request's time index column to the unit of the table the
|
||||
/// request targets: the existing table's own unit when the table exists,
|
||||
/// and `physical_unit` (of the request's selected physical table) for new
|
||||
/// tables.
|
||||
async fn align_metric_rows_per_destination(
|
||||
&self,
|
||||
ctx: &QueryContextRef,
|
||||
physical_unit: Option<TimeUnit>,
|
||||
requests: &mut RowInsertRequests,
|
||||
) -> Result<()> {
|
||||
for request in &mut requests.inserts {
|
||||
let Some(rows) = request.rows.as_mut() else {
|
||||
continue;
|
||||
};
|
||||
let target_unit = match self
|
||||
.get_table(
|
||||
ctx.current_catalog(),
|
||||
&ctx.current_schema(),
|
||||
&request.table_name,
|
||||
)
|
||||
.await?
|
||||
{
|
||||
Some(table) => table_time_index_unit(&table),
|
||||
None => physical_unit,
|
||||
};
|
||||
if let Some(target_unit) = target_unit {
|
||||
convert_rows_time_unit(rows, target_unit)?;
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn get_table(
|
||||
&self,
|
||||
catalog: &str,
|
||||
@@ -1592,6 +1711,84 @@ fn request_is_native_histogram(request_schema: &[ColumnSchema]) -> bool {
|
||||
)
|
||||
}
|
||||
|
||||
/// Returns the table's time index unit, if any. A metric table without a
|
||||
/// timestamp column is left alone by the unit alignment: the metric engine
|
||||
/// rejects it anyway.
|
||||
fn table_time_index_unit(table: &TableRef) -> Option<TimeUnit> {
|
||||
table
|
||||
.table_info()
|
||||
.meta
|
||||
.schema
|
||||
.timestamp_column()
|
||||
.and_then(|col| col.data_type.as_timestamp().map(|ts| ts.unit()))
|
||||
}
|
||||
|
||||
fn convert_rows_time_unit(rows: &mut Rows, target_unit: TimeUnit) -> Result<()> {
|
||||
let Some(ts_index) = rows
|
||||
.schema
|
||||
.iter()
|
||||
.position(|col| col.semantic_type == SemanticType::Timestamp as i32)
|
||||
else {
|
||||
return Ok(());
|
||||
};
|
||||
let Some(source_unit) = ColumnDataType::try_from(rows.schema[ts_index].datatype)
|
||||
.ok()
|
||||
.and_then(api::helper::timestamp_unit)
|
||||
else {
|
||||
return Ok(());
|
||||
};
|
||||
if source_unit == target_unit {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
rows.schema[ts_index].datatype = api::helper::timestamp_datatype(target_unit) as i32;
|
||||
// Timestamp columns never carry a datatype extension.
|
||||
rows.schema[ts_index].datatype_extension = None;
|
||||
|
||||
// Note: the schema is rewritten before the rows are converted, so an
|
||||
// overflow error mid-batch leaves this request half-converted. That is
|
||||
// harmless: the error aborts the whole insert request.
|
||||
//
|
||||
// `validate_column_count_match` guarantees every row carries exactly one
|
||||
// value per schema column, so the time index position is directly in
|
||||
// bounds; no per-value search is needed.
|
||||
for row in &mut rows.rows {
|
||||
debug_assert_eq!(row.values.len(), rows.schema.len());
|
||||
let value = &mut row.values[ts_index];
|
||||
let Some(value_data) = value.value_data.take() else {
|
||||
continue;
|
||||
};
|
||||
value.value_data =
|
||||
convert_timestamp_value_data(value_data, source_unit, target_unit, ts_index)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn convert_timestamp_value_data(
|
||||
value_data: ValueData,
|
||||
source_unit: TimeUnit,
|
||||
target_unit: TimeUnit,
|
||||
column_index: usize,
|
||||
) -> Result<Option<ValueData>> {
|
||||
let timestamp = match value_data {
|
||||
ValueData::TimestampSecondValue(v) => Timestamp::new_second(v),
|
||||
ValueData::TimestampMillisecondValue(v) => Timestamp::new_millisecond(v),
|
||||
ValueData::TimestampMicrosecondValue(v) => Timestamp::new_microsecond(v),
|
||||
ValueData::TimestampNanosecondValue(v) => Timestamp::new_nanosecond(v),
|
||||
// Null or non-timestamp value; nothing to convert.
|
||||
other => return Ok(Some(other)),
|
||||
};
|
||||
let converted = timestamp
|
||||
.convert_to(target_unit)
|
||||
.with_context(|| InvalidInsertRequestSnafu {
|
||||
reason: format!(
|
||||
"timestamp column {column_index} value {} in unit {source_unit:?} overflows when converting to unit {target_unit:?}",
|
||||
timestamp.value()
|
||||
),
|
||||
})?;
|
||||
Ok(api::helper::to_grpc_value(datatypes::value::Value::Timestamp(converted)).value_data)
|
||||
}
|
||||
|
||||
fn table_is_native_histogram(table: &TableRef) -> bool {
|
||||
let mut fields = table.field_columns();
|
||||
let Some(col) = fields.next() else {
|
||||
@@ -1909,7 +2106,7 @@ mod tests {
|
||||
|
||||
use api::helper::ColumnDataTypeWrapper;
|
||||
use api::v1::helper::{field_column_schema, time_index_column_schema};
|
||||
use api::v1::{RowInsertRequest, Rows, Value};
|
||||
use api::v1::{Row, RowInsertRequest, Rows, Value};
|
||||
use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME};
|
||||
use common_meta::cache::new_table_flownode_set_cache;
|
||||
use common_meta::ddl::test_util::datanode_handler::NaiveDatanodeHandler;
|
||||
@@ -1976,6 +2173,269 @@ mod tests {
|
||||
))
|
||||
}
|
||||
|
||||
fn make_metric_physical_table_ref_with_time_unit(unit: TimeUnit) -> TableRef {
|
||||
let schema = datatypes::schema::SchemaBuilder::try_from_columns(vec![
|
||||
ColumnSchema::new(
|
||||
greptime_timestamp(),
|
||||
ConcreteDataType::timestamp_datatype(unit),
|
||||
false,
|
||||
)
|
||||
.with_time_index(true),
|
||||
ColumnSchema::new(greptime_value(), ConcreteDataType::float64_datatype(), true),
|
||||
])
|
||||
.unwrap()
|
||||
.build()
|
||||
.unwrap();
|
||||
let meta = TableMetaBuilder::empty()
|
||||
.schema(Arc::new(schema))
|
||||
.primary_key_indices(vec![])
|
||||
.value_indices(vec![1])
|
||||
.engine("metric")
|
||||
.next_column_id(0)
|
||||
.options(Default::default())
|
||||
.created_on(Default::default())
|
||||
.build()
|
||||
.unwrap();
|
||||
let info = Arc::new(
|
||||
TableInfoBuilder::default()
|
||||
.table_id(1)
|
||||
.table_version(0)
|
||||
.name("greptime_physical_table")
|
||||
.schema_name(DEFAULT_SCHEMA_NAME)
|
||||
.catalog_name(DEFAULT_CATALOG_NAME)
|
||||
.desc(None)
|
||||
.table_type(TableType::Base)
|
||||
.meta(meta)
|
||||
.build()
|
||||
.unwrap(),
|
||||
);
|
||||
Arc::new(table::Table::new(
|
||||
info,
|
||||
table::metadata::FilterPushDownType::Unsupported,
|
||||
Arc::new(DummyDataSource),
|
||||
))
|
||||
}
|
||||
|
||||
fn ms_row_insert_request(timestamp_ms: i64) -> RowInsertRequest {
|
||||
ms_row_insert_request_named("my_metric", timestamp_ms)
|
||||
}
|
||||
|
||||
fn ms_row_insert_request_named(table: &str, timestamp_ms: i64) -> RowInsertRequest {
|
||||
RowInsertRequest {
|
||||
table_name: table.to_string(),
|
||||
rows: Some(Rows {
|
||||
schema: vec![
|
||||
time_index_column_schema(
|
||||
greptime_timestamp(),
|
||||
ColumnDataType::TimestampMillisecond,
|
||||
),
|
||||
field_column_schema(greptime_value(), ColumnDataType::Float64),
|
||||
],
|
||||
rows: vec![Row {
|
||||
values: vec![
|
||||
Value {
|
||||
value_data: Some(ValueData::TimestampMillisecondValue(timestamp_ms)),
|
||||
},
|
||||
Value {
|
||||
value_data: Some(ValueData::F64Value(1.0)),
|
||||
},
|
||||
],
|
||||
}],
|
||||
}),
|
||||
}
|
||||
}
|
||||
|
||||
/// Converts the time index column of each request to `target_unit`.
|
||||
fn convert_row_insert_requests_time_unit(
|
||||
requests: &mut RowInsertRequests,
|
||||
target_unit: TimeUnit,
|
||||
) -> Result<()> {
|
||||
for request in &mut requests.inserts {
|
||||
let Some(rows) = request.rows.as_mut() else {
|
||||
continue;
|
||||
};
|
||||
convert_rows_time_unit(rows, target_unit)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_convert_row_insert_requests_time_unit_noop_when_matching() {
|
||||
let mut requests = RowInsertRequests {
|
||||
inserts: vec![ms_row_insert_request(123)],
|
||||
};
|
||||
convert_row_insert_requests_time_unit(&mut requests, TimeUnit::Millisecond).unwrap();
|
||||
let rows = requests.inserts[0].rows.as_ref().unwrap();
|
||||
assert_eq!(
|
||||
rows.schema[0].datatype,
|
||||
ColumnDataType::TimestampMillisecond as i32
|
||||
);
|
||||
assert!(matches!(
|
||||
rows.rows[0].values[0].value_data,
|
||||
Some(ValueData::TimestampMillisecondValue(123))
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_convert_row_insert_requests_time_unit_widens_losslessly() {
|
||||
let mut requests = RowInsertRequests {
|
||||
inserts: vec![ms_row_insert_request(123)],
|
||||
};
|
||||
convert_row_insert_requests_time_unit(&mut requests, TimeUnit::Microsecond).unwrap();
|
||||
let rows = requests.inserts[0].rows.as_ref().unwrap();
|
||||
assert_eq!(
|
||||
rows.schema[0].datatype,
|
||||
ColumnDataType::TimestampMicrosecond as i32
|
||||
);
|
||||
assert!(matches!(
|
||||
rows.rows[0].values[0].value_data,
|
||||
Some(ValueData::TimestampMicrosecondValue(123_000))
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_convert_row_insert_requests_time_unit_truncates_on_narrowing() {
|
||||
// 123_456_789 ns floors to 123_456 us and 123 ms; negative values
|
||||
// floor towards negative infinity, matching `Timestamp::convert_to`.
|
||||
let requests = |value_ns: i64| RowInsertRequests {
|
||||
inserts: vec![RowInsertRequest {
|
||||
table_name: "my_metric".to_string(),
|
||||
rows: Some(Rows {
|
||||
schema: vec![
|
||||
time_index_column_schema(
|
||||
greptime_timestamp(),
|
||||
ColumnDataType::TimestampNanosecond,
|
||||
),
|
||||
field_column_schema(greptime_value(), ColumnDataType::Float64),
|
||||
],
|
||||
rows: vec![Row {
|
||||
values: vec![
|
||||
Value {
|
||||
value_data: Some(ValueData::TimestampNanosecondValue(value_ns)),
|
||||
},
|
||||
Value {
|
||||
value_data: Some(ValueData::F64Value(1.0)),
|
||||
},
|
||||
],
|
||||
}],
|
||||
}),
|
||||
}],
|
||||
};
|
||||
|
||||
let mut reqs = requests(123_456_789);
|
||||
convert_row_insert_requests_time_unit(&mut reqs, TimeUnit::Microsecond).unwrap();
|
||||
let rows = reqs.inserts[0].rows.as_ref().unwrap();
|
||||
assert!(matches!(
|
||||
rows.rows[0].values[0].value_data,
|
||||
Some(ValueData::TimestampMicrosecondValue(123_456))
|
||||
));
|
||||
|
||||
let mut reqs = requests(123_456_789);
|
||||
convert_row_insert_requests_time_unit(&mut reqs, TimeUnit::Millisecond).unwrap();
|
||||
let rows = reqs.inserts[0].rows.as_ref().unwrap();
|
||||
assert!(matches!(
|
||||
rows.rows[0].values[0].value_data,
|
||||
Some(ValueData::TimestampMillisecondValue(123))
|
||||
));
|
||||
|
||||
let mut reqs = requests(-123_456_789);
|
||||
convert_row_insert_requests_time_unit(&mut reqs, TimeUnit::Microsecond).unwrap();
|
||||
let rows = reqs.inserts[0].rows.as_ref().unwrap();
|
||||
assert!(matches!(
|
||||
rows.rows[0].values[0].value_data,
|
||||
Some(ValueData::TimestampMicrosecondValue(-123_457))
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_convert_row_insert_requests_time_unit_overflow() {
|
||||
let mut requests = RowInsertRequests {
|
||||
inserts: vec![ms_row_insert_request(i64::MAX)],
|
||||
};
|
||||
let err =
|
||||
convert_row_insert_requests_time_unit(&mut requests, TimeUnit::Nanosecond).unwrap_err();
|
||||
assert!(err.to_string().contains("overflows"), "{err}");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_table_time_index_unit() {
|
||||
assert_eq!(
|
||||
table_time_index_unit(&make_metric_physical_table_ref_with_time_unit(
|
||||
TimeUnit::Microsecond
|
||||
)),
|
||||
Some(TimeUnit::Microsecond)
|
||||
);
|
||||
assert_eq!(
|
||||
table_time_index_unit(&make_metric_physical_table_ref_with_time_unit(
|
||||
TimeUnit::Millisecond
|
||||
)),
|
||||
Some(TimeUnit::Millisecond)
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_align_metric_rows_per_destination() {
|
||||
use catalog::RegisterTableRequest;
|
||||
use catalog::memory::MemoryCatalogManager;
|
||||
|
||||
// An existing millisecond logical table `existing`, plus a
|
||||
// microsecond physical table selected by the request: the existing
|
||||
// table's request keeps the millisecond unit (it is bound to another
|
||||
// physical table), while the new table's request is converted to the
|
||||
// selected physical table's microsecond unit.
|
||||
let mut inserter = batcher_test_inserter().await;
|
||||
let existing = make_metric_physical_table_ref_with_time_unit(TimeUnit::Millisecond);
|
||||
let phy_us = make_metric_physical_table_ref_with_time_unit(TimeUnit::Microsecond);
|
||||
let catalog = MemoryCatalogManager::with_default_setup();
|
||||
for (table_name, table_id, table) in [("existing", 1, existing), ("phy_us", 2, phy_us)] {
|
||||
catalog
|
||||
.register_table_sync(RegisterTableRequest {
|
||||
catalog: DEFAULT_CATALOG_NAME.to_string(),
|
||||
schema: DEFAULT_SCHEMA_NAME.to_string(),
|
||||
table_name: table_name.to_string(),
|
||||
table_id,
|
||||
table,
|
||||
})
|
||||
.unwrap();
|
||||
}
|
||||
inserter.catalog_manager = catalog;
|
||||
|
||||
let ctx = Arc::new(QueryContext::with(
|
||||
DEFAULT_CATALOG_NAME,
|
||||
DEFAULT_SCHEMA_NAME,
|
||||
));
|
||||
let mut requests = RowInsertRequests {
|
||||
inserts: vec![
|
||||
ms_row_insert_request_named("existing", 123),
|
||||
ms_row_insert_request_named("fresh", 123),
|
||||
],
|
||||
};
|
||||
inserter
|
||||
.align_metric_row_inserts_time_unit(&ctx, "phy_us", &mut requests)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let existing_rows = requests.inserts[0].rows.as_ref().unwrap();
|
||||
assert_eq!(
|
||||
existing_rows.schema[0].datatype,
|
||||
ColumnDataType::TimestampMillisecond as i32
|
||||
);
|
||||
assert!(matches!(
|
||||
existing_rows.rows[0].values[0].value_data,
|
||||
Some(ValueData::TimestampMillisecondValue(123))
|
||||
));
|
||||
|
||||
let fresh_rows = requests.inserts[1].rows.as_ref().unwrap();
|
||||
assert_eq!(
|
||||
fresh_rows.schema[0].datatype,
|
||||
ColumnDataType::TimestampMicrosecond as i32
|
||||
);
|
||||
assert!(matches!(
|
||||
fresh_rows.rows[0].values[0].value_data,
|
||||
Some(ValueData::TimestampMicrosecondValue(123_000))
|
||||
));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_accommodate_existing_schema_and_reject_kind_changes() {
|
||||
let ts_name = "my_ts";
|
||||
|
||||
@@ -21,6 +21,7 @@ mod tables;
|
||||
#[cfg(test)]
|
||||
mod test_util;
|
||||
|
||||
use std::collections::HashSet;
|
||||
use std::future::{Future, ready};
|
||||
use std::num::NonZeroUsize;
|
||||
use std::sync::Arc;
|
||||
@@ -35,6 +36,7 @@ use common_batcher::worker_registry::WorkerRegistry;
|
||||
use common_meta::cache::TableFlownodeSetCacheRef;
|
||||
use common_meta::node_manager::NodeManagerRef;
|
||||
use common_query::prelude::GREPTIME_PHYSICAL_TABLE;
|
||||
use common_time::timestamp::TimeUnit;
|
||||
use meter_core::data::MeterRecord;
|
||||
use meter_macros::write_meter;
|
||||
use partition::manager::PartitionRuleManagerRef;
|
||||
@@ -67,15 +69,15 @@ const PHYSICAL_TABLE_KEY: &str = "physical_table";
|
||||
const WORKER_IDLE_TIMEOUT_MULTIPLIER: u32 = 3;
|
||||
|
||||
#[derive(Debug, Clone, Hash, Eq, PartialEq)]
|
||||
struct BatchKey {
|
||||
catalog: String,
|
||||
schema: String,
|
||||
physical_table: String,
|
||||
skip_wal: bool,
|
||||
pub(crate) struct BatchKey {
|
||||
pub(crate) catalog: String,
|
||||
pub(crate) schema: String,
|
||||
pub(crate) physical_table: String,
|
||||
pub(crate) skip_wal: bool,
|
||||
}
|
||||
|
||||
// Requests can share a batch only when their write target and WAL policy match.
|
||||
fn batch_key_from_ctx(ctx: &QueryContextRef) -> BatchKey {
|
||||
pub(crate) fn batch_key_from_ctx(ctx: &QueryContextRef) -> BatchKey {
|
||||
let physical_table = ctx
|
||||
.extension(PHYSICAL_TABLE_KEY)
|
||||
.unwrap_or(GREPTIME_PHYSICAL_TABLE)
|
||||
@@ -171,6 +173,90 @@ impl LogicalTablePendingRowsBatcher {
|
||||
.map(|(rows, ())| rows)
|
||||
}
|
||||
|
||||
/// Returns whether the bulk path can accept `batches`:
|
||||
/// - the physical metric table selected by each context must have a
|
||||
/// millisecond time index (or not exist yet — it is auto-created with
|
||||
/// the millisecond unit): new tables are created on it and the bulk
|
||||
/// encode produces millisecond batches;
|
||||
/// - every existing destination table must itself have a millisecond time
|
||||
/// index: it may be bound to another physical table with another unit,
|
||||
/// and the bulk encoder would build millisecond arrays against its
|
||||
/// schema.
|
||||
///
|
||||
/// Incompatible requests must stay on the ordinary insert path, which
|
||||
/// converts the requests to each destination table's unit. Write targets
|
||||
/// and destination tables are resolved once per distinct name.
|
||||
pub(crate) async fn accepts_bulk_time_indexes(
|
||||
&self,
|
||||
batches: impl Iterator<Item = &(QueryContextRef, RowInsertRequests)>,
|
||||
) -> bool {
|
||||
let mut checked_targets = HashSet::new();
|
||||
let mut checked_tables = HashSet::new();
|
||||
for (ctx, requests) in batches {
|
||||
let key = batch_key_from_ctx(ctx);
|
||||
if checked_targets.insert((key.catalog, key.schema, key.physical_table))
|
||||
&& !self.accepts_physical_table_time_index(ctx).await
|
||||
{
|
||||
return false;
|
||||
}
|
||||
for request in &requests.inserts {
|
||||
if !checked_tables.insert((ctx.current_schema(), request.table_name.clone())) {
|
||||
continue;
|
||||
}
|
||||
let Ok(Some(table)) = self
|
||||
.catalog_manager
|
||||
.table(
|
||||
ctx.current_catalog(),
|
||||
&ctx.current_schema(),
|
||||
&request.table_name,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
else {
|
||||
// New table: governed by the selected physical table's
|
||||
// unit, checked above.
|
||||
continue;
|
||||
};
|
||||
if table
|
||||
.table_info()
|
||||
.meta
|
||||
.schema
|
||||
.timestamp_column()
|
||||
.and_then(|col| col.data_type.as_timestamp().map(|ts| ts.unit()))
|
||||
.is_none_or(|unit| unit != TimeUnit::Millisecond)
|
||||
{
|
||||
return false;
|
||||
}
|
||||
}
|
||||
}
|
||||
true
|
||||
}
|
||||
|
||||
/// Returns whether the physical metric table resolved from `ctx` can use
|
||||
/// the bulk path. The bulk encode produces millisecond timestamp batches
|
||||
/// only, so a physical table with another time index unit (e.g. created
|
||||
/// as `TIMESTAMP(6)`) must stay on the ordinary insert path, which
|
||||
/// converts the requests to the physical table's unit. A missing physical
|
||||
/// table is accepted: it is auto-created with the millisecond unit.
|
||||
async fn accepts_physical_table_time_index(&self, ctx: &QueryContextRef) -> bool {
|
||||
let key = batch_key_from_ctx(ctx);
|
||||
let Ok(Some(table)) = self
|
||||
.catalog_manager
|
||||
.table(&key.catalog, &key.schema, &key.physical_table, None)
|
||||
.await
|
||||
else {
|
||||
return true;
|
||||
};
|
||||
table
|
||||
.table_info()
|
||||
.meta
|
||||
.schema
|
||||
.timestamp_column()
|
||||
.and_then(|col| col.data_type.as_timestamp().map(|ts| ts.unit()))
|
||||
.map(|unit| unit == TimeUnit::Millisecond)
|
||||
.unwrap_or(true)
|
||||
}
|
||||
|
||||
/// Submits with request-level accounting after schema preparation and before
|
||||
/// queue admission. Acknowledgement follows the global batching policy.
|
||||
pub async fn submit_with<T, F>(
|
||||
|
||||
@@ -397,27 +397,45 @@ async fn write_prometheus_rows_with_progress(
|
||||
mut batches: Vec<PromWriteBatch>,
|
||||
) -> std::result::Result<PromWriteOutcome, PromWriteError> {
|
||||
if prom_store_with_metric_engine && let Some(batcher) = pending_rows_batcher {
|
||||
// Preflight before the bulk eligibility decision: pre_write hooks
|
||||
// may redirect the contexts (e.g. to a per-tenant schema), so
|
||||
// eligibility must be evaluated against the prepared batches —
|
||||
// and the fallback below reuses them without re-running hooks or
|
||||
// admission.
|
||||
preflight_prometheus_rows(&prom_store_handler, &mut batches)
|
||||
.await
|
||||
.map_err(|error| PromWriteError {
|
||||
error,
|
||||
rows_written: 0,
|
||||
})?;
|
||||
let mut rows_written = 0;
|
||||
for (temp_ctx, reqs) in batches {
|
||||
let rows = batcher
|
||||
.submit(reqs, temp_ctx)
|
||||
.await
|
||||
.map_err(|error| PromWriteError {
|
||||
error,
|
||||
rows_written,
|
||||
})?;
|
||||
rows_written += rows;
|
||||
// The bulk encode produces millisecond batches only; write targets
|
||||
// and existing destination tables with another time index unit stay
|
||||
// on the ordinary insert path, which converts the requests to each
|
||||
// destination table's unit.
|
||||
if batcher.accepts_bulk_time_indexes(batches.iter()).await {
|
||||
let mut rows_written = 0;
|
||||
for (temp_ctx, reqs) in batches {
|
||||
let rows =
|
||||
batcher
|
||||
.submit(reqs, temp_ctx)
|
||||
.await
|
||||
.map_err(|error| PromWriteError {
|
||||
error,
|
||||
rows_written,
|
||||
})?;
|
||||
rows_written += rows;
|
||||
}
|
||||
return Ok(PromWriteOutcome {
|
||||
write_cost: 0,
|
||||
rows_written,
|
||||
});
|
||||
}
|
||||
return Ok(PromWriteOutcome {
|
||||
write_cost: 0,
|
||||
rows_written,
|
||||
});
|
||||
return write_prepared_prometheus_rows_with_progress(
|
||||
prom_store_handler,
|
||||
batches,
|
||||
prom_store_with_metric_engine,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
let row_counts = batches
|
||||
@@ -456,6 +474,33 @@ async fn write_prometheus_rows_with_progress(
|
||||
})
|
||||
}
|
||||
|
||||
/// Writes already-preflighted batches through the ordinary (prepared) write
|
||||
/// path; the pre_write hooks and admission have already run.
|
||||
async fn write_prepared_prometheus_rows_with_progress(
|
||||
prom_store_handler: PromStoreProtocolHandlerRef,
|
||||
batches: Vec<PromWriteBatch>,
|
||||
prom_store_with_metric_engine: bool,
|
||||
) -> std::result::Result<PromWriteOutcome, PromWriteError> {
|
||||
let mut write_cost = 0;
|
||||
let mut rows_written = 0;
|
||||
for (ctx, request) in batches {
|
||||
let rows = prom_write_row_count(&request);
|
||||
let output = prom_store_handler
|
||||
.write_prepared(request, ctx, prom_store_with_metric_engine)
|
||||
.await
|
||||
.map_err(|error| PromWriteError {
|
||||
error,
|
||||
rows_written,
|
||||
})?;
|
||||
write_cost += output.meta.cost;
|
||||
rows_written += rows;
|
||||
}
|
||||
Ok(PromWriteOutcome {
|
||||
write_cost,
|
||||
rows_written,
|
||||
})
|
||||
}
|
||||
|
||||
async fn write_prometheus_v2_rows_with_progress(
|
||||
prom_store_handler: PromStoreProtocolHandlerRef,
|
||||
pending_rows_batcher: Option<Arc<LogicalTablePendingRowsBatcher>>,
|
||||
@@ -483,20 +528,44 @@ async fn write_prometheus_v2_rows_with_progress(
|
||||
});
|
||||
}
|
||||
|
||||
let sample_batch_count = sample_batches.len();
|
||||
let mut batches = sample_batches;
|
||||
batches.extend(histogram_batches);
|
||||
|
||||
if prom_store_with_metric_engine && let Some(batcher) = pending_rows_batcher {
|
||||
return write_batched_prometheus_v2_rows_with_progress(
|
||||
// Same ordering as the v1 path: preflight (which may redirect the
|
||||
// contexts) before the bulk eligibility decision, and the fallback
|
||||
// reuses the prepared batches without re-running hooks or admission.
|
||||
preflight_prometheus_rows(&prom_store_handler, &mut batches)
|
||||
.await
|
||||
.map_err(|error| PromWriteV2Error {
|
||||
error,
|
||||
samples_written: 0,
|
||||
histograms_written: 0,
|
||||
})?;
|
||||
// The bulk encode produces millisecond batches only; write targets
|
||||
// and existing destination tables with another time index unit stay
|
||||
// on the ordinary insert path, which converts the requests to each
|
||||
// destination table's unit.
|
||||
if batcher.accepts_bulk_time_indexes(batches.iter()).await {
|
||||
return write_batched_prometheus_v2_rows_with_progress(
|
||||
prom_store_handler,
|
||||
batcher.as_ref(),
|
||||
prom_store_with_metric_engine,
|
||||
sample_batch_count,
|
||||
batches,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
return write_prepared_prometheus_v2_rows_with_progress(
|
||||
prom_store_handler,
|
||||
batcher.as_ref(),
|
||||
batches,
|
||||
sample_batch_count,
|
||||
prom_store_with_metric_engine,
|
||||
sample_batches,
|
||||
histogram_batches,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
let sample_batch_count = sample_batches.len();
|
||||
let mut batches = sample_batches;
|
||||
batches.extend(histogram_batches);
|
||||
let row_counts = batches
|
||||
.iter()
|
||||
.map(|(_, request)| prom_write_row_count(request))
|
||||
@@ -544,24 +613,49 @@ async fn write_prometheus_v2_rows_with_progress(
|
||||
})
|
||||
}
|
||||
|
||||
/// Writes already-preflighted batches through the ordinary (prepared) write
|
||||
/// path with v2 partial-progress accounting; the pre_write hooks and
|
||||
/// admission have already run.
|
||||
async fn write_prepared_prometheus_v2_rows_with_progress(
|
||||
prom_store_handler: PromStoreProtocolHandlerRef,
|
||||
batches: Vec<PromWriteBatch>,
|
||||
sample_batch_count: usize,
|
||||
prom_store_with_metric_engine: bool,
|
||||
) -> std::result::Result<PromWriteV2Outcome, PromWriteV2Error> {
|
||||
let mut write_cost = 0;
|
||||
let mut samples_written = 0;
|
||||
let mut histograms_written = 0;
|
||||
for (index, (ctx, request)) in batches.into_iter().enumerate() {
|
||||
let rows = prom_write_row_count(&request);
|
||||
let output = prom_store_handler
|
||||
.write_prepared(request, ctx, prom_store_with_metric_engine)
|
||||
.await
|
||||
.map_err(|error| PromWriteV2Error {
|
||||
error,
|
||||
samples_written,
|
||||
histograms_written,
|
||||
})?;
|
||||
write_cost += output.meta.cost;
|
||||
if index < sample_batch_count {
|
||||
samples_written += rows;
|
||||
} else {
|
||||
histograms_written += rows;
|
||||
}
|
||||
}
|
||||
Ok(PromWriteV2Outcome {
|
||||
write_cost,
|
||||
samples_written,
|
||||
histograms_written,
|
||||
})
|
||||
}
|
||||
|
||||
async fn write_batched_prometheus_v2_rows_with_progress<B: PromWriteBatcher + ?Sized>(
|
||||
prom_store_handler: PromStoreProtocolHandlerRef,
|
||||
batcher: &B,
|
||||
prom_store_with_metric_engine: bool,
|
||||
sample_batches: Vec<PromWriteBatch>,
|
||||
histogram_batches: Vec<PromWriteBatch>,
|
||||
sample_batch_count: usize,
|
||||
batches: Vec<PromWriteBatch>,
|
||||
) -> std::result::Result<PromWriteV2Outcome, PromWriteV2Error> {
|
||||
let sample_batch_count = sample_batches.len();
|
||||
let mut batches = sample_batches;
|
||||
batches.extend(histogram_batches);
|
||||
preflight_prometheus_rows(&prom_store_handler, &mut batches)
|
||||
.await
|
||||
.map_err(|error| PromWriteV2Error {
|
||||
error,
|
||||
samples_written: 0,
|
||||
histograms_written: 0,
|
||||
})?;
|
||||
|
||||
let mut samples_written = 0;
|
||||
let mut histograms_written = 0;
|
||||
let mut write_cost = 0;
|
||||
@@ -862,14 +956,18 @@ mod tests {
|
||||
events: events.clone(),
|
||||
};
|
||||
|
||||
let Ok(outcome) = write_batched_prometheus_v2_rows_with_progress(
|
||||
handler,
|
||||
&batcher,
|
||||
true,
|
||||
vec![test_prom_write_batch("sample")],
|
||||
vec![test_prom_write_batch("histogram")],
|
||||
)
|
||||
.await
|
||||
let mut batches = vec![
|
||||
test_prom_write_batch("sample"),
|
||||
test_prom_write_batch("histogram"),
|
||||
];
|
||||
// The caller preflights before choosing the bulk path; the batched
|
||||
// writer consumes the prepared batches.
|
||||
preflight_prometheus_rows(&handler, &mut batches)
|
||||
.await
|
||||
.unwrap();
|
||||
let Ok(outcome) =
|
||||
write_batched_prometheus_v2_rows_with_progress(handler, &batcher, true, 1, batches)
|
||||
.await
|
||||
else {
|
||||
panic!("mixed remote write should succeed")
|
||||
};
|
||||
|
||||
@@ -15,7 +15,7 @@
|
||||
use std::collections::BTreeMap;
|
||||
|
||||
use ahash::{HashMap, HashMapExt};
|
||||
use api::helper::ColumnDataTypeWrapper;
|
||||
use api::helper::{ColumnDataTypeWrapper, timestamp_unit};
|
||||
use api::v1::column_data_type_extension::TypeExt;
|
||||
use api::v1::column_def::options_from_column_schema;
|
||||
use api::v1::value::ValueData;
|
||||
@@ -898,16 +898,6 @@ fn is_timestamp_type(datatype: ColumnDataType) -> bool {
|
||||
timestamp_unit(datatype).is_some()
|
||||
}
|
||||
|
||||
fn timestamp_unit(datatype: ColumnDataType) -> Option<TimeUnit> {
|
||||
match datatype {
|
||||
ColumnDataType::TimestampSecond => Some(TimeUnit::Second),
|
||||
ColumnDataType::TimestampMillisecond => Some(TimeUnit::Millisecond),
|
||||
ColumnDataType::TimestampMicrosecond => Some(TimeUnit::Microsecond),
|
||||
ColumnDataType::TimestampNanosecond => Some(TimeUnit::Nanosecond),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
fn parse_export_logs_service_request_to_rows(
|
||||
request: ExportLogsServiceRequest,
|
||||
select_info: Box<SelectInfo>,
|
||||
|
||||
@@ -998,9 +998,15 @@ fn write_timestamp(
|
||||
table: &mut TableData,
|
||||
row: &mut Vec<Value>,
|
||||
time_nano: i64,
|
||||
legacy_mode: bool,
|
||||
metric_ctx: &OtlpMetricCtx,
|
||||
) -> Result<()> {
|
||||
if legacy_mode {
|
||||
// Keep the full nanosecond precision whenever the request is headed for
|
||||
// the metric engine: `Inserter::handle_metric_row_inserts` converts the
|
||||
// timestamps to the physical table's time index unit (which may be
|
||||
// micro/nanosecond). Only the non-metric prometheus-compatible path is
|
||||
// fixed to milliseconds, to keep auto-created mito tables on the
|
||||
// millisecond time index.
|
||||
if metric_ctx.is_legacy || metric_ctx.with_metric_engine {
|
||||
row_writer::write_ts_to_nanos(
|
||||
table,
|
||||
greptime_timestamp(),
|
||||
@@ -1102,7 +1108,7 @@ fn write_tags_and_timestamp(
|
||||
)?;
|
||||
}
|
||||
|
||||
write_timestamp(table, row, timestamp_nanos, metric_ctx.is_legacy)?;
|
||||
write_timestamp(table, row, timestamp_nanos, metric_ctx)?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
@@ -1541,6 +1547,7 @@ fn encode_summary(
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use api::v1::ColumnDataType;
|
||||
use common_query::prelude::set_default_prefix;
|
||||
use otel_arrow_rust::proto::opentelemetry::common::v1::AnyValue;
|
||||
use otel_arrow_rust::proto::opentelemetry::common::v1::any_value::Value as Val;
|
||||
@@ -2625,6 +2632,81 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_metric_engine_path_keeps_nanosecond_precision() {
|
||||
let time_unix_nano = 1_704_067_200_123_456_789u64;
|
||||
let request = metrics_request(vec![Metric {
|
||||
name: "my_gauge".to_string(),
|
||||
data: Some(metric::Data::Gauge(Gauge {
|
||||
data_points: vec![NumberDataPoint {
|
||||
time_unix_nano,
|
||||
value: Some(Value::AsDouble(1.0)),
|
||||
..Default::default()
|
||||
}],
|
||||
})),
|
||||
..Default::default()
|
||||
}]);
|
||||
|
||||
// The metric engine path keeps the full nanosecond precision here;
|
||||
// `Inserter::handle_metric_row_inserts` converts the timestamps to
|
||||
// the physical table's time index unit afterwards.
|
||||
let mut metric_ctx = OtlpMetricCtx {
|
||||
with_metric_engine: true,
|
||||
..Default::default()
|
||||
};
|
||||
let MetricsConversion { requests, .. } =
|
||||
to_grpc_insert_requests(request, &mut metric_ctx).unwrap();
|
||||
|
||||
let rows = requests.inserts[0].rows.as_ref().unwrap();
|
||||
let ts_index = rows
|
||||
.schema
|
||||
.iter()
|
||||
.position(|column| column.column_name == greptime_timestamp())
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
rows.schema[ts_index].datatype,
|
||||
ColumnDataType::TimestampNanosecond as i32
|
||||
);
|
||||
assert!(matches!(
|
||||
rows.rows[0].values[ts_index].value_data,
|
||||
Some(ValueData::TimestampNanosecondValue(
|
||||
1_704_067_200_123_456_789
|
||||
))
|
||||
));
|
||||
|
||||
// The non-metric prometheus-compatible path stays millisecond so
|
||||
// auto-created mito tables keep the millisecond time index.
|
||||
let mut compat_ctx = OtlpMetricCtx::default();
|
||||
let request = metrics_request(vec![Metric {
|
||||
name: "my_gauge".to_string(),
|
||||
data: Some(metric::Data::Gauge(Gauge {
|
||||
data_points: vec![NumberDataPoint {
|
||||
time_unix_nano,
|
||||
value: Some(Value::AsDouble(1.0)),
|
||||
..Default::default()
|
||||
}],
|
||||
})),
|
||||
..Default::default()
|
||||
}]);
|
||||
let MetricsConversion { requests, .. } =
|
||||
to_grpc_insert_requests(request, &mut compat_ctx).unwrap();
|
||||
|
||||
let rows = requests.inserts[0].rows.as_ref().unwrap();
|
||||
let ts_index = rows
|
||||
.schema
|
||||
.iter()
|
||||
.position(|column| column.column_name == greptime_timestamp())
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
rows.schema[ts_index].datatype,
|
||||
ColumnDataType::TimestampMillisecond as i32
|
||||
);
|
||||
assert!(matches!(
|
||||
rows.rows[0].values[ts_index].value_data,
|
||||
Some(ValueData::TimestampMillisecondValue(1_704_067_200_123))
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_exponential_histogram_rejection_metrics() {
|
||||
// Other conversion tests update the same process-global counters.
|
||||
|
||||
@@ -18,14 +18,19 @@ use std::cmp::Ordering;
|
||||
use std::collections::HashMap;
|
||||
use std::collections::hash_map::DefaultHasher;
|
||||
use std::hash::{Hash, Hasher};
|
||||
use std::sync::Arc;
|
||||
|
||||
use api::prom_store::remote::label_matcher::Type as MatcherType;
|
||||
use api::prom_store::remote::{Label, Query, Sample, TimeSeries, WriteRequest};
|
||||
use api::v1::RowInsertRequests;
|
||||
use arrow::array::{
|
||||
Array, AsArray, DictionaryArray, LargeStringArray, StringArray, StringViewArray,
|
||||
Array, ArrayRef, AsArray, DictionaryArray, LargeStringArray, StringArray, StringViewArray,
|
||||
};
|
||||
use arrow::compute::kernels::cast as casts;
|
||||
use arrow::datatypes::{
|
||||
DataType, Float64Type, TimeUnit, TimestampMicrosecondType, TimestampMillisecondType,
|
||||
TimestampNanosecondType, UInt32Type,
|
||||
};
|
||||
use arrow::datatypes::{Float64Type, TimestampMillisecondType, UInt32Type};
|
||||
use common_grpc::precision::Precision;
|
||||
use common_query::prelude::{greptime_timestamp, greptime_value};
|
||||
use common_recordbatch::{RecordBatch, RecordBatches};
|
||||
@@ -398,6 +403,50 @@ fn append_recordbatch_to_timeseries(
|
||||
.with_context(|| error::InvalidPromRemoteReadQueryResultSnafu {
|
||||
msg: format!("missing timestamp column '{timestamp_column_name}' in query result"),
|
||||
})?;
|
||||
// The Prometheus remote read wire format carries millisecond timestamps,
|
||||
// while the table's time index can use any time unit (e.g. a metric
|
||||
// physical table created with TIMESTAMP(6)). Narrowing floors towards
|
||||
// negative infinity, consistent with `Timestamp::convert_to` on the
|
||||
// ingestion path; arrow's cast would truncate towards zero, so a
|
||||
// pre-epoch sample (e.g. -1001us) would round to -1ms instead of -2ms
|
||||
// and disagree with what ingesting the same instant into a millisecond
|
||||
// table would have stored.
|
||||
let ts_column: ArrayRef = match ts_column.data_type() {
|
||||
DataType::Timestamp(TimeUnit::Millisecond, _) => ts_column.clone(),
|
||||
// Second -> millisecond is a widening (exact) conversion.
|
||||
DataType::Timestamp(TimeUnit::Second, _) => casts::cast(
|
||||
ts_column,
|
||||
&DataType::Timestamp(TimeUnit::Millisecond, None),
|
||||
)
|
||||
.map_err(|e| {
|
||||
error::InvalidPromRemoteReadQueryResultSnafu {
|
||||
msg: format!(
|
||||
"failed to cast timestamp column '{timestamp_column_name}' of datatype {:?} to millisecond: {e}",
|
||||
ts_column.data_type()
|
||||
),
|
||||
}
|
||||
.build()
|
||||
})?,
|
||||
DataType::Timestamp(TimeUnit::Microsecond, _) => Arc::new(
|
||||
ts_column
|
||||
.as_primitive::<TimestampMicrosecondType>()
|
||||
.unary::<_, TimestampMillisecondType>(|v| v.div_euclid(1_000)),
|
||||
),
|
||||
DataType::Timestamp(TimeUnit::Nanosecond, _) => Arc::new(
|
||||
ts_column
|
||||
.as_primitive::<TimestampNanosecondType>()
|
||||
.unary::<_, TimestampMillisecondType>(|v| v.div_euclid(1_000_000)),
|
||||
),
|
||||
_ => {
|
||||
return error::InvalidPromRemoteReadQueryResultSnafu {
|
||||
msg: format!(
|
||||
"Expect timestamp column of datatype Timestamp(Millisecond), actual {:?}",
|
||||
ts_column.data_type()
|
||||
),
|
||||
}
|
||||
.fail();
|
||||
}
|
||||
};
|
||||
let ts_column = ts_column
|
||||
.as_primitive_opt::<TimestampMillisecondType>()
|
||||
.with_context(|| error::InvalidPromRemoteReadQueryResultSnafu {
|
||||
|
||||
@@ -19,6 +19,10 @@ mod test {
|
||||
use client::{DEFAULT_CATALOG_NAME, OutputData};
|
||||
use common_recordbatch::RecordBatches;
|
||||
use datatypes::arrow::array::AsArray;
|
||||
use datatypes::arrow::datatypes::{
|
||||
DataType, TimeUnit, TimestampMicrosecondType, TimestampMillisecondType,
|
||||
TimestampNanosecondType, TimestampSecondType,
|
||||
};
|
||||
use frontend::instance::Instance;
|
||||
use opentelemetry_proto::tonic::collector::trace::v1::ExportTraceServiceRequest;
|
||||
use otel_arrow_rust::proto::opentelemetry::collector::metrics::v1::ExportMetricsServiceRequest;
|
||||
@@ -37,6 +41,7 @@ mod test {
|
||||
use servers::query_handler::OpenTelemetryProtocolHandler;
|
||||
use servers::query_handler::sql::SqlQueryHandler;
|
||||
use session::context::QueryContext;
|
||||
use session::protocol_ctx::{OtlpMetricCtx, ProtocolCtx};
|
||||
|
||||
use crate::standalone::GreptimeDbStandaloneBuilder;
|
||||
use crate::tests;
|
||||
@@ -1164,6 +1169,162 @@ WITH(
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
pub async fn test_otlp_metrics_into_microsecond_physical_table_on_standalone() {
|
||||
let standalone = GreptimeDbStandaloneBuilder::new("test_otlp_us_physical")
|
||||
.build()
|
||||
.await;
|
||||
|
||||
test_otlp_metrics_into_non_millisecond_physical_table(
|
||||
standalone.fe_instance(),
|
||||
"TIMESTAMP(6)",
|
||||
TimeUnit::Microsecond,
|
||||
1_704_067_200_123_456,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
pub async fn test_otlp_metrics_into_microsecond_physical_table_on_distributed() {
|
||||
let instance = tests::create_distributed_instance("test_otlp_us_physical_dist").await;
|
||||
|
||||
test_otlp_metrics_into_non_millisecond_physical_table(
|
||||
&instance.frontend(),
|
||||
"TIMESTAMP(6)",
|
||||
TimeUnit::Microsecond,
|
||||
1_704_067_200_123_456,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
pub async fn test_otlp_metrics_into_seconds_physical_table_on_standalone() {
|
||||
let standalone = GreptimeDbStandaloneBuilder::new("test_otlp_s_physical")
|
||||
.build()
|
||||
.await;
|
||||
|
||||
// Narrowing truncates: 1704067200123456789ns -> 1704067200s.
|
||||
test_otlp_metrics_into_non_millisecond_physical_table(
|
||||
standalone.fe_instance(),
|
||||
"TIMESTAMP(0)",
|
||||
TimeUnit::Second,
|
||||
1_704_067_200,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
pub async fn test_otlp_metrics_into_nanoseconds_physical_table_on_standalone() {
|
||||
let standalone = GreptimeDbStandaloneBuilder::new("test_otlp_ns_physical")
|
||||
.build()
|
||||
.await;
|
||||
|
||||
// Same unit: the full nanosecond precision is kept verbatim.
|
||||
test_otlp_metrics_into_non_millisecond_physical_table(
|
||||
standalone.fe_instance(),
|
||||
"TIMESTAMP(9)",
|
||||
TimeUnit::Nanosecond,
|
||||
1_704_067_200_123_456_789,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
/// Regression test for <https://github.com/GreptimeTeam/greptimedb/issues/9231>:
|
||||
/// a physical metric table pre-created with a non-millisecond time index
|
||||
/// accepts OTLP ingestion; the nanosecond samples are converted to the
|
||||
/// physical table's unit (kept verbatim for a nanosecond table,
|
||||
/// truncating the sub-unit part for micro/second tables).
|
||||
async fn test_otlp_metrics_into_non_millisecond_physical_table(
|
||||
instance: &Arc<Instance>,
|
||||
sql_ts_type: &str,
|
||||
expected_unit: TimeUnit,
|
||||
expected_value: i64,
|
||||
) {
|
||||
let db = "otlp_non_ms_physical";
|
||||
let mut ctx = QueryContext::with(DEFAULT_CATALOG_NAME, db);
|
||||
// Route the request to the metric engine, like the OTLP HTTP handler
|
||||
// does with the default `prom_store.with_metric_engine = true`.
|
||||
ctx.set_protocol_ctx(ProtocolCtx::OtlpMetric(OtlpMetricCtx {
|
||||
with_metric_engine: true,
|
||||
..Default::default()
|
||||
}));
|
||||
let ctx = Arc::new(ctx);
|
||||
assert!(
|
||||
SqlQueryHandler::do_query(
|
||||
instance.as_ref(),
|
||||
&format!("CREATE DATABASE IF NOT EXISTS {db}"),
|
||||
ctx.clone(),
|
||||
)
|
||||
.await
|
||||
.first()
|
||||
.unwrap()
|
||||
.is_ok()
|
||||
);
|
||||
|
||||
let mut output = instance
|
||||
.do_query(
|
||||
&format!(
|
||||
"CREATE TABLE greptime_physical_table (\
|
||||
greptime_timestamp {sql_ts_type} NOT NULL, \
|
||||
greptime_value DOUBLE NULL, \
|
||||
TIME INDEX (greptime_timestamp)) \
|
||||
ENGINE = metric WITH ('physical_metric_table' = 'true')"
|
||||
),
|
||||
ctx.clone(),
|
||||
)
|
||||
.await;
|
||||
assert!(output.remove(0).is_ok());
|
||||
|
||||
let request = ExportMetricsServiceRequest {
|
||||
resource_metrics: vec![ResourceMetrics {
|
||||
scope_metrics: vec![ScopeMetrics {
|
||||
metrics: vec![Metric {
|
||||
name: "my_gauge".to_string(),
|
||||
data: Some(metric::Data::Gauge(Gauge {
|
||||
data_points: vec![NumberDataPoint {
|
||||
attributes: vec![keyvalue("host", "h1")],
|
||||
time_unix_nano: 1_704_067_200_123_456_789,
|
||||
value: Some(Value::AsDouble(1.0)),
|
||||
..Default::default()
|
||||
}],
|
||||
})),
|
||||
..Default::default()
|
||||
}],
|
||||
..Default::default()
|
||||
}],
|
||||
..Default::default()
|
||||
}],
|
||||
};
|
||||
instance.metrics(request, ctx.clone()).await.unwrap();
|
||||
|
||||
let mut output = instance
|
||||
.do_query("SELECT greptime_timestamp FROM my_gauge", ctx.clone())
|
||||
.await;
|
||||
let OutputData::Stream(stream) = output.remove(0).unwrap().data else {
|
||||
unreachable!()
|
||||
};
|
||||
let batches = RecordBatches::try_collect(stream).await.unwrap().take();
|
||||
assert_eq!(batches[0].num_rows(), 1);
|
||||
// The logical table's time index keeps the physical table's unit.
|
||||
let ts_column = batches[0].column(0);
|
||||
assert_eq!(
|
||||
ts_column.data_type(),
|
||||
&DataType::Timestamp(expected_unit, None),
|
||||
"unexpected time index type"
|
||||
);
|
||||
let stored = match expected_unit {
|
||||
TimeUnit::Second => ts_column.as_primitive::<TimestampSecondType>().value(0),
|
||||
TimeUnit::Millisecond => ts_column
|
||||
.as_primitive::<TimestampMillisecondType>()
|
||||
.value(0),
|
||||
TimeUnit::Microsecond => ts_column
|
||||
.as_primitive::<TimestampMicrosecondType>()
|
||||
.value(0),
|
||||
TimeUnit::Nanosecond => ts_column.as_primitive::<TimestampNanosecondType>().value(0),
|
||||
};
|
||||
assert_eq!(stored, expected_value);
|
||||
}
|
||||
|
||||
fn build_request() -> ExportMetricsServiceRequest {
|
||||
let data_points = vec![
|
||||
NumberDataPoint {
|
||||
|
||||
@@ -18,9 +18,19 @@ mod tests {
|
||||
|
||||
use api::prom_store::remote::label_matcher::Type as MatcherType;
|
||||
use api::prom_store::remote::{
|
||||
Label, LabelMatcher, Query, ReadRequest, ReadResponse, Sample, WriteRequest,
|
||||
Label, LabelMatcher, Query, ReadRequest, ReadResponse, Sample, TimeSeries, WriteRequest,
|
||||
};
|
||||
use api::v1::value::ValueData;
|
||||
use api::v1::{
|
||||
ColumnDataType, Row, RowInsertRequest, RowInsertRequests, Rows, SemanticType, Value,
|
||||
};
|
||||
use client::OutputData;
|
||||
use common_catalog::consts::DEFAULT_CATALOG_NAME;
|
||||
use datatypes::arrow::array::AsArray;
|
||||
use datatypes::arrow::datatypes::{
|
||||
DataType, TimeUnit, TimestampMicrosecondType, TimestampMillisecondType,
|
||||
TimestampNanosecondType, TimestampSecondType,
|
||||
};
|
||||
use frontend::instance::Instance;
|
||||
use prost::Message;
|
||||
use servers::http::prom_store::PHYSICAL_TABLE_PARAM;
|
||||
@@ -83,6 +93,432 @@ mod tests {
|
||||
.await;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn test_standalone_prom_store_remote_rw_microsecond_physical_table() {
|
||||
common_telemetry::init_default_ut_logging();
|
||||
let standalone =
|
||||
GreptimeDbStandaloneBuilder::new("test_prom_store_remote_rw_us_physical_table")
|
||||
.build()
|
||||
.await;
|
||||
let instance = standalone.fe_instance();
|
||||
|
||||
test_prom_store_remote_rw_non_millisecond_physical_table(
|
||||
instance,
|
||||
"TIMESTAMP(6)",
|
||||
TimeUnit::Microsecond,
|
||||
1_000_000,
|
||||
2_000_000,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn test_distributed_prom_store_remote_rw_microsecond_physical_table() {
|
||||
common_telemetry::init_default_ut_logging();
|
||||
let distributed =
|
||||
tests::create_distributed_instance("test_prom_store_remote_rw_us_physical_table").await;
|
||||
test_prom_store_remote_rw_non_millisecond_physical_table(
|
||||
&distributed.frontend(),
|
||||
"TIMESTAMP(6)",
|
||||
TimeUnit::Microsecond,
|
||||
1_000_000,
|
||||
2_000_000,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn test_standalone_prom_store_remote_rw_seconds_physical_table() {
|
||||
common_telemetry::init_default_ut_logging();
|
||||
let standalone =
|
||||
GreptimeDbStandaloneBuilder::new("test_prom_store_remote_rw_s_physical_table")
|
||||
.build()
|
||||
.await;
|
||||
let instance = standalone.fe_instance();
|
||||
|
||||
// Narrowing truncates: 1000ms/2000ms -> 1s/2s.
|
||||
test_prom_store_remote_rw_non_millisecond_physical_table(
|
||||
instance,
|
||||
"TIMESTAMP(0)",
|
||||
TimeUnit::Second,
|
||||
1,
|
||||
2,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn test_standalone_prom_store_remote_rw_nanoseconds_physical_table() {
|
||||
common_telemetry::init_default_ut_logging();
|
||||
let standalone =
|
||||
GreptimeDbStandaloneBuilder::new("test_prom_store_remote_rw_ns_physical_table")
|
||||
.build()
|
||||
.await;
|
||||
let instance = standalone.fe_instance();
|
||||
|
||||
// Lossless widening: 1000ms/2000ms -> 1e9ns/2e9ns.
|
||||
test_prom_store_remote_rw_non_millisecond_physical_table(
|
||||
instance,
|
||||
"TIMESTAMP(9)",
|
||||
TimeUnit::Nanosecond,
|
||||
1_000_000_000,
|
||||
2_000_000_000,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
/// Regression test for <https://github.com/GreptimeTeam/greptimedb/issues/9231>:
|
||||
/// prometheus remote write/read against a physical metric table
|
||||
/// pre-created with a non-millisecond time index. Millisecond samples are
|
||||
/// converted to the physical table's unit on write (lossless widening,
|
||||
/// truncating narrowing) and narrowed back to milliseconds on remote read.
|
||||
async fn test_prom_store_remote_rw_non_millisecond_physical_table(
|
||||
instance: &Arc<Instance>,
|
||||
sql_ts_type: &str,
|
||||
expected_unit: TimeUnit,
|
||||
expected_first: i64,
|
||||
expected_second: i64,
|
||||
) {
|
||||
let db = "prometheus_non_ms";
|
||||
let ctx = Arc::new(QueryContext::with(DEFAULT_CATALOG_NAME, db));
|
||||
assert!(
|
||||
SqlQueryHandler::do_query(
|
||||
instance.as_ref(),
|
||||
&format!("CREATE DATABASE IF NOT EXISTS {db}"),
|
||||
ctx.clone(),
|
||||
)
|
||||
.await
|
||||
.first()
|
||||
.unwrap()
|
||||
.is_ok()
|
||||
);
|
||||
|
||||
let mut output = instance
|
||||
.do_query(
|
||||
&format!(
|
||||
"CREATE TABLE greptime_physical_table (\
|
||||
greptime_timestamp {sql_ts_type} NOT NULL, \
|
||||
greptime_value DOUBLE NULL, \
|
||||
TIME INDEX (greptime_timestamp)) \
|
||||
ENGINE = metric WITH ('physical_metric_table' = 'true')"
|
||||
),
|
||||
ctx.clone(),
|
||||
)
|
||||
.await;
|
||||
assert!(output.remove(0).is_ok());
|
||||
|
||||
let write_request = WriteRequest {
|
||||
timeseries: vec![prom_store::mock_timeseries()[0].clone()],
|
||||
..Default::default()
|
||||
};
|
||||
let (row_inserts, _) = to_grpc_row_insert_requests(&write_request).unwrap();
|
||||
instance
|
||||
.write(row_inserts, ctx.clone(), true)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let read_request = ReadRequest {
|
||||
queries: vec![Query {
|
||||
start_timestamp_ms: 1000,
|
||||
end_timestamp_ms: 2000,
|
||||
matchers: vec![LabelMatcher {
|
||||
name: prom_store::METRIC_NAME_LABEL.to_string(),
|
||||
value: "metric1".to_string(),
|
||||
r#type: 0,
|
||||
}],
|
||||
..Default::default()
|
||||
}],
|
||||
..Default::default()
|
||||
};
|
||||
let resp = instance.read(read_request, ctx.clone()).await.unwrap();
|
||||
let body = prom_store::snappy_decompress(&resp.body).unwrap();
|
||||
let read_response = ReadResponse::decode(&body[..]).unwrap();
|
||||
assert_eq!(1, read_response.results.len());
|
||||
assert_eq!(1, read_response.results[0].timeseries.len());
|
||||
let timeseries = &read_response.results[0].timeseries[0];
|
||||
assert_eq!(
|
||||
timeseries.samples,
|
||||
vec![
|
||||
Sample {
|
||||
value: 1.0,
|
||||
timestamp: 1000,
|
||||
},
|
||||
Sample {
|
||||
value: 2.0,
|
||||
timestamp: 2000,
|
||||
},
|
||||
]
|
||||
);
|
||||
|
||||
// The stored timestamps keep the physical table's time index unit.
|
||||
let mut output = instance
|
||||
.do_query("SELECT greptime_timestamp FROM metric1", ctx.clone())
|
||||
.await;
|
||||
let OutputData::Stream(stream) = output.remove(0).unwrap().data else {
|
||||
unreachable!()
|
||||
};
|
||||
let batches = common_recordbatch::RecordBatches::try_collect(stream)
|
||||
.await
|
||||
.unwrap()
|
||||
.take();
|
||||
assert_eq!(batches[0].num_rows(), 2);
|
||||
let ts_column = batches[0].column(0);
|
||||
assert_eq!(
|
||||
ts_column.data_type(),
|
||||
&DataType::Timestamp(expected_unit, None),
|
||||
"unexpected time index type"
|
||||
);
|
||||
let stored = |row: usize| match expected_unit {
|
||||
TimeUnit::Second => ts_column.as_primitive::<TimestampSecondType>().value(row),
|
||||
TimeUnit::Millisecond => ts_column
|
||||
.as_primitive::<TimestampMillisecondType>()
|
||||
.value(row),
|
||||
TimeUnit::Microsecond => ts_column
|
||||
.as_primitive::<TimestampMicrosecondType>()
|
||||
.value(row),
|
||||
TimeUnit::Nanosecond => ts_column
|
||||
.as_primitive::<TimestampNanosecondType>()
|
||||
.value(row),
|
||||
};
|
||||
assert_eq!(stored(0), expected_first);
|
||||
assert_eq!(stored(1), expected_second);
|
||||
|
||||
// Negative, non-aligned timestamps must floor towards negative
|
||||
// infinity on remote read, matching `Timestamp::convert_to` on the
|
||||
// ingestion path: -1001us/-1001000ns -> -2ms, not -1ms.
|
||||
if matches!(expected_unit, TimeUnit::Microsecond | TimeUnit::Nanosecond) {
|
||||
let (ts_datatype, negative_value) = if expected_unit == TimeUnit::Microsecond {
|
||||
(
|
||||
ColumnDataType::TimestampMicrosecond,
|
||||
ValueData::TimestampMicrosecondValue(-1001),
|
||||
)
|
||||
} else {
|
||||
(
|
||||
ColumnDataType::TimestampNanosecond,
|
||||
ValueData::TimestampNanosecondValue(-1_001_000),
|
||||
)
|
||||
};
|
||||
let negative_request = RowInsertRequests {
|
||||
inserts: vec![RowInsertRequest {
|
||||
table_name: "metric1".to_string(),
|
||||
rows: Some(Rows {
|
||||
schema: vec![
|
||||
api::v1::ColumnSchema {
|
||||
column_name: "greptime_timestamp".to_string(),
|
||||
datatype: ts_datatype as i32,
|
||||
semantic_type: SemanticType::Timestamp as i32,
|
||||
datatype_extension: None,
|
||||
options: None,
|
||||
},
|
||||
api::v1::ColumnSchema {
|
||||
column_name: "greptime_value".to_string(),
|
||||
datatype: ColumnDataType::Float64 as i32,
|
||||
semantic_type: SemanticType::Field as i32,
|
||||
datatype_extension: None,
|
||||
options: None,
|
||||
},
|
||||
api::v1::ColumnSchema {
|
||||
column_name: "job".to_string(),
|
||||
datatype: ColumnDataType::String as i32,
|
||||
semantic_type: SemanticType::Tag as i32,
|
||||
datatype_extension: None,
|
||||
options: None,
|
||||
},
|
||||
],
|
||||
rows: vec![Row {
|
||||
values: vec![
|
||||
Value {
|
||||
value_data: Some(negative_value),
|
||||
},
|
||||
Value {
|
||||
value_data: Some(ValueData::F64Value(3.0)),
|
||||
},
|
||||
Value {
|
||||
value_data: Some(ValueData::StringValue("spark".to_string())),
|
||||
},
|
||||
],
|
||||
}],
|
||||
}),
|
||||
}],
|
||||
};
|
||||
instance
|
||||
.write(negative_request, ctx.clone(), true)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let read_request = ReadRequest {
|
||||
queries: vec![Query {
|
||||
start_timestamp_ms: -2,
|
||||
end_timestamp_ms: 0,
|
||||
matchers: vec![LabelMatcher {
|
||||
name: prom_store::METRIC_NAME_LABEL.to_string(),
|
||||
value: "metric1".to_string(),
|
||||
r#type: 0,
|
||||
}],
|
||||
..Default::default()
|
||||
}],
|
||||
..Default::default()
|
||||
};
|
||||
let resp = instance.read(read_request, ctx.clone()).await.unwrap();
|
||||
let body = prom_store::snappy_decompress(&resp.body).unwrap();
|
||||
let read_response = ReadResponse::decode(&body[..]).unwrap();
|
||||
assert_eq!(1, read_response.results.len());
|
||||
assert_eq!(
|
||||
read_response.results[0].timeseries[0].samples,
|
||||
vec![Sample {
|
||||
value: 3.0,
|
||||
timestamp: -2,
|
||||
}],
|
||||
"remote read must floor negative non-aligned timestamps"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn test_standalone_prom_store_write_existing_logical_table_of_other_physical_unit() {
|
||||
common_telemetry::init_default_ut_logging();
|
||||
let standalone = GreptimeDbStandaloneBuilder::new(
|
||||
"test_prom_store_write_existing_logical_other_physical_unit",
|
||||
)
|
||||
.build()
|
||||
.await;
|
||||
let instance = standalone.fe_instance();
|
||||
|
||||
// Two physical metric tables with different time index units.
|
||||
let db = "prometheus_mixed_units";
|
||||
let ctx = Arc::new(QueryContext::with(DEFAULT_CATALOG_NAME, db));
|
||||
assert!(
|
||||
SqlQueryHandler::do_query(
|
||||
instance.as_ref(),
|
||||
&format!("CREATE DATABASE IF NOT EXISTS {db}"),
|
||||
ctx.clone(),
|
||||
)
|
||||
.await
|
||||
.first()
|
||||
.unwrap()
|
||||
.is_ok()
|
||||
);
|
||||
let mut output = instance
|
||||
.do_query(
|
||||
"CREATE TABLE phy_us (greptime_timestamp TIMESTAMP(6) NOT NULL, \
|
||||
greptime_value DOUBLE NULL, TIME INDEX (greptime_timestamp)) \
|
||||
ENGINE = metric WITH ('physical_metric_table' = 'true')",
|
||||
ctx.clone(),
|
||||
)
|
||||
.await;
|
||||
assert!(output.remove(0).is_ok());
|
||||
|
||||
// First create the logical table `shared_metric` on the default
|
||||
// (millisecond) physical table.
|
||||
let shared_series = TimeSeries {
|
||||
labels: vec![
|
||||
Label {
|
||||
name: prom_store::METRIC_NAME_LABEL.to_string(),
|
||||
value: "shared_metric".to_string(),
|
||||
},
|
||||
Label {
|
||||
name: "job".to_string(),
|
||||
value: "demo".to_string(),
|
||||
},
|
||||
],
|
||||
samples: vec![Sample {
|
||||
value: 1.0,
|
||||
timestamp: 1000,
|
||||
}],
|
||||
..Default::default()
|
||||
};
|
||||
let (row_inserts, _) = to_grpc_row_insert_requests(&WriteRequest {
|
||||
timeseries: vec![shared_series.clone()],
|
||||
..Default::default()
|
||||
})
|
||||
.unwrap();
|
||||
instance
|
||||
.write(row_inserts, ctx.clone(), true)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
// Now select the microsecond physical table while writing to BOTH the
|
||||
// existing millisecond `shared_metric` (bound to the default physical
|
||||
// table) and a new table: the existing table's request must keep the
|
||||
// millisecond unit, and only the new table uses microsecond.
|
||||
let mut hint_ctx = QueryContext::with(DEFAULT_CATALOG_NAME, db);
|
||||
hint_ctx.set_extension(PHYSICAL_TABLE_PARAM, "phy_us".to_string());
|
||||
let hint_ctx = Arc::new(hint_ctx);
|
||||
let (row_inserts, _) = to_grpc_row_insert_requests(&WriteRequest {
|
||||
timeseries: vec![
|
||||
TimeSeries {
|
||||
samples: vec![Sample {
|
||||
value: 1.5,
|
||||
timestamp: 1500,
|
||||
}],
|
||||
..shared_series
|
||||
},
|
||||
TimeSeries {
|
||||
labels: vec![
|
||||
Label {
|
||||
name: prom_store::METRIC_NAME_LABEL.to_string(),
|
||||
value: "fresh_us_metric".to_string(),
|
||||
},
|
||||
Label {
|
||||
name: "job".to_string(),
|
||||
value: "demo".to_string(),
|
||||
},
|
||||
],
|
||||
samples: vec![Sample {
|
||||
value: 2.5,
|
||||
timestamp: 1500,
|
||||
}],
|
||||
..Default::default()
|
||||
},
|
||||
],
|
||||
..Default::default()
|
||||
})
|
||||
.unwrap();
|
||||
instance
|
||||
.write(row_inserts, hint_ctx.clone(), true)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
// The existing table keeps its millisecond unit and precision.
|
||||
let mut output = instance
|
||||
.do_query(
|
||||
"SELECT greptime_timestamp FROM shared_metric ORDER BY greptime_timestamp",
|
||||
ctx.clone(),
|
||||
)
|
||||
.await;
|
||||
let OutputData::Stream(stream) = output.remove(0).unwrap().data else {
|
||||
unreachable!()
|
||||
};
|
||||
let batches = common_recordbatch::RecordBatches::try_collect(stream)
|
||||
.await
|
||||
.unwrap()
|
||||
.take();
|
||||
let shared_ts = batches[0]
|
||||
.column(0)
|
||||
.as_primitive::<TimestampMillisecondType>();
|
||||
assert_eq!((shared_ts.value(0), shared_ts.value(1)), (1000, 1500));
|
||||
|
||||
// The new table is created on the selected physical table with the
|
||||
// microsecond unit.
|
||||
let mut output = instance
|
||||
.do_query(
|
||||
"SELECT greptime_timestamp FROM fresh_us_metric",
|
||||
ctx.clone(),
|
||||
)
|
||||
.await;
|
||||
let OutputData::Stream(stream) = output.remove(0).unwrap().data else {
|
||||
unreachable!()
|
||||
};
|
||||
let batches = common_recordbatch::RecordBatches::try_collect(stream)
|
||||
.await
|
||||
.unwrap()
|
||||
.take();
|
||||
let fresh_ts = batches[0]
|
||||
.column(0)
|
||||
.as_primitive::<TimestampMicrosecondType>();
|
||||
assert_eq!(fresh_ts.value(0), 1_500_000);
|
||||
}
|
||||
|
||||
async fn test_prom_store_remote_rw(instance: &Arc<Instance>, physical_table: Option<String>) {
|
||||
let write_request = WriteRequest {
|
||||
timeseries: prom_store::mock_timeseries(),
|
||||
|
||||
@@ -15,18 +15,20 @@
|
||||
use std::collections::BTreeMap;
|
||||
use std::io::Write;
|
||||
use std::str::FromStr;
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
|
||||
use api::greptime_proto::io::prometheus::write::v2::histogram::{Count, ZeroCount};
|
||||
use api::greptime_proto::io::prometheus::write::v2::metadata::MetricType as RemoteWriteV2MetricType;
|
||||
use api::greptime_proto::io::prometheus::write::v2::{
|
||||
BucketSpan, Histogram, Metadata as RemoteWriteV2Metadata, Sample as RemoteWriteV2Sample,
|
||||
TimeSeries as RemoteWriteV2TimeSeries,
|
||||
BucketSpan, Histogram, Metadata as RemoteWriteV2Metadata, Request as RemoteWriteV2Request,
|
||||
Sample as RemoteWriteV2Sample, TimeSeries as RemoteWriteV2TimeSeries,
|
||||
};
|
||||
use api::prom_store::remote::label_matcher::Type as MatcherType;
|
||||
use api::prom_store::remote::{
|
||||
Label, LabelMatcher, Query, ReadRequest, ReadResponse, Sample, TimeSeries, WriteRequest,
|
||||
};
|
||||
use api::v1::RowInsertRequests;
|
||||
use auth::{UserProviderRef, user_provider_from_option};
|
||||
use axum::http::{HeaderMap, HeaderName, HeaderValue, StatusCode};
|
||||
use base64::prelude::{BASE64_STANDARD, Engine as _};
|
||||
@@ -69,7 +71,9 @@ use servers::http::header::constants::{
|
||||
};
|
||||
use servers::http::header::{GREPTIME_DB_HEADER_NAME, GREPTIME_TIMEZONE_HEADER_NAME};
|
||||
use servers::http::otlp::GoogleRpcStatus;
|
||||
use servers::http::prometheus::{Column, PrometheusJsonResponse, PrometheusResponse};
|
||||
use servers::http::prometheus::{
|
||||
Column, PromQueryResult, PrometheusJsonResponse, PrometheusResponse,
|
||||
};
|
||||
use servers::http::result::error_result::ErrorResponse;
|
||||
use servers::http::result::greptime_result_v1::GreptimedbV1Response;
|
||||
use servers::http::result::influxdb_result_v1::{InfluxdbOutput, InfluxdbV1Response};
|
||||
@@ -79,11 +83,13 @@ use servers::prom_remote_write::v2::test_util as remote_write_v2;
|
||||
use servers::prom_remote_write::validation::PromValidationMode;
|
||||
use servers::prom_store::{self, mock_timeseries_new_label};
|
||||
use servers::request_memory_limiter::ServerMemoryLimiter;
|
||||
use session::context::QueryContextRef;
|
||||
use standalone::options::StandaloneOptions;
|
||||
use table::table_name::TableName;
|
||||
use tests_integration::test_util::{
|
||||
MockInstanceImpl, StorageType, assert_wal_delta, build_test_prom_server, setup_test_http_app,
|
||||
setup_test_http_app_with_frontend, setup_test_http_app_with_frontend_and_slow_query_threshold,
|
||||
MockInstanceImpl, StorageType, TestGuard, assert_wal_delta, build_test_prom_server,
|
||||
setup_test_http_app, setup_test_http_app_with_frontend,
|
||||
setup_test_http_app_with_frontend_and_slow_query_threshold,
|
||||
setup_test_http_app_with_frontend_and_user_provider, setup_test_prom_app_with_frontend,
|
||||
setup_test_prom_app_with_frontend_batched,
|
||||
};
|
||||
@@ -157,6 +163,7 @@ macro_rules! http_tests {
|
||||
test_http_analyze_stream_tql,
|
||||
test_http_sql_slow_query,
|
||||
test_prometheus_promql_api,
|
||||
test_promql_over_non_millisecond_physical_tables,
|
||||
test_prometheus_label_replace_response,
|
||||
test_prom_http_api,
|
||||
test_config_api,
|
||||
@@ -165,6 +172,9 @@ macro_rules! http_tests {
|
||||
test_prometheus_remote_write_v2,
|
||||
test_prometheus_remote_write_v2_native_histogram,
|
||||
test_prometheus_remote_write_batched,
|
||||
test_prometheus_remote_write_batched_mixed_time_index_units,
|
||||
test_prometheus_remote_write_batched_interceptor_time_index_units,
|
||||
test_prometheus_remote_write_v2_batched_interceptor_time_index_units,
|
||||
test_prometheus_remote_special_labels,
|
||||
test_prometheus_remote_schema_labels,
|
||||
test_prometheus_remote_write_with_pipeline,
|
||||
@@ -1713,6 +1723,197 @@ pub async fn test_prom_http_api(store_type: StorageType) {
|
||||
guard.remove_all().await;
|
||||
}
|
||||
|
||||
/// PromQL must behave identically regardless of the physical metric table's
|
||||
/// time index unit: the same samples are remote-written into physical tables
|
||||
/// pre-created with second/micro/nano time indexes (plus a millisecond
|
||||
/// baseline), and every PromQL response must match the baseline exactly.
|
||||
/// Regression test for <https://github.com/GreptimeTeam/greptimedb/issues/9231>.
|
||||
pub async fn test_promql_over_non_millisecond_physical_tables(store_type: StorageType) {
|
||||
common_telemetry::init_default_ut_logging();
|
||||
let (app, mut guard) =
|
||||
setup_test_prom_app_with_frontend(store_type, "promql_non_ms_units").await;
|
||||
let client = TestClient::new(app).await;
|
||||
|
||||
let units = [
|
||||
("ms", "timestamp(3)"),
|
||||
("second", "timestamp(0)"),
|
||||
("micro", "timestamp(6)"),
|
||||
("nano", "timestamp(9)"),
|
||||
];
|
||||
|
||||
for (suffix, ts_type) in units {
|
||||
let db = format!("promql_units_{suffix}");
|
||||
let res = client
|
||||
.get(&format!("/v1/sql?db=public&sql=create database {db}"))
|
||||
.send()
|
||||
.await;
|
||||
assert_eq!(res.status(), StatusCode::OK, "create database {db}");
|
||||
|
||||
let res = client
|
||||
.get(&format!(
|
||||
"/v1/sql?db={db}&sql=CREATE TABLE greptime_physical_table \
|
||||
(greptime_timestamp {ts_type} NOT NULL, greptime_value DOUBLE NULL, \
|
||||
TIME INDEX (greptime_timestamp)) \
|
||||
ENGINE = metric WITH ('physical_metric_table' = 'true')"
|
||||
))
|
||||
.send()
|
||||
.await;
|
||||
assert_eq!(
|
||||
res.status(),
|
||||
StatusCode::OK,
|
||||
"create physical table with {ts_type} in {db}"
|
||||
);
|
||||
|
||||
// unit_gauge{job="demo"} = 1.0 @ 1000ms, 2.0 @ 2000ms
|
||||
let write_request = WriteRequest {
|
||||
timeseries: vec![TimeSeries {
|
||||
labels: vec![
|
||||
Label {
|
||||
name: prom_store::METRIC_NAME_LABEL.to_string(),
|
||||
value: "unit_gauge".to_string(),
|
||||
},
|
||||
Label {
|
||||
name: "job".to_string(),
|
||||
value: "demo".to_string(),
|
||||
},
|
||||
],
|
||||
samples: vec![
|
||||
Sample {
|
||||
value: 1.0,
|
||||
timestamp: 1000,
|
||||
},
|
||||
Sample {
|
||||
value: 2.0,
|
||||
timestamp: 2000,
|
||||
},
|
||||
],
|
||||
..Default::default()
|
||||
}],
|
||||
..Default::default()
|
||||
};
|
||||
let compressed = prom_store::snappy_compress(&write_request.encode_to_vec()).unwrap();
|
||||
let write_url = format!("/v1/prometheus/write?db={db}");
|
||||
let res = client
|
||||
.post(write_url.as_str())
|
||||
.header("Content-Encoding", "snappy")
|
||||
.body(compressed)
|
||||
.send()
|
||||
.await;
|
||||
assert_eq!(
|
||||
res.status(),
|
||||
StatusCode::NO_CONTENT,
|
||||
"remote write into {db}"
|
||||
);
|
||||
}
|
||||
|
||||
let query_paths = [
|
||||
"query?query=unit_gauge&time=2",
|
||||
"query_range?query=unit_gauge&start=0&end=5&step=1",
|
||||
"query?query=avg_over_time(unit_gauge[1m])&time=2",
|
||||
"query?query=rate(unit_gauge[1m])&time=2",
|
||||
];
|
||||
for path in query_paths {
|
||||
let mut baseline: Option<(String, PrometheusResponse)> = None;
|
||||
for (suffix, _) in units {
|
||||
let db = format!("promql_units_{suffix}");
|
||||
let res = client
|
||||
.get(&format!("/v1/prometheus/api/v1/{path}&db={db}"))
|
||||
.send()
|
||||
.await;
|
||||
assert_eq!(res.status(), StatusCode::OK, "promql `{path}` on {db}");
|
||||
let body = res.json::<PrometheusJsonResponse>().await;
|
||||
assert_eq!(body.status, "success", "promql `{path}` on {db}");
|
||||
match &baseline {
|
||||
None => baseline = Some((suffix.to_string(), body.data)),
|
||||
Some((baseline_suffix, baseline_data)) => {
|
||||
assert_eq!(
|
||||
&body.data, baseline_data,
|
||||
"promql `{path}` on {db} differs from the {baseline_suffix} baseline"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Pin the baseline responses so the cross-unit equality above cannot pass
|
||||
// vacuously on wrong data: the samples are 1.0@1s and 2.0@2s, and
|
||||
// `rate(unit_gauge[1m])@2s` follows Prometheus's extrapolated rate:
|
||||
// Δv=1 over a 1s sampled interval extrapolated to 1.5s, divided by the
|
||||
// 60s window -> 1.5/60 = 0.025 per second.
|
||||
let res = client
|
||||
.get("/v1/prometheus/api/v1/query?db=promql_units_ms&query=unit_gauge&time=2")
|
||||
.send()
|
||||
.await;
|
||||
assert_eq!(res.status(), StatusCode::OK);
|
||||
let body = res.json::<PrometheusJsonResponse>().await;
|
||||
assert_eq!(
|
||||
body.data,
|
||||
serde_json::from_value::<PrometheusResponse>(json!({
|
||||
"resultType": "vector",
|
||||
"result": [{
|
||||
"metric": {"__name__": "unit_gauge", "job": "demo"},
|
||||
"value": [2.0, "2.0"]
|
||||
}]
|
||||
}))
|
||||
.unwrap()
|
||||
);
|
||||
|
||||
let res = client
|
||||
.get("/v1/prometheus/api/v1/query_range?db=promql_units_ms&query=unit_gauge&start=0&end=5&step=1")
|
||||
.send()
|
||||
.await;
|
||||
assert_eq!(res.status(), StatusCode::OK);
|
||||
let body = res.json::<PrometheusJsonResponse>().await;
|
||||
assert_eq!(
|
||||
body.data,
|
||||
serde_json::from_value::<PrometheusResponse>(json!({
|
||||
"resultType": "matrix",
|
||||
"result": [{
|
||||
"metric": {"__name__": "unit_gauge", "job": "demo"},
|
||||
"values": [
|
||||
[1.0, "1.0"], [2.0, "2.0"], [3.0, "2.0"], [4.0, "2.0"], [5.0, "2.0"]
|
||||
]
|
||||
}]
|
||||
}))
|
||||
.unwrap()
|
||||
);
|
||||
|
||||
let res = client
|
||||
.get("/v1/prometheus/api/v1/query?db=promql_units_ms&query=avg_over_time(unit_gauge[1m])&time=2")
|
||||
.send()
|
||||
.await;
|
||||
assert_eq!(res.status(), StatusCode::OK);
|
||||
let body = res.json::<PrometheusJsonResponse>().await;
|
||||
assert_eq!(
|
||||
body.data,
|
||||
serde_json::from_value::<PrometheusResponse>(json!({
|
||||
"resultType": "vector",
|
||||
"result": [{
|
||||
"metric": {"__name__": "unit_gauge", "job": "demo"},
|
||||
"value": [2.0, "1.5"]
|
||||
}]
|
||||
}))
|
||||
.unwrap()
|
||||
);
|
||||
|
||||
let res = client
|
||||
.get("/v1/prometheus/api/v1/query?db=promql_units_ms&query=rate(unit_gauge[1m])&time=2")
|
||||
.send()
|
||||
.await;
|
||||
assert_eq!(res.status(), StatusCode::OK);
|
||||
let body = res.json::<PrometheusJsonResponse>().await;
|
||||
let rate = match body.data {
|
||||
PrometheusResponse::PromData(data) => match data.result {
|
||||
PromQueryResult::Vector(v) => v[0].value.as_ref().unwrap().1.clone(),
|
||||
other => panic!("expected vector, got {other:?}"),
|
||||
},
|
||||
other => panic!("expected prom data, got {other:?}"),
|
||||
};
|
||||
assert_eq!(rate, "0.025");
|
||||
|
||||
guard.remove_all().await;
|
||||
}
|
||||
|
||||
pub async fn test_metrics_api(store_type: StorageType) {
|
||||
common_telemetry::init_default_ut_logging();
|
||||
let (app, mut guard) = setup_test_http_app(store_type, "metrics_api").await;
|
||||
@@ -3593,6 +3794,277 @@ async fn write_prometheus_skip_wal_sample(
|
||||
assert_eq!(request.send().await.status(), StatusCode::NO_CONTENT);
|
||||
}
|
||||
|
||||
/// Batched remote write against a logical table bound to a non-millisecond
|
||||
/// physical table: the bulk guard must reject the destination (the bulk
|
||||
/// encode only produces millisecond batches) and fall back to the ordinary
|
||||
/// insert path, which converts the requests to the table's unit.
|
||||
pub async fn test_prometheus_remote_write_batched_mixed_time_index_units(store_type: StorageType) {
|
||||
common_telemetry::init_default_ut_logging();
|
||||
let (app, mut guard) =
|
||||
setup_test_prom_app_with_frontend_batched(store_type, "prom_rw_batched_mixed_units").await;
|
||||
let client = TestClient::new(app).await;
|
||||
|
||||
let res = client
|
||||
.get("/v1/sql?db=public&sql=CREATE TABLE phy_us (greptime_timestamp TIMESTAMP(6) NOT NULL, greptime_value DOUBLE NULL, TIME INDEX (greptime_timestamp)) ENGINE = metric WITH ('physical_metric_table' = 'true')")
|
||||
.send()
|
||||
.await;
|
||||
assert_eq!(res.status(), StatusCode::OK);
|
||||
|
||||
let write = |metric: &str, value: f64, timestamp: i64| {
|
||||
let write_request = WriteRequest {
|
||||
timeseries: vec![TimeSeries {
|
||||
labels: vec![
|
||||
Label {
|
||||
name: prom_store::METRIC_NAME_LABEL.to_string(),
|
||||
value: metric.to_string(),
|
||||
},
|
||||
Label {
|
||||
name: "job".to_string(),
|
||||
value: "demo".to_string(),
|
||||
},
|
||||
],
|
||||
samples: vec![Sample { value, timestamp }],
|
||||
..Default::default()
|
||||
}],
|
||||
..Default::default()
|
||||
};
|
||||
prom_store::snappy_compress(&write_request.encode_to_vec()).unwrap()
|
||||
};
|
||||
|
||||
// Create the logical table on the microsecond physical table. The bulk
|
||||
// guard rejects the microsecond selected physical table, so this write
|
||||
// takes the ordinary insert path.
|
||||
let res = client
|
||||
.post("/v1/prometheus/write?physical_table=phy_us")
|
||||
.header("Content-Encoding", "snappy")
|
||||
.body(write("us_metric", 2.5, 1500))
|
||||
.send()
|
||||
.await;
|
||||
assert_eq!(res.status(), StatusCode::NO_CONTENT);
|
||||
|
||||
// Write the same existing table again while selecting the default
|
||||
// (millisecond) physical table: the guard must reject the existing
|
||||
// microsecond destination and fall back to the ordinary insert path —
|
||||
// without the destination check the bulk path would build millisecond
|
||||
// arrays against the microsecond schema and fail the write.
|
||||
let res = client
|
||||
.post("/v1/prometheus/write")
|
||||
.header("Content-Encoding", "snappy")
|
||||
.body(write("us_metric", 3.5, 2000))
|
||||
.send()
|
||||
.await;
|
||||
assert_eq!(res.status(), StatusCode::NO_CONTENT);
|
||||
|
||||
// Both samples are stored on the microsecond time index.
|
||||
validate_data(
|
||||
"prom_rw_batched_mixed_units",
|
||||
&client,
|
||||
"SELECT COUNT(*), MAX(greptime_value) FROM us_metric",
|
||||
"[[2,3.5]]",
|
||||
)
|
||||
.await;
|
||||
|
||||
guard.remove_all().await;
|
||||
}
|
||||
|
||||
/// A Prometheus write interceptor that redirects every remote write to the
|
||||
/// `tenant_redirect` schema, like a per-tenant redirection: the incoming
|
||||
/// context targets `public` while the destination tables live elsewhere.
|
||||
struct PromSchemaRedirectInterceptor;
|
||||
|
||||
impl servers::interceptor::PromStoreProtocolInterceptor for PromSchemaRedirectInterceptor {
|
||||
type Error = servers::error::Error;
|
||||
|
||||
fn pre_write(
|
||||
&self,
|
||||
_write_req: &RowInsertRequests,
|
||||
ctx: QueryContextRef,
|
||||
) -> servers::error::Result<()> {
|
||||
ctx.set_current_schema("tenant_redirect");
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
async fn setup_redirecting_batched_prom_app(
|
||||
store_type: StorageType,
|
||||
name: &str,
|
||||
) -> (TestClient, TestGuard) {
|
||||
let plugins = Plugins::default();
|
||||
plugins.insert::<servers::interceptor::PromStoreProtocolInterceptorRef<servers::error::Error>>(
|
||||
Arc::new(PromSchemaRedirectInterceptor),
|
||||
);
|
||||
let standalone = tests_integration::standalone::GreptimeDbStandaloneBuilder::new(name)
|
||||
.with_default_store_type(store_type)
|
||||
.with_plugin(plugins)
|
||||
.build()
|
||||
.await;
|
||||
let server = build_test_prom_server(standalone.fe_instance().clone(), true)
|
||||
.with_greptime_config_options(standalone.opts.datanode_options().to_toml().unwrap())
|
||||
.build();
|
||||
let client = TestClient::new(server.build(server.make_app()).unwrap()).await;
|
||||
(client, standalone.guard)
|
||||
}
|
||||
|
||||
async fn create_redirect_schema_with_microsecond_table(client: &TestClient, table: &str) {
|
||||
for sql in [
|
||||
"create database if not exists tenant_redirect",
|
||||
"CREATE TABLE tenant_redirect.phy_us \
|
||||
(ts timestamp(6) time index, val double, host string primary key) \
|
||||
engine=metric with ('physical_metric_table' = 'true')",
|
||||
&format!(
|
||||
"CREATE TABLE tenant_redirect.{table} \
|
||||
(ts timestamp(6) time index, val double, host string primary key) \
|
||||
engine=metric with ('on_physical_table' = 'phy_us')"
|
||||
),
|
||||
] {
|
||||
let res = client
|
||||
.get(format!("/v1/sql?sql={sql}").as_str())
|
||||
.send()
|
||||
.await;
|
||||
assert_eq!(res.status(), StatusCode::OK, "setup: {sql}");
|
||||
}
|
||||
}
|
||||
|
||||
/// Regression test: with batching enabled and a pre_write interceptor that
|
||||
/// redirects the context to a schema holding same-named non-millisecond
|
||||
/// tables, the bulk eligibility must be evaluated against the redirected
|
||||
/// destinations (after preflight) and fall back to the ordinary insert path
|
||||
/// instead of failing in the bulk encode.
|
||||
pub async fn test_prometheus_remote_write_batched_interceptor_time_index_units(
|
||||
store_type: StorageType,
|
||||
) {
|
||||
common_telemetry::init_default_ut_logging();
|
||||
let (client, mut guard) =
|
||||
setup_redirecting_batched_prom_app(store_type, "prom_rw_batched_interceptor_units").await;
|
||||
create_redirect_schema_with_microsecond_table(&client, "intercept_metric").await;
|
||||
|
||||
let write_request = WriteRequest {
|
||||
timeseries: vec![TimeSeries {
|
||||
labels: vec![
|
||||
Label {
|
||||
name: prom_store::METRIC_NAME_LABEL.to_string(),
|
||||
value: "intercept_metric".to_string(),
|
||||
},
|
||||
Label {
|
||||
name: "job".to_string(),
|
||||
value: "demo".to_string(),
|
||||
},
|
||||
],
|
||||
samples: vec![Sample {
|
||||
value: 1.0,
|
||||
timestamp: 1000,
|
||||
}],
|
||||
..Default::default()
|
||||
}],
|
||||
..Default::default()
|
||||
};
|
||||
let compressed = prom_store::snappy_compress(&write_request.encode_to_vec()).unwrap();
|
||||
let res = client
|
||||
.post("/v1/prometheus/write")
|
||||
.header("Content-Encoding", "snappy")
|
||||
.body(compressed)
|
||||
.send()
|
||||
.await;
|
||||
assert_eq!(res.status(), StatusCode::NO_CONTENT);
|
||||
|
||||
validate_data(
|
||||
"prom_rw_batched_interceptor_units",
|
||||
&client,
|
||||
"SELECT COUNT(*) FROM tenant_redirect.intercept_metric",
|
||||
"[[1]]",
|
||||
)
|
||||
.await;
|
||||
|
||||
guard.remove_all().await;
|
||||
}
|
||||
|
||||
/// The v2 variant of the interceptor regression test: a mixed samples-plus-
|
||||
/// histograms request whose existing sample destination lives in the
|
||||
/// redirected schema with a microsecond time index must fall back to the
|
||||
/// ordinary insert path (both series land), not fail in the bulk encode.
|
||||
pub async fn test_prometheus_remote_write_v2_batched_interceptor_time_index_units(
|
||||
store_type: StorageType,
|
||||
) {
|
||||
common_telemetry::init_default_ut_logging();
|
||||
let (client, mut guard) =
|
||||
setup_redirecting_batched_prom_app(store_type, "prom_rw_v2_batched_interceptor_units")
|
||||
.await;
|
||||
create_redirect_schema_with_microsecond_table(&client, "v2_mixed_sample").await;
|
||||
|
||||
// One sample series targeting the existing microsecond table, plus one
|
||||
// histogram series for a new (millisecond) table.
|
||||
let mut symbols = vec![String::new()];
|
||||
let mut symbol = |value: &str| {
|
||||
symbols.push(value.to_string());
|
||||
(symbols.len() - 1) as u32
|
||||
};
|
||||
let name_ref = symbol("__name__");
|
||||
let sample_ref = symbol("v2_mixed_sample");
|
||||
let job_ref = symbol("job");
|
||||
let demo_ref = symbol("demo");
|
||||
let name_ref2 = symbol("__name__");
|
||||
let histo_ref = symbol("v2_mixed_histo");
|
||||
let request = RemoteWriteV2Request {
|
||||
symbols,
|
||||
timeseries: vec![
|
||||
RemoteWriteV2TimeSeries {
|
||||
labels_refs: vec![name_ref, sample_ref, job_ref, demo_ref],
|
||||
samples: vec![RemoteWriteV2Sample {
|
||||
value: 1.0,
|
||||
timestamp: 1000,
|
||||
start_timestamp: 0,
|
||||
}],
|
||||
..Default::default()
|
||||
},
|
||||
RemoteWriteV2TimeSeries {
|
||||
labels_refs: vec![name_ref2, histo_ref, job_ref, demo_ref],
|
||||
histograms: vec![Histogram {
|
||||
count: Some(Count::CountInt(1)),
|
||||
sum: 1.0,
|
||||
positive_spans: vec![BucketSpan {
|
||||
offset: 0,
|
||||
length: 1,
|
||||
}],
|
||||
positive_deltas: vec![1],
|
||||
timestamp: 1000,
|
||||
start_timestamp: 1000,
|
||||
..Default::default()
|
||||
}],
|
||||
..Default::default()
|
||||
},
|
||||
],
|
||||
};
|
||||
let compressed = prom_store::snappy_compress(&request.encode_to_vec()).unwrap();
|
||||
let res = client
|
||||
.post("/v1/prometheus/write")
|
||||
.header(
|
||||
"Content-Type",
|
||||
"application/x-protobuf;proto=io.prometheus.write.v2.Request",
|
||||
)
|
||||
.header("Content-Encoding", "snappy")
|
||||
.body(compressed)
|
||||
.send()
|
||||
.await;
|
||||
assert_eq!(res.status(), StatusCode::NO_CONTENT);
|
||||
|
||||
validate_data(
|
||||
"prom_rw_v2_batched_interceptor_units",
|
||||
&client,
|
||||
"SELECT COUNT(*) FROM tenant_redirect.v2_mixed_sample",
|
||||
"[[1]]",
|
||||
)
|
||||
.await;
|
||||
validate_data(
|
||||
"prom_rw_v2_batched_interceptor_units",
|
||||
&client,
|
||||
"SELECT COUNT(*) FROM tenant_redirect.v2_mixed_histo",
|
||||
"[[1]]",
|
||||
)
|
||||
.await;
|
||||
|
||||
guard.remove_all().await;
|
||||
}
|
||||
|
||||
/// Covers the batched (pending-rows-batcher) Prometheus remote write path, which
|
||||
/// bypasses `PromStoreProtocolHandler::write`. Verifies the metric table is created
|
||||
/// asynchronously and still carries the Prometheus semantic identity stamped on the
|
||||
|
||||
Reference in New Issue
Block a user