diff --git a/Cargo.lock b/Cargo.lock index 122d2a62433..e00e30e28a6 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5500,7 +5500,6 @@ dependencies = [ "humantime", "humantime-serde", "hyper-util", - "itertools 0.14.0", "lazy_static", "log-query", "meta-client", diff --git a/config/config.md b/config/config.md index 3a6e187375f..3a0c1729af5 100644 --- a/config/config.md +++ b/config/config.md @@ -71,6 +71,9 @@ | `influxdb.default_merge_mode` | String | `last_non_null` | Default merge mode for tables automatically created by InfluxDB protocol.
Available values: "last_non_null", "last_row". | | `jaeger` | -- | -- | Jaeger protocol options. | | `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.trace_ingest_chunk_size` | Integer | `128` | Maximum spans per trace ingest chunk. Set to 0 to disable splitting. | | `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. | @@ -300,6 +303,9 @@ | `influxdb.default_merge_mode` | String | `last_non_null` | Default merge mode for tables automatically created by InfluxDB protocol.
Available values: "last_non_null", "last_row". | | `jaeger` | -- | -- | Jaeger protocol options. | | `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.trace_ingest_chunk_size` | Integer | `128` | Maximum spans per trace ingest chunk. Set to 0 to disable splitting. | | `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. | diff --git a/config/frontend.example.toml b/config/frontend.example.toml index badba500623..9fc65ece8af 100644 --- a/config/frontend.example.toml +++ b/config/frontend.example.toml @@ -219,6 +219,13 @@ default_merge_mode = "last_non_null" ## Whether to enable Jaeger protocol in HTTP API. enable = true +## OpenTelemetry protocol options. +[otlp] +## Whether to enable OpenTelemetry protocol in HTTP API. +enable = true +## Maximum spans per trace ingest chunk. Set to 0 to disable splitting. +trace_ingest_chunk_size = 128 + ## Prometheus remote storage options [prom_store] ## Whether to enable Prometheus remote write and read in HTTP API. diff --git a/config/standalone.example.toml b/config/standalone.example.toml index de3fa7a49d8..074b827a2b5 100644 --- a/config/standalone.example.toml +++ b/config/standalone.example.toml @@ -186,6 +186,13 @@ default_merge_mode = "last_non_null" ## Whether to enable Jaeger protocol in HTTP API. enable = true +## OpenTelemetry protocol options. +[otlp] +## Whether to enable OpenTelemetry protocol in HTTP API. +enable = true +## Maximum spans per trace ingest chunk. Set to 0 to disable splitting. +trace_ingest_chunk_size = 128 + ## Prometheus remote storage options [prom_store] ## Whether to enable Prometheus remote write and read in HTTP API. diff --git a/src/frontend/Cargo.toml b/src/frontend/Cargo.toml index aa9276a0cba..3e18ff1d3d7 100644 --- a/src/frontend/Cargo.toml +++ b/src/frontend/Cargo.toml @@ -52,7 +52,6 @@ futures.workspace = true hostname.workspace = true humantime.workspace = true humantime-serde.workspace = true -itertools.workspace = true lazy_static.workspace = true log-query.workspace = true meta-client.workspace = true diff --git a/src/frontend/src/instance.rs b/src/frontend/src/instance.rs index 8a6bc062f25..7242d601c1b 100644 --- a/src/frontend/src/instance.rs +++ b/src/frontend/src/instance.rs @@ -124,6 +124,7 @@ pub struct Instance { process_manager: ProcessManagerRef, slow_query_options: SlowQueryOptions, influxdb_default_merge_mode: InfluxdbMergeMode, + trace_ingest_chunk_size: usize, suspend: Arc, // cache for otlp metrics diff --git a/src/frontend/src/instance/builder.rs b/src/frontend/src/instance/builder.rs index 37bfc3013bf..a9e7968219c 100644 --- a/src/frontend/src/instance/builder.rs +++ b/src/frontend/src/instance/builder.rs @@ -314,6 +314,7 @@ impl FrontendBuilder { otlp_metrics_table_legacy_cache: DashMap::new(), slow_query_options: self.options.slow_query.clone(), influxdb_default_merge_mode: self.options.influxdb.default_merge_mode, + trace_ingest_chunk_size: self.options.otlp.trace_ingest_chunk_size, suspend: Arc::new(AtomicBool::new(false)), }) } diff --git a/src/frontend/src/instance/otlp.rs b/src/frontend/src/instance/otlp.rs index 9672b05301b..16e3b996064 100644 --- a/src/frontend/src/instance/otlp.rs +++ b/src/frontend/src/instance/otlp.rs @@ -26,7 +26,6 @@ use common_error::ext::{BoxedError, ErrorExt}; use common_error::status_code::StatusCode; use common_query::prelude::GREPTIME_PHYSICAL_TABLE; use common_telemetry::{tracing, warn}; -use itertools::Itertools; use opentelemetry_proto::tonic::collector::logs::v1::ExportLogsServiceRequest; use opentelemetry_proto::tonic::collector::trace::v1::ExportTraceServiceRequest; use otel_arrow_rust::proto::opentelemetry::collector::metrics::v1::ExportMetricsServiceRequest; @@ -63,7 +62,6 @@ use crate::metrics::{ pub mod trace_semconv; pub mod trace_types; -const TRACE_INGEST_CHUNK_SIZE: usize = 64; const TRACE_FAILURE_MESSAGE_LIMIT: usize = 4; #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -311,13 +309,7 @@ impl Instance { }; for group in groups { - let chunks = group - .spans - .into_iter() - .chunks(TRACE_INGEST_CHUNK_SIZE) - .into_iter() - .map(|chunk| chunk.collect::>()) - .collect::>(); + let chunks = chunk_owned(group.spans, self.trace_ingest_chunk_size); for chunk in chunks { self.ingest_trace_chunk(&ingest_ctx, chunk, main_ctx.clone(), &mut ingest_state) .await?; @@ -843,6 +835,23 @@ impl Instance { } } +fn chunk_owned(items: Vec, chunk_size: usize) -> Vec> { + if items.is_empty() { + return Vec::new(); + } + + if chunk_size == 0 { + return vec![items]; + } + + let mut chunks = Vec::with_capacity(items.len().div_ceil(chunk_size)); + let mut iter = items.into_iter(); + while iter.len() > 0 { + chunks.push(iter.by_ref().take(chunk_size).collect()); + } + chunks +} + /// Preserve the original alter failure status so chunk retry behavior stays correct. fn wrap_trace_alter_failure(err: E) -> servers::error::Error where @@ -898,9 +907,21 @@ mod tests { use common_error::status_code::StatusCode; use servers::query_handler::TraceIngestOutcome; - use super::{ChunkFailureReaction, Instance, wrap_trace_alter_failure}; + use super::{ChunkFailureReaction, Instance, chunk_owned, wrap_trace_alter_failure}; use crate::metrics::OTLP_TRACES_FAILURE_COUNT; + #[test] + fn test_chunk_owned() { + let chunks = chunk_owned(vec![1, 2, 3], 2); + assert_eq!(chunks.iter().map(Vec::len).collect::>(), vec![2, 1]); + + let chunks = chunk_owned(vec![1, 2, 3], 0); + assert_eq!(chunks.len(), 1); + assert_eq!(chunks[0].len(), 3); + + assert!(chunk_owned::(Vec::new(), 0).is_empty()); + } + #[test] fn test_classify_trace_chunk_failure() { assert_eq!( diff --git a/src/frontend/src/service_config/otlp.rs b/src/frontend/src/service_config/otlp.rs index 76ba465afb8..0c1efb7b80f 100644 --- a/src/frontend/src/service_config/otlp.rs +++ b/src/frontend/src/service_config/otlp.rs @@ -14,14 +14,22 @@ use serde::{Deserialize, Serialize}; +const DEFAULT_TRACE_INGEST_CHUNK_SIZE: usize = 128; + #[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)] +#[serde(default)] pub struct OtlpOptions { pub enable: bool, + /// Maximum spans per trace ingest chunk. Set to 0 to disable splitting. + pub trace_ingest_chunk_size: usize, } impl Default for OtlpOptions { fn default() -> Self { - Self { enable: true } + Self { + enable: true, + trace_ingest_chunk_size: DEFAULT_TRACE_INGEST_CHUNK_SIZE, + } } } @@ -33,5 +41,20 @@ mod tests { fn test_otlp_options() { let default = OtlpOptions::default(); assert!(default.enable); + assert_eq!( + default.trace_ingest_chunk_size, + DEFAULT_TRACE_INGEST_CHUNK_SIZE + ); + + let options: OtlpOptions = toml::from_str("enable = false").unwrap(); + assert!(!options.enable); + assert_eq!( + options.trace_ingest_chunk_size, + DEFAULT_TRACE_INGEST_CHUNK_SIZE + ); + + let options: OtlpOptions = toml::from_str("trace_ingest_chunk_size = 0").unwrap(); + assert!(options.enable); + assert_eq!(options.trace_ingest_chunk_size, 0); } }