chore: make OTLP trace ingest chunk size configurable (#8455)

* chore: expose trace chunk to config

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

* chore: change default value to 128

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

* fix: cr issue

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

---------

Signed-off-by: shuiyisong <xixing.sys@gmail.com>
This commit is contained in:
shuiyisong
2026-07-09 08:22:40 +00:00
committed by GitHub
parent 6ae687dc8e
commit 7764d2f054
9 changed files with 77 additions and 13 deletions
Generated
-1
View File
@@ -5500,7 +5500,6 @@ dependencies = [
"humantime",
"humantime-serde",
"hyper-util",
"itertools 0.14.0",
"lazy_static",
"log-query",
"meta-client",
+6
View File
@@ -71,6 +71,9 @@
| `influxdb.default_merge_mode` | String | `last_non_null` | Default merge mode for tables automatically created by InfluxDB protocol.<br/>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.<br/>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. |
+7
View File
@@ -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.
+7
View File
@@ -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.
-1
View File
@@ -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
+1
View File
@@ -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<AtomicBool>,
// cache for otlp metrics
+1
View File
@@ -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)),
})
}
+31 -10
View File
@@ -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::<Vec<_>>())
.collect::<Vec<_>>();
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<T>(items: Vec<T>, chunk_size: usize) -> Vec<Vec<T>> {
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<E>(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<_>>(), 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::<i32>(Vec::new(), 0).is_empty());
}
#[test]
fn test_classify_trace_chunk_failure() {
assert_eq!(
+24 -1
View File
@@ -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);
}
}