feat: enable native histogram ingestion by default (#9301)

* feat: enable native histogram ingestion by default

Signed-off-by: shuiyisong <xixing.sys@gmail.com>

* chore: add comments

Signed-off-by: shuiyisong <xixing.sys@gmail.com>

* fix: test

Signed-off-by: shuiyisong <xixing.sys@gmail.com>

---------

Signed-off-by: shuiyisong <xixing.sys@gmail.com>
This commit is contained in:
shuiyisong
2026-09-23 13:03:51 +00:00
committed by GitHub
parent 4c98fa5265
commit 20dde2601f
29 changed files with 429 additions and 466 deletions
-4
View File
@@ -93,14 +93,12 @@
| `jaeger.enable` | Bool | `true` | Whether to enable Jaeger protocol in HTTP API. |
| `otlp` | -- | -- | OpenTelemetry protocol options. |
| `otlp.enable` | Bool | `true` | Whether to enable OpenTelemetry protocol in HTTP API. |
| `otlp.experimental_enable_exponential_histogram` | Bool | `false` | Experimental: enable cumulative OTLP exponential histogram ingestion. |
| `otlp.trace_ingest_chunk_size` | Integer | `512` | Maximum spans per trace ingest chunk. Set to 0 to disable splitting. |
| `otlp.experimental_enable_resource_info` | Bool | `false` | Whether to synthesize the `greptime_otel_resource_info` table from OTLP metric<br/>resource attributes, so metrics-only services reach the semantic graph. |
| `prom_store` | -- | -- | Prometheus remote storage options |
| `prom_store.enable` | Bool | `true` | Whether to enable Prometheus remote write and read in HTTP API. |
| `prom_store.with_metric_engine` | Bool | `true` | Whether to store the data from Prometheus remote write in metric engine. |
| `prom_store.prom_validation_mode` | String | `strict` | Whether to enable validation for Prometheus remote write requests.<br/>Available options:<br/>- strict: deny invalid UTF-8 strings (default).<br/>- lossy: allow invalid UTF-8 strings, replace invalid characters with REPLACEMENT_CHARACTER(U+FFFD).<br/>- unchecked: do not valid strings. |
| `prom_store.experimental_enable_prometheus_native_histogram` | Bool | `false` | Experimental: enable Prometheus remote write v2 native histogram ingestion. |
| `wal` | -- | -- | The WAL options. |
| `wal.provider` | String | `raft_engine` | The provider of the WAL.<br/>- `raft_engine`: the wal is stored in the local file system by raft-engine.<br/>- `kafka`: it's remote wal that data is stored in Kafka.<br/>- `experimental_object_store`: the wal is stored as objects in an object store.<br/>**Notes: experimental and not supported yet.** |
| `wal.dir` | String | Unset | The directory to store the WAL files.<br/>**It's only used when the provider is `raft_engine`**. |
@@ -361,14 +359,12 @@
| `jaeger.enable` | Bool | `true` | Whether to enable Jaeger protocol in HTTP API. |
| `otlp` | -- | -- | OpenTelemetry protocol options. |
| `otlp.enable` | Bool | `true` | Whether to enable OpenTelemetry protocol in HTTP API. |
| `otlp.experimental_enable_exponential_histogram` | Bool | `false` | Experimental: enable cumulative OTLP exponential histogram ingestion. |
| `otlp.trace_ingest_chunk_size` | Integer | `512` | Maximum spans per trace ingest chunk. Set to 0 to disable splitting. |
| `otlp.experimental_enable_resource_info` | Bool | `false` | Whether to synthesize the `greptime_otel_resource_info` table from OTLP metric<br/>resource attributes, so metrics-only services reach the semantic graph. |
| `prom_store` | -- | -- | Prometheus remote storage options |
| `prom_store.enable` | Bool | `true` | Whether to enable Prometheus remote write and read in HTTP API. |
| `prom_store.with_metric_engine` | Bool | `true` | Whether to store the data from Prometheus remote write in metric engine. |
| `prom_store.prom_validation_mode` | String | `strict` | Whether to enable validation for Prometheus remote write requests.<br/>Available options:<br/>- strict: deny invalid UTF-8 strings (default).<br/>- lossy: allow invalid UTF-8 strings, replace invalid characters with REPLACEMENT_CHARACTER(U+FFFD).<br/>- unchecked: do not valid strings. |
| `prom_store.experimental_enable_prometheus_native_histogram` | Bool | `false` | Experimental: enable Prometheus remote write v2 native histogram ingestion. |
| `meta_client` | -- | -- | The metasrv client options. |
| `meta_client.metasrv_addrs` | Array | -- | The addresses of the metasrv. |
| `meta_client.timeout` | String | `3s` | Operation timeout. |
-4
View File
@@ -295,8 +295,6 @@ enable = true
[otlp]
## Whether to enable OpenTelemetry protocol in HTTP API.
enable = true
## Experimental: enable cumulative OTLP exponential histogram ingestion.
experimental_enable_exponential_histogram = false
## Maximum spans per trace ingest chunk. Set to 0 to disable splitting.
trace_ingest_chunk_size = 512
## Whether to synthesize the `greptime_otel_resource_info` table from OTLP metric
@@ -315,8 +313,6 @@ with_metric_engine = true
## - lossy: allow invalid UTF-8 strings, replace invalid characters with REPLACEMENT_CHARACTER(U+FFFD).
## - unchecked: do not valid strings.
prom_validation_mode = "strict"
## Experimental: enable Prometheus remote write v2 native histogram ingestion.
experimental_enable_prometheus_native_histogram = false
## The metasrv client options.
-4
View File
@@ -274,8 +274,6 @@ enable = true
[otlp]
## Whether to enable OpenTelemetry protocol in HTTP API.
enable = true
## Experimental: enable cumulative OTLP exponential histogram ingestion.
experimental_enable_exponential_histogram = false
## Maximum spans per trace ingest chunk. Set to 0 to disable splitting.
trace_ingest_chunk_size = 512
## Whether to synthesize the `greptime_otel_resource_info` table from OTLP metric
@@ -294,8 +292,6 @@ with_metric_engine = true
## - lossy: allow invalid UTF-8 strings, replace invalid characters with REPLACEMENT_CHARACTER(U+FFFD).
## - unchecked: do not valid strings.
prom_validation_mode = "strict"
## Experimental: enable Prometheus remote write v2 native histogram ingestion.
experimental_enable_prometheus_native_histogram = false
## The WAL options.
+12 -8
View File
@@ -13,8 +13,8 @@ field and evaluates them as first-class PromQL samples. Prometheus Remote Write
ingestion transport: accepted cumulative points are normalized into the same
Struct before persistence and are queried only as native histograms.
Native histograms are experimental. Prometheus Remote Write 2.0 and cumulative
OTLP exponential histograms are supported behind separate configuration gates.
Prometheus Remote Write 2.0 native histograms and cumulative OTLP/HTTP
exponential histograms are enabled by default.
Remote Write 1.0 histogram payloads are rejected instead of being acknowledged and dropped.
Native-histogram Remote Read is deferred; the existing Remote Read path
continues to return scalar samples only.
@@ -95,8 +95,7 @@ that send native histograms must use Remote Write 2.0.
## Remote Write 2.0
Remote Write 2.0 accepts integer and float native histograms while
`prom_store.experimental_enable_prometheus_native_histogram` is enabled. Supported
Remote Write 2.0 accepts integer and float native histograms by default. Supported
exponential schemas are `-4` through `8`; schema `-53` represents native
histograms with custom buckets.
@@ -121,10 +120,8 @@ explicit if sampled responses are implemented.
## OTLP
OTLP exponential histograms are accepted when
`otlp.experimental_enable_exponential_histogram` is enabled. The option defaults
to false and applies to OTLP/HTTP. Disabled points are rejected rather than
silently acknowledged. OTel Arrow exponential histograms are rejected because
OTLP/HTTP exponential histograms are accepted by default.
OTel Arrow exponential histograms are rejected because
the current Arrow wire format omits `zero_threshold`; accepting them would
silently change the distribution. Cumulative temporality is required; delta and
unspecified exponential histograms are rejected before their points are
@@ -284,3 +281,10 @@ native-histogram data remains readable.
- Versioned two-phase native-histogram aggregation state.
- Exact Prometheus summation compensation if measured precision differences
justify the extra aggregate state.
## Configuration migration
The `prom_store.experimental_enable_prometheus_native_histogram` and
`otlp.experimental_enable_exponential_histogram` options have been removed.
Remove them from existing configurations; both ingestion paths are always enabled
when their protocol is enabled. Old option values are ignored, including `false`.
-1
View File
@@ -1139,7 +1139,6 @@ mod tests {
fn test_toml() {
let opts = StandaloneOptions::default();
let toml_string = toml::to_string(&opts).unwrap();
assert!(toml_string.contains("experimental_enable_exponential_histogram = false"));
let parsed: StandaloneOptions = toml::from_str(&toml_string).unwrap();
assert_eq!(parsed.otlp, opts.otlp);
}
+27 -45
View File
@@ -199,12 +199,6 @@ fn test_load_frontend_example_config() {
let options =
GreptimeOptions::<FrontendOptions>::load_layered_options(example_config.to_str(), "")
.unwrap();
assert!(
!options
.component
.otlp
.experimental_enable_exponential_histogram
);
let expected = GreptimeOptions::<FrontendOptions> {
component: FrontendOptions {
pending_rows_batcher: PendingRowsBatcherOptions {
@@ -386,12 +380,6 @@ fn test_load_standalone_example_config() {
let options =
GreptimeOptions::<StandaloneOptions>::load_layered_options(example_config.to_str(), "")
.unwrap();
assert!(
!options
.component
.otlp
.experimental_enable_exponential_histogram
);
let expected = GreptimeOptions::<StandaloneOptions> {
component: StandaloneOptions {
pending_rows_batcher: PendingRowsBatcherOptions {
@@ -448,40 +436,34 @@ fn test_load_standalone_example_config() {
}
#[test]
fn test_load_otlp_exponential_histogram_option() {
let config = tempfile::NamedTempFile::new().unwrap();
std::fs::write(
config.path(),
"[otlp]\nexperimental_enable_exponential_histogram = true\n",
)
.unwrap();
fn test_load_removed_histogram_options() {
for enabled in [false, true] {
let config = tempfile::NamedTempFile::new().unwrap();
std::fs::write(
config.path(),
format!(
"[otlp]\nexperimental_enable_exponential_histogram = {enabled}\n\
[prom_store]\nenable = true\nwith_metric_engine = true\n\
experimental_enable_prometheus_native_histogram = {enabled}\n"
),
)
.unwrap();
let frontend =
GreptimeOptions::<FrontendOptions>::load_layered_options(config.path().to_str(), "")
.unwrap();
assert!(
frontend
.component
.otlp
.experimental_enable_exponential_histogram
);
let standalone =
GreptimeOptions::<StandaloneOptions>::load_layered_options(config.path().to_str(), "")
.unwrap();
assert!(
standalone
.component
.otlp
.experimental_enable_exponential_histogram
);
assert!(
standalone
.component
.frontend_options()
.otlp
.experimental_enable_exponential_histogram
);
let frontend =
GreptimeOptions::<FrontendOptions>::load_layered_options(config.path().to_str(), "")
.unwrap();
let standalone =
GreptimeOptions::<StandaloneOptions>::load_layered_options(config.path().to_str(), "")
.unwrap();
let defaults = FrontendOptions::default();
for options in [frontend.component, standalone.component.frontend_options()] {
assert_eq!(options.otlp, defaults.otlp);
assert_eq!(options.prom_store, defaults.prom_store);
let serialized = toml::to_string(&options).unwrap();
assert!(!serialized.contains("experimental_enable_exponential_histogram"));
assert!(!serialized.contains("experimental_enable_prometheus_native_histogram"));
}
}
}
#[test]
-1
View File
@@ -314,7 +314,6 @@ max_batch_rows = 25
let enabled: FrontendOptions = toml::from_str("experimental_metric_export = true").unwrap();
assert!(enabled.experimental_metric_export);
let toml_string = toml::to_string(&opts).unwrap();
assert!(toml_string.contains("experimental_enable_exponential_histogram = false"));
let parsed: FrontendOptions = toml::from_str(&toml_string).unwrap();
assert_eq!(parsed.otlp, opts.otlp);
assert_eq!(parsed.influxdb, opts.influxdb);
+2 -9
View File
@@ -146,19 +146,14 @@ where
Some(self.instance.clone()),
opts.prom_store.with_metric_engine,
opts.prom_store.prom_validation_mode,
opts.prom_store
.experimental_enable_prometheus_native_histogram,
pending_rows_batcher,
)
.with_prometheus_handler(self.instance.clone());
}
if opts.otlp.enable {
builder = builder.with_otlp_handler(
self.instance.clone(),
opts.prom_store.with_metric_engine,
opts.otlp.experimental_enable_exponential_histogram,
);
builder = builder
.with_otlp_handler(self.instance.clone(), opts.prom_store.with_metric_engine);
}
if opts.jaeger.enable {
@@ -623,8 +618,6 @@ mod tests {
opts.prom_store.pending_rows_flush_interval = Duration::from_secs(2);
opts.prom_store.with_metric_engine = metric_engine;
opts.prom_store.enable = prom_enabled;
opts.prom_store
.experimental_enable_prometheus_native_histogram = true;
let shared = &mut opts.pending_rows_batcher.table;
shared.protocols = vec![if selected {
BatchingProtocol::Prom
-10
View File
@@ -20,7 +20,6 @@ const DEFAULT_TRACE_INGEST_CHUNK_SIZE: usize = 512;
#[serde(default)]
pub struct OtlpOptions {
pub enable: bool,
pub experimental_enable_exponential_histogram: bool,
/// Maximum spans per trace ingest chunk. Set to 0 to disable splitting.
pub trace_ingest_chunk_size: usize,
/// Whether to synthesize the `greptime_otel_resource_info` descriptor
@@ -34,7 +33,6 @@ impl Default for OtlpOptions {
fn default() -> Self {
Self {
enable: true,
experimental_enable_exponential_histogram: false,
trace_ingest_chunk_size: DEFAULT_TRACE_INGEST_CHUNK_SIZE,
experimental_enable_resource_info: false,
}
@@ -49,13 +47,11 @@ mod tests {
fn test_otlp_options() {
let default = OtlpOptions::default();
assert!(default.enable);
assert!(!default.experimental_enable_exponential_histogram);
assert_eq!(default.trace_ingest_chunk_size, 512);
assert!(!default.experimental_enable_resource_info);
let options: OtlpOptions = toml::from_str("enable = false").unwrap();
assert!(!options.enable);
assert!(!options.experimental_enable_exponential_histogram);
assert_eq!(
options.trace_ingest_chunk_size,
DEFAULT_TRACE_INGEST_CHUNK_SIZE
@@ -63,15 +59,9 @@ mod tests {
let options: OtlpOptions = toml::from_str("trace_ingest_chunk_size = 0").unwrap();
assert!(options.enable);
assert!(!options.experimental_enable_exponential_histogram);
assert_eq!(options.trace_ingest_chunk_size, 0);
let options: OtlpOptions =
toml::from_str("experimental_enable_exponential_histogram = true").unwrap();
assert!(options.experimental_enable_exponential_histogram);
let serialized = toml::to_string(&options).unwrap();
assert!(serialized.contains("experimental_enable_exponential_histogram = true"));
assert_eq!(toml::from_str::<OtlpOptions>(&serialized).unwrap(), options);
}
}
@@ -27,9 +27,6 @@ pub struct PromStoreOptions {
/// Validation mode while decoding Prometheus remote write requests.
#[serde(default)]
pub prom_validation_mode: PromValidationMode,
/// Enables experimental Prometheus remote write v2 native histogram ingestion.
#[serde(default)]
pub experimental_enable_prometheus_native_histogram: bool,
#[serde(default, with = "humantime_serde")]
pub pending_rows_flush_interval: Duration,
#[serde(default = "default_max_batch_rows")]
@@ -89,7 +86,6 @@ impl Default for PromStoreOptions {
enable: true,
with_metric_engine: true,
prom_validation_mode: PromValidationMode::Strict,
experimental_enable_prometheus_native_histogram: false,
pending_rows_flush_interval: Duration::ZERO,
max_batch_rows: default_max_batch_rows(),
max_concurrent_flushes: default_max_concurrent_flushes(),
@@ -152,7 +148,6 @@ mod tests {
assert!(default.enable);
assert!(default.with_metric_engine);
assert_eq!(default.prom_validation_mode, PromValidationMode::Strict);
assert!(!default.experimental_enable_prometheus_native_histogram);
assert_eq!(default.pending_rows_flush_interval, Duration::ZERO);
assert_eq!(default.max_batch_rows, default_max_batch_rows());
assert_eq!(
+5 -16
View File
@@ -18,10 +18,9 @@ use api::v1::{ArrowIpc, SemanticType};
use bytes::Bytes;
use common_grpc::flight::{FlightEncoder, FlightMessage};
use datatypes::arrow::record_batch::RecordBatch;
use snafu::{OptionExt, ResultExt, ensure};
use snafu::{OptionExt, ensure};
use store_api::codec::PrimaryKeyEncoding;
use store_api::metadata::RegionMetadataRef;
use store_api::region_engine::RegionEngine;
use store_api::region_request::{AffectedRows, RegionBulkInsertsRequest, RegionRequest};
use store_api::storage::RegionId;
@@ -74,9 +73,9 @@ impl MetricEngineInner {
region_id: RegionId,
mut request: RegionBulkInsertsRequest,
) -> Result<AffectedRows> {
// Simply set the aligned schema to the data region schema version to avoid filling missing columns
// because that schema should be constant and callers have ensured request has the same schema.
request.aligned_schema_version = Some(self.physical_schema_version(region_id).await?);
// Other logical tables can add fields to the shared physical schema.
// Let Mito fill fields absent from this batch before writing it.
request.aligned_schema_version = None;
self.data_region
.write_data(region_id, RegionRequest::BulkInserts(request))
.await
@@ -121,7 +120,6 @@ impl MetricEngineInner {
let (schema, data_header, payload) = record_batch_to_ipc(&modified_batch)?;
let partition_expr_version = request.partition_expr_version;
let aligned_schema_version = Some(self.physical_schema_version(data_region_id).await?);
let request = RegionBulkInsertsRequest {
skip_wal: request.skip_wal,
@@ -133,22 +131,13 @@ impl MetricEngineInner {
payload,
},
partition_expr_version,
aligned_schema_version,
aligned_schema_version: None,
};
self.data_region
.write_data(data_region_id, RegionRequest::BulkInserts(request))
.await
}
async fn physical_schema_version(&self, region_id: RegionId) -> Result<u64> {
Ok(self
.mito
.get_metadata(region_id)
.await
.context(error::MitoReadOperationSnafu)?
.schema_version)
}
fn resolve_tag_columns_from_metadata(
&self,
logical_region_id: RegionId,
+7 -1
View File
@@ -20,7 +20,7 @@ use std::time::{Duration, Instant};
use api::helper::{ColumnDataTypeWrapper, to_grpc_value};
use api::v1::bulk_wal_entry::Body;
use api::v1::{ArrowIpc, BulkWalEntry, Mutation, OpType};
use api::v1::{ArrowIpc, BulkWalEntry, Mutation, OpType, SemanticType};
use bytes::Bytes;
use common_grpc::flight::{FlightDecoder, FlightEncoder, FlightMessage};
use common_recordbatch::DfRecordBatch as RecordBatch;
@@ -218,6 +218,12 @@ impl BulkPart {
// Finds columns that need to be filled
let mut columns_to_fill = Vec::new();
for column_meta in &region_metadata.column_metadatas {
// Sparse bulk batches already carry tags in the encoded primary key.
if region_metadata.primary_key_encoding == PrimaryKeyEncoding::Sparse
&& column_meta.semantic_type == SemanticType::Tag
{
continue;
}
// TODO(yingwen): Returns error if it is impure default after we support filling
// bulk insert request in the frontend
if !batch_columns.contains(column_meta.column_schema.name.as_str()) {
+33 -2
View File
@@ -37,7 +37,7 @@ use datatypes::arrow::array::{
Array, ArrayRef, BinaryArray, DictionaryArray, UInt32Array, UInt64Array,
};
use datatypes::arrow::compute::kernels::take::take;
use datatypes::arrow::datatypes::{Schema, SchemaRef};
use datatypes::arrow::datatypes::{DataType as ArrowDataType, Schema, SchemaRef};
use datatypes::arrow::record_batch::RecordBatch;
use datatypes::prelude::{ConcreteDataType, DataType};
use datatypes::value::ValueRef;
@@ -407,11 +407,42 @@ impl FlatReadFormat {
};
// First, apply flat format conversion.
let batch = match &self.parquet_adapter {
let mut batch = match &self.parquet_adapter {
ParquetAdapter::Flat(_) => record_batch,
ParquetAdapter::PrimaryKeyToFlat(p) => p.convert_batch(record_batch)?,
};
// Normalize nested field names and metadata to the SST's region metadata
// before schema compatibility and merging with memtables. This removes
// Parquet-added field IDs; equals_datatype also permits nested name differences.
for index in 0..batch.num_columns() {
let array = batch.column(index);
if !matches!(array.data_type(), ArrowDataType::Struct(_)) {
continue;
}
let field = batch.schema_ref().field(index);
let Some(column) = self.metadata().column_by_name(field.name()) else {
continue;
};
let target = column.column_schema.data_type.as_arrow_type();
if array.data_type() != &target && array.data_type().equals_datatype(&target) {
let array =
datatypes::arrow::compute::cast(array, &target).context(ComputeArrowSnafu)?;
let mut fields = batch.schema().fields().to_vec();
fields[index] = Arc::new(field.clone().with_data_type(target));
let mut columns = batch.columns().to_vec();
columns[index] = array;
batch = RecordBatch::try_new(
Arc::new(Schema::new_with_metadata(
fields,
batch.schema().metadata().clone(),
)),
columns,
)
.context(NewRecordBatchSnafu)?;
}
}
// Then apply sequence override if provided
let Some(override_array) = override_sequence_array else {
return Ok(batch);
+3 -4
View File
@@ -288,10 +288,9 @@ fn bench_prom_v1_vs_v2_decode(c: &mut Criterion) {
group.bench_with_input(BenchmarkId::new("v2_manual", name), v2_bytes, |b, bytes| {
b.iter(|| {
black_box(
remote_write_v2::decode_uncompressed_write_requests(
black_box(bytes.as_slice()),
true,
)
remote_write_v2::decode_uncompressed_write_requests(black_box(
bytes.as_slice(),
))
.unwrap(),
);
});
+1 -10
View File
@@ -728,7 +728,6 @@ impl HttpServerBuilder {
pipeline_handler: Option<PipelineHandlerRef>,
prom_store_with_metric_engine: bool,
prom_validation_mode: PromValidationMode,
experimental_enable_prometheus_native_histogram: bool,
pending_rows_batcher: Option<Arc<LogicalTablePendingRowsBatcher>>,
) -> Self {
let state = PromStoreState {
@@ -736,7 +735,6 @@ impl HttpServerBuilder {
pipeline_handler,
prom_store_with_metric_engine,
prom_validation_mode,
experimental_enable_prometheus_native_histogram,
pending_rows_batcher,
};
@@ -763,16 +761,11 @@ impl HttpServerBuilder {
self,
handler: OpenTelemetryProtocolHandlerRef,
with_metric_engine: bool,
experimental_enable_exponential_histogram: bool,
) -> Self {
Self {
router: self.router.nest(
&format!("/{HTTP_API_VERSION}/otlp"),
HttpServer::route_otlp(
handler,
with_metric_engine,
experimental_enable_exponential_histogram,
),
HttpServer::route_otlp(handler, with_metric_engine),
),
..self
}
@@ -1500,7 +1493,6 @@ impl HttpServer {
fn route_otlp<S>(
otlp_handler: OpenTelemetryProtocolHandlerRef,
with_metric_engine: bool,
experimental_enable_exponential_histogram: bool,
) -> Router<S> {
Router::new()
.route("/v1/metrics", routing::post(otlp::metrics))
@@ -1516,7 +1508,6 @@ impl HttpServer {
))
.with_state(OtlpState {
with_metric_engine,
experimental_enable_exponential_histogram,
handler: otlp_handler,
})
}
-3
View File
@@ -77,7 +77,6 @@ fn content_type_to_string(content_type: Option<&TypedHeader<ContentType>>) -> St
#[derive(Clone)]
pub struct OtlpState {
pub with_metric_engine: bool,
pub experimental_enable_exponential_histogram: bool,
pub handler: OpenTelemetryProtocolHandlerRef,
}
@@ -108,7 +107,6 @@ pub async fn metrics(
let OtlpState {
with_metric_engine,
experimental_enable_exponential_histogram,
handler,
} = state;
@@ -117,7 +115,6 @@ pub async fn metrics(
resource_attrs: http_opts.resource_attrs,
promote_scope_attrs: http_opts.promote_scope_attrs,
with_metric_engine,
experimental_enable_exponential_histogram,
// set by the frontend from its config
is_legacy: false,
resource_info: false,
+1 -9
View File
@@ -74,7 +74,6 @@ pub struct PromStoreState {
pub pipeline_handler: Option<PipelineHandlerRef>,
pub prom_store_with_metric_engine: bool,
pub prom_validation_mode: PromValidationMode,
pub experimental_enable_prometheus_native_histogram: bool,
pub pending_rows_batcher: Option<Arc<LogicalTablePendingRowsBatcher>>,
}
@@ -146,7 +145,6 @@ async fn remote_write_v1(
pipeline_handler,
prom_store_with_metric_engine,
prom_validation_mode,
experimental_enable_prometheus_native_histogram: _,
pending_rows_batcher,
} = state;
@@ -228,7 +226,6 @@ async fn remote_write_v2(
pipeline_handler: _,
prom_store_with_metric_engine,
prom_validation_mode: _,
experimental_enable_prometheus_native_histogram,
pending_rows_batcher,
} = state;
@@ -243,11 +240,7 @@ async fn remote_write_v2(
let (db, mut query_ctx, _timer) =
prepare_remote_write_context(&params, query_ctx, REMOTE_WRITE_V2_VERSION);
let req = match decode_remote_write_v2(
is_zstd,
body,
experimental_enable_prometheus_native_histogram,
) {
let req = match decode_remote_write_v2(is_zstd, body) {
Ok(req) => req,
Err(error) => return Ok(remote_write_v2_error_response(error, 0, 0, 0)),
};
@@ -1017,7 +1010,6 @@ mod tests {
pipeline_handler: None,
prom_store_with_metric_engine: false,
prom_validation_mode: PromValidationMode::Strict,
experimental_enable_prometheus_native_histogram: false,
pending_rows_batcher: None,
}
}
+86 -34
View File
@@ -24,6 +24,7 @@ use otel_arrow_rust::proto::opentelemetry::arrow::v1::arrow_metrics_service_serv
use otel_arrow_rust::proto::opentelemetry::arrow::v1::{
BatchArrowRecords, BatchStatus, StatusCode as ArrowStatusCode,
};
use otel_arrow_rust::proto::opentelemetry::collector::metrics::v1::ExportMetricsServiceRequest;
use otel_arrow_rust::proto::opentelemetry::metrics::v1::metric;
use session::protocol_ctx::{OtlpMetricCtx, ProtocolCtx};
use tonic::metadata::{Entry, MetadataValue};
@@ -50,23 +51,42 @@ impl<T> OtelArrowServiceHandler<T> {
}
}
/// Removes unsupported Arrow histograms before ingestion, preserving other metrics.
fn remove_exponential_histograms(request: &mut ExportMetricsServiceRequest) -> bool {
let mut has_data_points = false;
for scope in request
.resource_metrics
.iter_mut()
.flat_map(|resource| &mut resource.scope_metrics)
{
scope.metrics.retain(|item| {
if let Some(metric::Data::ExponentialHistogram(histogram)) = &item.data {
has_data_points |= !histogram.data_points.is_empty();
false
} else {
true
}
});
}
has_data_points
}
fn batch_status(
batch_id: i64,
outcome: MetricsIngestOutcome,
has_exponential_histogram_data_points: bool,
) -> BatchStatus {
let status_code = if outcome.accepted_data_points == 0 && outcome.rejected_data_points > 0 {
let status_code = if outcome.accepted_data_points == 0
&& (outcome.rejected_data_points > 0 || has_exponential_histogram_data_points)
{
ArrowStatusCode::InvalidArgument
} else {
ArrowStatusCode::Ok
};
let status_message = match outcome.error_message {
// Arrow keeps the feature gate off, so these fail before per-point validation.
Some(_) if has_exponential_histogram_data_points => {
EXPONENTIAL_HISTOGRAM_UNSUPPORTED.to_string()
}
Some(message) => message,
None => String::new(),
let status_message = if has_exponential_histogram_data_points {
EXPONENTIAL_HISTOGRAM_UNSUPPORTED.to_string()
} else {
outcome.error_message.unwrap_or_default()
};
BatchStatus {
batch_id,
@@ -117,7 +137,7 @@ impl ArrowMetricsService for OtelArrowServiceHandler<OpenTelemetryProtocolHandle
}
};
let batch_id = batch.batch_id;
let request = match consumer.consume_metrics_batches(&mut batch).map_err(|e| {
let mut request = match consumer.consume_metrics_batches(&mut batch).map_err(|e| {
error::HandleOtelArrowRequestSnafu {
err_msg: e.to_string(),
}
@@ -137,18 +157,8 @@ impl ArrowMetricsService for OtelArrowServiceHandler<OpenTelemetryProtocolHandle
return;
}
};
let has_exponential_histogram_data_points = request
.resource_metrics
.iter()
.flat_map(|resource| &resource.scope_metrics)
.flat_map(|scope| &scope.metrics)
.any(|item| {
matches!(
item.data.as_ref(),
Some(metric::Data::ExponentialHistogram(histogram))
if !histogram.data_points.is_empty()
)
});
let has_exponential_histogram_data_points =
remove_exponential_histograms(&mut request);
let outcome = match handler.metrics(request, query_ctx.clone()).await {
Ok(outcome) => outcome,
Err(error::Error::InvalidOtlpMetricInput { reason }) => {
@@ -202,19 +212,61 @@ mod tests {
use super::*;
#[test]
fn batch_status_explains_arrow_exponential_histogram_limit() {
let status = batch_status(
7,
MetricsIngestOutcome {
rejected_data_points: 1,
error_message: Some("internal OTLP rejection detail".to_string()),
..Default::default()
},
true,
);
fn removes_arrow_exponential_histograms_and_preserves_other_metrics() {
use otel_arrow_rust::proto::opentelemetry::metrics::v1::{
ExponentialHistogram, ExponentialHistogramDataPoint, Gauge, Metric, ResourceMetrics,
ScopeMetrics,
};
assert_eq!(7, status.batch_id);
assert_eq!(ArrowStatusCode::InvalidArgument as i32, status.status_code);
assert_eq!(EXPONENTIAL_HISTOGRAM_UNSUPPORTED, status.status_message);
let mut request = ExportMetricsServiceRequest {
resource_metrics: vec![ResourceMetrics {
scope_metrics: vec![ScopeMetrics {
metrics: vec![
Metric {
data: Some(metric::Data::ExponentialHistogram(ExponentialHistogram {
data_points: vec![ExponentialHistogramDataPoint::default()],
..Default::default()
})),
..Default::default()
},
Metric {
data: Some(metric::Data::Gauge(Gauge::default())),
..Default::default()
},
],
..Default::default()
}],
..Default::default()
}],
};
assert!(remove_exponential_histograms(&mut request));
let metrics = &request.resource_metrics[0].scope_metrics[0].metrics;
assert_eq!(metrics.len(), 1);
assert!(matches!(metrics[0].data, Some(metric::Data::Gauge(_))));
assert!(!remove_exponential_histograms(&mut request));
}
#[test]
fn batch_status_explains_arrow_exponential_histogram_limit() {
for accepted_data_points in [0, 1] {
let status = batch_status(
7,
MetricsIngestOutcome {
accepted_data_points,
..Default::default()
},
true,
);
assert_eq!(7, status.batch_id);
assert_eq!(
if accepted_data_points == 0 {
ArrowStatusCode::InvalidArgument
} else {
ArrowStatusCode::Ok
} as i32,
status.status_code
);
assert_eq!(EXPONENTIAL_HISTOGRAM_UNSUPPORTED, status.status_message);
}
}
}
+18 -34
View File
@@ -138,7 +138,7 @@ pub fn to_grpc_insert_requests(
&& !metric_ctx.is_legacy
&& let Some(r) = resource.resource.as_ref()
{
resource_info.observe(&r.attributes, resource, metric_ctx);
resource_info.observe(&r.attributes, resource);
}
let resource_attrs = resource.resource.as_ref().map(|r| {
@@ -630,7 +630,7 @@ fn encode_exponential_histogram(
metric_ctx: &OtlpMetricCtx,
outcome: &mut MetricsIngestOutcome,
) -> Result<bool> {
if let Err(rejection) = exponential_histogram_gate(histogram, metric_ctx) {
if let Err(rejection) = exponential_histogram_gate(histogram) {
reject_data_points(outcome, histogram.data_points.len(), || {
rejection.message(name)
})?;
@@ -684,7 +684,6 @@ fn encode_exponential_histogram(
}
pub(crate) enum ExponentialHistogramRejection {
Disabled,
DeltaTemporality,
UnspecifiedTemporality,
}
@@ -692,9 +691,6 @@ pub(crate) enum ExponentialHistogramRejection {
impl ExponentialHistogramRejection {
fn message(&self, name: &str) -> String {
match self {
Self::Disabled => format!(
"metric `{name}` uses OTLP exponential histograms; set otlp.experimental_enable_exponential_histogram = true to enable ingestion"
),
Self::DeltaTemporality => format!(
"metric `{name}` uses delta OTLP exponential histograms; only cumulative temporality is supported"
),
@@ -710,11 +706,7 @@ impl ExponentialHistogramRejection {
/// Individual points can still fail [`exponential_histogram_value`].
pub(crate) fn exponential_histogram_gate(
histogram: &ExponentialHistogram,
metric_ctx: &OtlpMetricCtx,
) -> std::result::Result<(), ExponentialHistogramRejection> {
if !metric_ctx.experimental_enable_exponential_histogram {
return Err(ExponentialHistogramRejection::Disabled);
}
match AggregationTemporality::try_from(histogram.aggregation_temporality) {
Ok(AggregationTemporality::Cumulative) => Ok(()),
Ok(AggregationTemporality::Delta) => Err(ExponentialHistogramRejection::DeltaTemporality),
@@ -2498,7 +2490,7 @@ mod tests {
}
#[test]
fn test_exponential_histogram_gate_and_partial_outcome() {
fn test_exponential_histogram_default_context_and_partial_outcome() {
let request = metrics_request(vec![
Metric {
name: "temperature".to_string(),
@@ -2512,6 +2504,11 @@ mod tests {
vec![exponential_point()],
AggregationTemporality::Cumulative,
),
exponential_metric(
"delta_latency",
vec![exponential_point()],
AggregationTemporality::Delta,
),
]);
let MetricsConversion {
requests,
@@ -2520,20 +2517,20 @@ mod tests {
..
} = to_grpc_insert_requests(request, &mut OtlpMetricCtx::default()).unwrap();
assert_eq!(outcome.accepted_data_points, 1);
assert_eq!(outcome.accepted_data_points, 2);
assert_eq!(outcome.rejected_data_points, 1);
assert!(
outcome
.error_message
.as_deref()
.unwrap()
.contains("otlp.experimental_enable_exponential_histogram")
.contains("only cumulative temporality is supported")
);
assert_eq!(requests.inserts.len(), 1);
assert_eq!(requests.inserts[0].table_name, "temperature");
assert_eq!(requests.inserts.len(), 2);
let semantics = decode(&semantic_index);
assert!(semantics.contains_key("temperature"));
assert!(!semantics.contains_key("latency"));
assert!(semantics.contains_key("latency"));
assert!(!semantics.contains_key("delta_latency"));
let empty = metrics_request(vec![exponential_metric(
"empty",
@@ -2566,10 +2563,7 @@ mod tests {
AggregationTemporality::Cumulative,
),
]);
let mut ctx = OtlpMetricCtx {
experimental_enable_exponential_histogram: true,
..Default::default()
};
let mut ctx = OtlpMetricCtx::default();
let error = to_grpc_insert_requests(request, &mut ctx).unwrap_err();
assert!(
@@ -2589,10 +2583,7 @@ mod tests {
),
histogram_metric("latency"),
]);
let mut ctx = OtlpMetricCtx {
experimental_enable_exponential_histogram: true,
..Default::default()
};
let mut ctx = OtlpMetricCtx::default();
let error = to_grpc_insert_requests(request, &mut ctx).unwrap_err();
assert!(
@@ -2657,10 +2648,7 @@ mod tests {
vec![stale.clone()],
temporality,
)]);
let mut ctx = OtlpMetricCtx {
experimental_enable_exponential_histogram: true,
..Default::default()
};
let mut ctx = OtlpMetricCtx::default();
let MetricsConversion {
requests,
rows,
@@ -2688,15 +2676,11 @@ mod tests {
vec![point],
AggregationTemporality::Cumulative,
)]);
let mut new_ctx = OtlpMetricCtx {
experimental_enable_exponential_histogram: true,
..Default::default()
};
let mut new_ctx = OtlpMetricCtx::default();
let new_requests = to_grpc_insert_requests(request.clone(), &mut new_ctx)
.unwrap()
.requests;
let mut legacy_ctx = OtlpMetricCtx {
experimental_enable_exponential_histogram: true,
is_legacy: true,
..Default::default()
};
@@ -2749,7 +2733,7 @@ mod tests {
let request = metrics_request(vec![exponential_metric(
"x".repeat(1_000),
vec![ExponentialHistogramDataPoint::default()],
AggregationTemporality::Cumulative,
AggregationTemporality::Delta,
)]);
let outcome = to_grpc_insert_requests(request, &mut OtlpMetricCtx::default())
.unwrap()
+12 -39
View File
@@ -30,7 +30,6 @@ use otel_arrow_rust::proto::opentelemetry::common::v1::KeyValue;
use otel_arrow_rust::proto::opentelemetry::metrics::v1::{
AggregationTemporality, ResourceMetrics, metric,
};
use session::protocol_ctx::OtlpMetricCtx;
use crate::error::Result;
use crate::otlp::metrics::{
@@ -90,12 +89,7 @@ pub struct ResourceInfoData {
impl ResourceInfoData {
/// Takes the raw attributes, before the promote filter runs on them.
pub fn observe(
&mut self,
raw_attrs: &[KeyValue],
resource: &ResourceMetrics,
metric_ctx: &OtlpMetricCtx,
) {
pub fn observe(&mut self, raw_attrs: &[KeyValue], resource: &ResourceMetrics) {
let mut tags = Vec::with_capacity(MAX_PROJECTED_TAGS);
let ServiceIdentity { job, instance } = service_identity(raw_attrs);
if let Some(job) = job {
@@ -119,7 +113,7 @@ impl ResourceInfoData {
tags.sort_unstable();
let mut observed: BTreeMap<i64, i64> = BTreeMap::new();
for_each_encoded_time(resource, metric_ctx, |ts| {
for_each_encoded_time(resource, |ts| {
let window = ts - ts.rem_euclid(SEMANTIC_GRAPH_WINDOW_NANOS);
observed
.entry(window)
@@ -176,11 +170,7 @@ impl ResourceInfoData {
/// Visits the times of the data points the encoder writes rows for, so a
/// resource is described exactly where it is measured rather than wherever
/// its request happens to reach.
fn for_each_encoded_time(
resource: &ResourceMetrics,
metric_ctx: &OtlpMetricCtx,
mut visit: impl FnMut(i64),
) {
fn for_each_encoded_time(resource: &ResourceMetrics, mut visit: impl FnMut(i64)) {
fn visit_all(points: impl Iterator<Item = u64>, visit: &mut impl FnMut(i64)) {
for ts in points {
visit(ts as i64);
@@ -214,7 +204,7 @@ fn for_each_encoded_time(
visit_all(s.data_points.iter().map(|p| p.time_unix_nano), &mut visit)
}
Some(metric::Data::ExponentialHistogram(h))
if exponential_histogram_gate(h, metric_ctx).is_ok() =>
if exponential_histogram_gate(h).is_ok() =>
{
for point in &h.data_points {
if let Ok((_, ts)) = exponential_histogram_value(point) {
@@ -285,7 +275,7 @@ mod tests {
kv("k8s.node.name", "node-a"),
kv("os.type", "linux"),
];
data.observe(&attrs, &gauge_at(&[100, 50]), &OtlpMetricCtx::default());
data.observe(&attrs, &gauge_at(&[100, 50]));
assert_eq!(data.rows.len(), 1);
let (tags, windows) = data.rows.iter().next().unwrap();
assert_eq!(windows.values().copied().collect::<Vec<_>>(), vec![100]);
@@ -298,19 +288,11 @@ mod tests {
.all(|(k, _)| k != "os.type" && k != "service.instance.id")
);
data.observe(
&[kv("host.id", "h-2")],
&gauge_at(&[10]),
&OtlpMetricCtx::default(),
);
data.observe(&[kv("host.id", "h-2")], &gauge_at(&[10]));
assert_eq!(data.rows.len(), 2);
let mut empty = ResourceInfoData::default();
empty.observe(
&[kv("os.type", "linux")],
&gauge_at(&[100]),
&OtlpMetricCtx::default(),
);
empty.observe(&[kv("os.type", "linux")], &gauge_at(&[100]));
assert!(empty.into_row_insert_requests().unwrap().is_none());
}
@@ -322,7 +304,6 @@ mod tests {
data.observe(
&[kv("service.name", "api")],
&gauge_at(&[window + 1, window + 2, 3 * window + 7]),
&OtlpMetricCtx::default(),
);
let windows = data.rows.values().next().unwrap();
@@ -352,20 +333,13 @@ mod tests {
}],
..Default::default()
};
let enabled = OtlpMetricCtx {
experimental_enable_exponential_histogram: true,
..Default::default()
};
for (resource, ctx) in [
(
exponential(AggregationTemporality::Cumulative),
OtlpMetricCtx::default(),
),
(exponential(AggregationTemporality::Delta), enabled),
for temporality in [
AggregationTemporality::Delta,
AggregationTemporality::Unspecified,
] {
let resource = exponential(temporality);
let mut data = ResourceInfoData::default();
data.observe(&[kv("service.name", "api")], &resource, &ctx);
data.observe(&[kv("service.name", "api")], &resource);
assert!(data.into_row_insert_requests().unwrap().is_none());
}
}
@@ -377,7 +351,6 @@ mod tests {
data.observe(
&[kv("service.name", "api"), kv("host.id", "h-1")],
&gauge_at(&[1_700_000_000_123_456_789]),
&OtlpMetricCtx::default(),
);
let requests = data.into_row_insert_requests().unwrap().unwrap();
assert_eq!(requests.inserts.len(), 1);
@@ -56,11 +56,7 @@ fn observe_ignores_rejected_classic_histogram_windows() {
..Default::default()
};
let mut data = ResourceInfoData::default();
data.observe(
&[kv("service.name", "api")],
&resource,
&OtlpMetricCtx::default(),
);
data.observe(&[kv("service.name", "api")], &resource);
let windows = data.rows.values().next().unwrap();
assert_eq!(
+38 -55
View File
@@ -130,7 +130,6 @@ pub(crate) struct RemoteWriteV2WriteRequests {
pub(crate) fn decode_remote_write_v2(
is_zstd: bool,
body: Bytes,
native_histograms_enabled: bool,
) -> Result<RemoteWriteV2WriteRequests> {
let decode_timer = crate::metrics::METRIC_HTTP_PROM_STORE_CODEC_ELAPSED
.with_label_values(&["decode", REMOTE_WRITE_V2_VERSION])
@@ -151,13 +150,10 @@ pub(crate) fn decode_remote_write_v2(
let _convert_timer = crate::metrics::METRIC_HTTP_PROM_STORE_CODEC_ELAPSED
.with_label_values(&["convert", REMOTE_WRITE_V2_VERSION])
.start_timer();
convert_remote_write_v2(request, native_histograms_enabled)
convert_remote_write_v2(request)
}
fn convert_remote_write_v2(
request: BorrowedRequest<'_>,
native_histograms_enabled: bool,
) -> Result<RemoteWriteV2WriteRequests> {
fn convert_remote_write_v2(request: BorrowedRequest<'_>) -> Result<RemoteWriteV2WriteRequests> {
ensure!(
request.symbols.first().copied() == Some(""),
error::InvalidPromRemoteRequestSnafu {
@@ -179,14 +175,6 @@ fn convert_remote_write_v2(
let counts = scan_series(series, &mut labels_refs, &mut metadata)
.context(error::DecodePromRemoteRequestSnafu)?;
ensure!(
native_histograms_enabled || counts.histograms == 0,
error::InvalidPromRemoteRequestSnafu {
msg: "prometheus remote write v2 native histogram ingestion is experimental; set prom_store.experimental_enable_prometheus_native_histogram = true to enable it"
.to_string(),
}
);
if counts.samples == 0 && counts.histograms == 0 {
decode_series_leaves(series, None, Vec::new(), 0, &mut scratch)?;
continue;
@@ -831,15 +819,14 @@ pub mod test_util {
request: Request,
) -> Result<(Vec<RowInsertRequest>, Vec<RowInsertRequest>, u64, u64)> {
let body = Bytes::from(snappy_compress(&request.encode_to_vec())?);
decode_write_requests(false, body, true)
decode_write_requests(false, body)
}
pub fn decode_write_requests(
is_zstd: bool,
body: Bytes,
native_histograms_enabled: bool,
) -> Result<(Vec<RowInsertRequest>, Vec<RowInsertRequest>, u64, u64)> {
let requests = super::decode_remote_write_v2(is_zstd, body, native_histograms_enabled)?;
let requests = super::decode_remote_write_v2(is_zstd, body)?;
Ok((
requests.samples.all_req().collect(),
requests.histograms.all_req().collect(),
@@ -850,11 +837,10 @@ pub mod test_util {
pub fn decode_uncompressed_write_requests(
body: &[u8],
native_histograms_enabled: bool,
) -> Result<(Vec<RowInsertRequest>, Vec<RowInsertRequest>, u64, u64)> {
let request =
super::BorrowedRequest::decode(body).context(error::DecodePromRemoteRequestSnafu)?;
let requests = super::convert_remote_write_v2(request, native_histograms_enabled)?;
let requests = super::convert_remote_write_v2(request)?;
Ok((
requests.samples.all_req().collect(),
requests.histograms.all_req().collect(),
@@ -953,12 +939,7 @@ mod tests {
assert_eq!(decoded.timeseries[0].samples.len(), 1);
assert_eq!(decoded.timeseries[0].samples[0].value, 42.0);
assert_eq!(decoded.timeseries[0].metadata.as_ref().unwrap().r#type, 1);
assert_eq!(
decode_remote_write_v2(true, body, true)
.unwrap()
.sample_count,
1
);
assert_eq!(decode_remote_write_v2(true, body).unwrap().sample_count, 1);
}
#[test]
@@ -990,7 +971,7 @@ mod tests {
wire.extend(encoded_message_field(5, &packed_u32_field(1, &[99])));
wire.extend(string_field(4, b"http_requests_total"));
let requests = decode_wire(&wire, true).unwrap();
let requests = decode_wire(&wire).unwrap();
assert_eq!(requests.sample_count, 2);
assert_eq!(requests.histogram_count, 0);
let rows = requests.samples.all_req().next().unwrap().rows.unwrap();
@@ -1028,7 +1009,7 @@ mod tests {
series.extend(packed_u32_field(1, &[1, 2]));
let wire = request_wire(&["", METRIC_NAME_LABEL, "metric"], &[series]);
let requests = decode_wire(&wire, true).unwrap();
let requests = decode_wire(&wire).unwrap();
assert_eq!(requests.histogram_count, 1);
let rows = requests.histograms.all_req().next().unwrap().rows.unwrap();
assert_eq!(
@@ -1093,7 +1074,7 @@ mod tests {
("invalid sample", invalid_sample),
("invalid histogram", invalid_histogram),
] {
let error = decode_wire_error(&wire, true, name);
let error = decode_wire_error(&wire, name);
assert!(
matches!(error, error::Error::DecodePromRemoteRequest { .. }),
"{name}: {error}"
@@ -1108,7 +1089,7 @@ mod tests {
series.extend(encoded_message_field(tag, &[0x08]));
let wire = request_wire(&["", METRIC_NAME_LABEL, "metric"], &[series]);
let error = decode_wire_error(&wire, true, "malformed ignored message");
let error = decode_wire_error(&wire, "malformed ignored message");
assert!(matches!(
error,
error::Error::DecodePromRemoteRequest { .. }
@@ -1174,8 +1155,8 @@ mod tests {
})
.unwrap();
assert_eq!(requests.sample_count, 0);
assert!(decode_wire(&[], true).is_err());
assert!(decode_wire(&[0x0a, 0x00], true).is_err());
assert!(decode_wire(&[]).is_err());
assert!(decode_wire(&[0x0a, 0x00]).is_err());
}
#[test]
@@ -1205,22 +1186,34 @@ mod tests {
}
#[test]
fn test_fused_decoder_pins_experimental_error_precedence() {
let histogram = series_wire(&[1, 2], 3, &Histogram::default().encode_to_vec());
let malformed_sample = series_wire(&[1, 2], 2, &[0x08]);
fn test_fused_decoder_rejects_malformed_series_with_histograms() {
let histogram = series_wire(
&[1, 2],
3,
&Histogram {
count: Some(Count::CountInt(0)),
zero_count: Some(ZeroCount::ZeroCountInt(0)),
..Default::default()
}
.encode_to_vec(),
);
let malformed_sample = series_wire(&[1, 3], 2, &[0x08]);
let wire = request_wire(
&["", METRIC_NAME_LABEL, "metric"],
&["", METRIC_NAME_LABEL, "histogram", "sample"],
&[histogram.clone(), malformed_sample.clone()],
);
let error = decode_wire_error(&wire, false, "histogram before malformed series");
assert!(error.to_string().contains("ingestion is experimental"));
let error = decode_wire_error(&wire, "histogram before malformed series");
assert!(matches!(
error,
error::Error::DecodePromRemoteRequest { .. }
));
let wire = request_wire(
&["", METRIC_NAME_LABEL, "metric"],
&["", METRIC_NAME_LABEL, "histogram", "sample"],
&[malformed_sample, histogram.clone()],
);
let error = decode_wire_error(&wire, false, "malformed series before histogram");
let error = decode_wire_error(&wire, "malformed series before histogram");
assert!(matches!(
error,
error::Error::DecodePromRemoteRequest { .. }
@@ -1231,7 +1224,7 @@ mod tests {
&["", METRIC_NAME_LABEL, "metric", "job", "api"],
&[missing_name, histogram],
);
let error = decode_wire_error(&wire, false, "conversion error before histogram");
let error = decode_wire_error(&wire, "conversion error before histogram");
assert!(error.to_string().contains("missing '__name__'"));
}
@@ -2026,16 +2019,13 @@ mod tests {
));
}
fn decode_wire(
wire: &[u8],
native_histograms_enabled: bool,
) -> Result<RemoteWriteV2WriteRequests> {
fn decode_wire(wire: &[u8]) -> Result<RemoteWriteV2WriteRequests> {
let body = Bytes::from(crate::prom_store::snappy_compress(wire).unwrap());
decode_remote_write_v2(false, body, native_histograms_enabled)
decode_remote_write_v2(false, body)
}
fn decode_wire_error(wire: &[u8], native_histograms_enabled: bool, name: &str) -> error::Error {
match decode_wire(wire, native_histograms_enabled) {
fn decode_wire_error(wire: &[u8], name: &str) -> error::Error {
match decode_wire(wire) {
Ok(_) => panic!("{name}: expected decoder error"),
Err(error) => error,
}
@@ -2104,16 +2094,9 @@ mod tests {
}
fn decode_test_request(request: Request) -> Result<RemoteWriteV2WriteRequests> {
decode_test_request_with_histograms(request, true)
}
fn decode_test_request_with_histograms(
request: Request,
native_histograms_enabled: bool,
) -> Result<RemoteWriteV2WriteRequests> {
let body =
Bytes::from(crate::prom_store::snappy_compress(&request.encode_to_vec()).unwrap());
decode_remote_write_v2(false, body, native_histograms_enabled)
decode_remote_write_v2(false, body)
}
fn assert_invalid(name: &str, request: Request, expected: &str) {
+3 -65
View File
@@ -209,34 +209,10 @@ fn make_test_app_with_write_capture(
make_test_app_with_write_failure(read_tx, write_tx, None)
}
fn make_test_app_with_native_histogram_write_capture(
read_tx: mpsc::Sender<(String, Vec<u8>)>,
write_tx: mpsc::Sender<RemoteWriteCapture>,
) -> Router {
make_test_app_with_write_failure_inner(read_tx, write_tx, None, true)
}
fn make_test_app_with_write_failure(
read_tx: mpsc::Sender<(String, Vec<u8>)>,
write_tx: mpsc::Sender<RemoteWriteCapture>,
fail_write_call: Option<usize>,
) -> Router {
make_test_app_with_write_failure_inner(read_tx, write_tx, fail_write_call, false)
}
fn make_test_app_with_native_histogram_write_failure(
read_tx: mpsc::Sender<(String, Vec<u8>)>,
write_tx: mpsc::Sender<RemoteWriteCapture>,
fail_write_call: Option<usize>,
) -> Router {
make_test_app_with_write_failure_inner(read_tx, write_tx, fail_write_call, true)
}
fn make_test_app_with_write_failure_inner(
read_tx: mpsc::Sender<(String, Vec<u8>)>,
write_tx: mpsc::Sender<RemoteWriteCapture>,
fail_write_call: Option<usize>,
experimental_enable_prometheus_native_histogram: bool,
) -> Router {
let http_opts = HttpOptions {
addr: format!("127.0.0.1:{}", ports::get_port()),
@@ -262,14 +238,7 @@ fn make_test_app_with_write_failure_inner(
BatchingProtocol::HttpSql,
])
.with_sql_handler(instance.clone())
.with_prom_handler(
instance,
None,
true,
PromValidationMode::Unchecked,
experimental_enable_prometheus_native_histogram,
None,
)
.with_prom_handler(instance, None, true, PromValidationMode::Unchecked, None)
.build();
server.build(server.make_app()).unwrap()
}
@@ -418,7 +387,7 @@ async fn test_prometheus_remote_write_v2_histogram_write_error_has_partial_writt
let (read_tx, _read_rx) = mpsc::channel(100);
let (write_tx, mut write_rx) = mpsc::channel(100);
let app = make_test_app_with_native_histogram_write_failure(read_tx, write_tx, Some(2));
let app = make_test_app_with_write_failure(read_tx, write_tx, Some(2));
let client = TestClient::new(app).await;
let mut write_request = remote_write_v2::request_with_labels_and_samples(
@@ -671,7 +640,7 @@ async fn test_prometheus_remote_write_v2_writes_histogram_only_series() {
let (read_tx, _read_rx) = mpsc::channel(100);
let (write_tx, mut write_rx) = mpsc::channel(100);
let app = make_test_app_with_native_histogram_write_capture(read_tx, write_tx);
let app = make_test_app_with_write_capture(read_tx, write_tx);
let client = TestClient::new(app).await;
let write_request = remote_write_v2::request_with_labels_and_histograms(
@@ -708,36 +677,6 @@ async fn test_prometheus_remote_write_v2_writes_histogram_only_series() {
assert!(write_rx.try_recv().is_err());
}
#[tokio::test]
async fn test_prometheus_remote_write_v2_rejects_native_histogram_when_disabled() {
common_telemetry::init_default_ut_logging();
let (read_tx, _read_rx) = mpsc::channel(100);
let (write_tx, mut write_rx) = mpsc::channel(100);
let app = make_test_app_with_write_capture(read_tx, write_tx);
let client = TestClient::new(app).await;
let write_request = remote_write_v2::request_with_labels_and_histograms(
vec![(
prom_store::METRIC_NAME_LABEL,
"http_request_duration_seconds",
)],
vec![remote_write_v2::histogram(1000)],
);
let result = post_remote_write_v2(&client, &write_request).await;
assert_eq!(result.status(), 400);
assert_remote_write_v2_written_headers_with_histograms(&result.headers(), "0", "0");
assert!(
result
.text()
.await
.contains("prom_store.experimental_enable_prometheus_native_histogram")
);
assert!(write_rx.try_recv().is_err());
}
#[tokio::test]
async fn test_prometheus_remote_write_rejects_unsupported_proto() {
common_telemetry::init_default_ut_logging();
@@ -873,7 +812,6 @@ async fn test_prom_batching_depends_on_protocol_and_metric_engine() {
None,
with_metric_engine,
PromValidationMode::Unchecked,
false,
None,
)
.build();
@@ -77,7 +77,7 @@ fn test_decode_remote_write_v2_native_histogram_dump() {
assert_eq!(histogram.timestamp, 1782358160412);
let (sample_inserts, histogram_inserts, sample_count, histogram_count) =
remote_write_v2::decode_write_requests(false, Bytes::from_static(BODY), true).unwrap();
remote_write_v2::decode_write_requests(false, Bytes::from_static(BODY)).unwrap();
assert!(sample_inserts.is_empty());
assert_eq!(sample_count, 0);
assert_eq!(histogram_count, 1);
-2
View File
@@ -42,7 +42,6 @@ impl ProtocolCtx {
/// If true, all scope attributes will be promoted to the final table schema.
/// Along with the scope name, scope version and scope schema URL.
/// - `with_metric_engine`
/// - `experimental_enable_exponential_histogram`
/// - `is_legacy`
/// If the user uses OTLP metrics ingestion before v0.16, it uses the old path.
/// So we call this path 'legacy'.
@@ -54,7 +53,6 @@ pub struct OtlpMetricCtx {
pub resource_attrs: HashSet<String>,
pub promote_scope_attrs: bool,
pub with_metric_engine: bool,
pub experimental_enable_exponential_histogram: bool,
pub is_legacy: bool,
/// Set from the server's `otlp.experimental_enable_resource_info`; off
/// means the resource descriptor is not synthesized at all.
+1 -2
View File
@@ -535,8 +535,7 @@ WITH(
.build()
.await;
let instance = standalone.fe_instance();
let mut options = standalone.opts.clone();
options.otlp.experimental_enable_exponential_histogram = true;
let options = standalone.opts.clone();
let services = Services::new(options.clone(), instance.clone(), Plugins::default());
let server = services
.http_server_builder(
+8 -39
View File
@@ -701,7 +701,7 @@ pub async fn setup_test_http_app_with_frontend_and_slow_query_threshold(
.with_log_ingest_handler(instance.fe_instance().clone(), None, None)
.with_logs_handler(instance.fe_instance().clone())
.with_influxdb_handler(instance.fe_instance().clone())
.with_otlp_handler(instance.fe_instance().clone(), true, false)
.with_otlp_handler(instance.fe_instance().clone(), true)
.with_jaeger_handler(instance.fe_instance().clone())
.with_greptime_config_options(instance.opts.to_toml().unwrap())
.build();
@@ -721,18 +721,6 @@ pub async fn setup_test_http_app_with_frontend_and_user_provider(
user_provider,
None,
None,
false,
)
.await
}
pub async fn setup_test_http_app_with_otlp_exponential_histogram(
store_type: StorageType,
name: &str,
enabled: bool,
) -> (Router, TestGuard) {
setup_test_http_app_with_frontend_and_custom_options(
store_type, name, None, None, None, enabled,
)
.await
}
@@ -743,7 +731,6 @@ pub async fn setup_test_http_app_with_frontend_and_custom_options(
user_provider: Option<UserProviderRef>,
http_opts: Option<HttpOptions>,
memory_limiter: Option<ServerMemoryLimiter>,
experimental_enable_exponential_histogram: bool,
) -> (Router, TestGuard) {
let plugins = Plugins::new();
if let Some(user_provider) = user_provider.clone() {
@@ -765,11 +752,7 @@ pub async fn setup_test_http_app_with_frontend_and_custom_options(
.with_log_ingest_handler(instance.fe_instance().clone(), None, None)
.with_logs_handler(instance.fe_instance().clone())
.with_influxdb_handler(instance.fe_instance().clone())
.with_otlp_handler(
instance.fe_instance().clone(),
true,
experimental_enable_exponential_histogram,
)
.with_otlp_handler(instance.fe_instance().clone(), true)
.with_prometheus_handler(instance.fe_instance().clone())
.with_jaeger_handler(instance.fe_instance().clone())
.with_dashboard_handler(instance.fe_instance().clone())
@@ -801,14 +784,7 @@ pub async fn setup_test_prom_app_with_frontend(
store_type: StorageType,
name: &str,
) -> (Router, TestGuard) {
setup_test_prom_app_with_frontend_inner(store_type, name, false, false).await
}
pub async fn setup_test_prom_app_with_frontend_native_histogram(
store_type: StorageType,
name: &str,
) -> (Router, TestGuard) {
setup_test_prom_app_with_frontend_inner(store_type, name, false, true).await
setup_test_prom_app_with_frontend_inner(store_type, name, false).await
}
/// Like [`setup_test_prom_app_with_frontend`] but enables the pending-rows batcher,
@@ -818,14 +794,13 @@ pub async fn setup_test_prom_app_with_frontend_batched(
store_type: StorageType,
name: &str,
) -> (Router, TestGuard) {
setup_test_prom_app_with_frontend_inner(store_type, name, true, false).await
setup_test_prom_app_with_frontend_inner(store_type, name, true).await
}
async fn setup_test_prom_app_with_frontend_inner(
store_type: StorageType,
name: &str,
enable_batcher: bool,
experimental_enable_prometheus_native_histogram: bool,
) -> (Router, TestGuard) {
unsafe {
std::env::set_var("TZ", "UTC");
@@ -874,13 +849,9 @@ async fn setup_test_prom_app_with_frontend_inner(
let sql = "INSERT INTO mito(host, val, ts) VALUES (1, 1.1, 0)";
run_sql(sql, &instance).await;
let http_server = build_test_prom_server(
instance.fe_instance().clone(),
enable_batcher,
experimental_enable_prometheus_native_histogram,
)
.with_greptime_config_options(instance.opts.datanode_options().to_toml().unwrap())
.build();
let http_server = build_test_prom_server(instance.fe_instance().clone(), enable_batcher)
.with_greptime_config_options(instance.opts.datanode_options().to_toml().unwrap())
.build();
let app = http_server.build(http_server.make_app()).unwrap();
(app, instance.guard)
}
@@ -889,7 +860,6 @@ async fn setup_test_prom_app_with_frontend_inner(
pub fn build_test_prom_server(
frontend_ref: Arc<Instance>,
enable_batcher: bool,
experimental_enable_prometheus_native_histogram: bool,
) -> HttpServerBuilder {
let http_opts = HttpOptions {
addr: format!("127.0.0.1:{}", ports::get_port()),
@@ -924,7 +894,6 @@ pub fn build_test_prom_server(
Some(frontend_ref.clone()),
true,
PromValidationMode::Strict,
experimental_enable_prometheus_native_histogram,
pending_rows_batcher,
)
.with_prometheus_handler(frontend_ref)
@@ -1279,7 +1248,7 @@ pub async fn setup_pg_server_with_prom_native_histogram(
let instance = setup_standalone_instance(name, store_type).await;
// Prometheus remote-write HTTP app with native histograms enabled.
let http_server = build_test_prom_server(instance.fe_instance().clone(), false, true)
let http_server = build_test_prom_server(instance.fe_instance().clone(), false)
.with_greptime_config_options(instance.opts.datanode_options().to_toml().unwrap())
.build();
let app = http_server.build(http_server.make_app()).unwrap();
@@ -3312,3 +3312,166 @@ CREATE TABLE b (
+-------+-----------------------------------+"#;
check_output_stream(output, expected).await;
}
#[rstest]
#[case::mito(false)]
#[case::batched_metric(true)]
#[tokio::test(flavor = "multi_thread")]
async fn test_histogram_ingestion_storage_lifecycle(#[case] metric_engine: bool) {
use api::greptime_proto::io::prometheus::write::v2::histogram::{Count, ZeroCount};
use api::greptime_proto::io::prometheus::write::v2::{BucketSpan, Histogram, Sample};
use axum::body::{Body, to_bytes};
use http::Request;
use prost::Message;
use servers::http::{HttpOptions, HttpServerBuilder};
use servers::prom_remote_write::v2::test_util as remote_write_v2;
use servers::prom_remote_write::validation::PromValidationMode;
use servers::prom_store::snappy_compress;
use tower::ServiceExt;
use crate::cluster::GreptimeDbClusterBuilder;
use crate::test_util::build_test_prom_server;
use crate::tests::test_util::{MockInstanceBuilder, RebuildableMockInstance, TestContext};
common_telemetry::init_default_ut_logging();
let builder = MockInstanceBuilder::Distributed(
GreptimeDbClusterBuilder::new(&format!("histogram_lifecycle_{metric_engine}"))
.await
.with_datanodes(3),
);
let mut context = TestContext::new(builder).await;
let make_router = |frontend: Arc<Instance>| {
let builder = if metric_engine {
build_test_prom_server(frontend.clone(), true)
} else {
HttpServerBuilder::new(HttpOptions::default())
.with_sql_handler(frontend.clone())
.with_prometheus_handler(frontend.clone())
.with_prom_handler(
frontend.clone(),
None,
false,
PromValidationMode::Strict,
None,
)
};
let server = builder.build();
server.build(server.make_app()).unwrap()
};
let snapshot = |frontend: Arc<Instance>| async move {
let mut result = Vec::new();
for table in ["lifecycle_sample", "lifecycle_histogram"] {
let output = execute_sql(
&frontend,
&format!("select * from {table} order by greptime_timestamp"),
)
.await;
let OutputData::Stream(stream) = output.data else {
panic!("expected query stream")
};
result.push(
util::collect_batches(stream)
.await
.unwrap()
.pretty_print()
.unwrap(),
);
}
result
};
let mut router = make_router(context.frontend());
for round in 0..2 {
let timestamp = 1_000 + round * 10_000;
let mut labels = vec![("__name__", "lifecycle_histogram"), ("job", "api")];
if round == 1 {
// Exercise schema evolution after the first SST already exists.
labels.push(("zone", "east"));
}
let sample = remote_write_v2::request_with_labels_and_samples(
vec![("__name__", "lifecycle_sample"), ("job", "api")],
vec![Sample {
value: 42.0,
timestamp,
..Default::default()
}],
);
let histograms = remote_write_v2::request_with_labels_and_histograms(
labels,
vec![Histogram {
timestamp,
start_timestamp: 500,
count: Some(Count::CountInt(9_007_199_254_740_993)),
zero_count: Some(ZeroCount::ZeroCountInt(1)),
zero_threshold: 0.001,
sum: 12.5,
positive_spans: vec![BucketSpan {
offset: 0,
length: 1,
}],
positive_deltas: vec![9_007_199_254_740_992],
..Default::default()
}],
);
// Create the scalar table first, then add histogram storage to its physical table.
for request in [sample, histograms] {
let response = router
.clone()
.oneshot(
Request::post("/v1/prometheus/write")
.header("Content-Encoding", "snappy")
.header(
"Content-Type",
"application/x-protobuf;proto=io.prometheus.write.v2.Request",
)
.body(Body::from(
snappy_compress(&request.encode_to_vec()).unwrap(),
))
.unwrap(),
)
.await
.unwrap();
assert_eq!(
response.status(),
204,
"{:?}",
to_bytes(response.into_body(), usize::MAX).await.unwrap()
);
}
let expected = snapshot(context.frontend()).await;
assert!(expected[1].contains("9007199254740993"), "{}", expected[1]);
// Rebuild before flushing, including after adding fields to the physical
// schema: recovery must retain the columns filled into bulk writes.
drop(router);
context.rebuild().await;
router = make_router(context.frontend());
assert_eq!(snapshot(context.frontend()).await, expected, "WAL recovery");
let storage_tables = if metric_engine {
vec!["greptime_physical_table"]
} else {
vec!["lifecycle_sample", "lifecycle_histogram"]
};
for table in &storage_tables {
execute_sql(
&context.frontend(),
&format!("admin flush_table('{table}')"),
)
.await;
}
assert_eq!(snapshot(context.frontend()).await, expected, "SST read");
if round == 1 {
for table in &storage_tables {
execute_sql(
&context.frontend(),
&format!("admin compact_table('{table}', 'strict_window', 'window=3600')"),
)
.await;
}
assert_eq!(snapshot(context.frontend()).await, expected, "compaction");
drop(router);
context.rebuild().await;
router = make_router(context.frontend());
assert_eq!(snapshot(context.frontend()).await, expected, "SST reopen");
}
}
}
+7 -54
View File
@@ -85,7 +85,7 @@ 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,
setup_test_http_app_with_frontend_and_user_provider, setup_test_prom_app_with_frontend,
setup_test_prom_app_with_frontend_batched, setup_test_prom_app_with_frontend_native_histogram,
setup_test_prom_app_with_frontend_batched,
};
use urlencoding::encode;
use yaml_rust::YamlLoader;
@@ -2492,7 +2492,6 @@ enable = true
[otlp]
enable = true
experimental_enable_exponential_histogram = false
trace_ingest_chunk_size = 512
experimental_enable_resource_info = true
@@ -2500,7 +2499,6 @@ experimental_enable_resource_info = true
enable = true
with_metric_engine = true
prom_validation_mode = "strict"
experimental_enable_prometheus_native_histogram = false
pending_rows_flush_interval = "0s"
max_batch_rows = 100000
max_concurrent_flushes = 256
@@ -3152,7 +3150,7 @@ pub async fn test_prometheus_remote_write_v2(store_type: StorageType) {
pub async fn test_prometheus_remote_write_v2_native_histogram(store_type: StorageType) {
common_telemetry::init_default_ut_logging();
let (app, mut guard) = setup_test_prom_app_with_frontend_native_histogram(
let (app, mut guard) = setup_test_prom_app_with_frontend(
store_type,
"prometheus_remote_write_v2_native_histogram",
)
@@ -3532,7 +3530,7 @@ async fn check_prometheus_remote_write_batched_skip_wal(distributed: bool, v2: b
common_telemetry::init_default_ut_logging();
let mut instance =
MockInstanceImpl::new(&format!("prom_bulk_skip_wal_v2_{v2}"), distributed).await;
let server = build_test_prom_server(instance.frontend(), true, false).build();
let server = build_test_prom_server(instance.frontend(), true).build();
let client = TestClient::new(server.build(server.make_app()).unwrap()).await;
write_prometheus_skip_wal_sample(&client, v2, 1000, None).await;
@@ -7464,7 +7462,6 @@ pub async fn test_otlp_exponential_histogram(store_type: StorageType) {
AggregationTemporality, ExponentialHistogram, ExponentialHistogramDataPoint, Metric,
ResourceMetrics, ScopeMetrics, exponential_histogram_data_point, metric,
};
use tests_integration::test_util::setup_test_http_app_with_otlp_exponential_histogram;
common_telemetry::init_default_ut_logging();
let req = ExportMetricsServiceRequest {
@@ -7511,44 +7508,8 @@ pub async fn test_otlp_exponential_histogram(store_type: StorageType) {
)]
};
let (app, mut guard) = setup_test_http_app_with_otlp_exponential_histogram(
store_type,
"test_otlp_exponential_histogram_disabled",
false,
)
.await;
let client = TestClient::new(app).await;
let res = send_req(
&client,
headers(),
"/v1/otlp/v1/metrics",
body.clone(),
false,
)
.await;
assert_eq!(StatusCode::BAD_REQUEST, res.status());
let status = GoogleRpcStatus::decode(res.bytes().await.as_ref()).unwrap();
assert_eq!(3, status.code);
assert!(
status
.message
.contains("otlp.experimental_enable_exponential_histogram")
);
validate_data(
"otlp_exponential_histogram_disabled_no_table",
&client,
"select count(*) from information_schema.tables where table_name = 'otlp_exponential_latency';",
"[[0]]",
)
.await;
guard.remove_all().await;
let (app, mut guard) = setup_test_http_app_with_otlp_exponential_histogram(
store_type,
"test_otlp_exponential_histogram_enabled",
true,
)
.await;
let (app, mut guard) =
setup_test_http_app_with_frontend(store_type, "test_otlp_exponential_histogram").await;
let client = TestClient::new(app).await;
let res = send_req(&client, headers(), "/v1/otlp/v1/metrics", body, false).await;
assert_eq!(StatusCode::OK, res.status());
@@ -12041,16 +12002,9 @@ async fn check_http_skip_wal(name: &str, cases: &[HttpWalCase], distributed: boo
.with_influxdb_handler(fe.clone())
.with_opentsdb_handler(fe.clone())
.with_log_ingest_handler(fe.clone(), None, None)
.with_otlp_handler(fe.clone(), true, false)
.with_otlp_handler(fe.clone(), true)
// The pending batcher uses BulkInsert, deliberately outside this PR.
.with_prom_handler(
fe.clone(),
Some(fe),
true,
PromValidationMode::Strict,
false,
None,
)
.with_prom_handler(fe.clone(), Some(fe), true, PromValidationMode::Strict, None)
.build();
let client = TestClient::new(server.build(server.make_app()).unwrap()).await;
for case in cases {
@@ -12395,7 +12349,6 @@ pub async fn test_http_memory_limit(store_type: StorageType) {
None,
Some(http_opts),
Some(memory_limiter),
false,
)
.await;