feat: share logical table batching with OTLP metrics (#9288)

* feat: share logical table batching with OTLP metrics

Signed-off-by: WenyXu <wenymedia@gmail.com>

* fix: unify pending rows batch acknowledgement policy

Signed-off-by: WenyXu <wenymedia@gmail.com>

* fix: align logical batcher example configuration expectations

Signed-off-by: WenyXu <wenymedia@gmail.com>

* fix: align batcher worker channel defaults to 65536

Signed-off-by: WenyXu <wenymedia@gmail.com>

---------

Signed-off-by: WenyXu <wenymedia@gmail.com>
This commit is contained in:
Weny Xu
2026-09-22 14:38:17 +00:00
committed by GitHub
parent d0f8f4b80c
commit 723da69b21
28 changed files with 1518 additions and 237 deletions
+20 -18
View File
@@ -34,7 +34,7 @@
| `runtime.experimental_workload_scheduler.sample_every_polls` | Integer | `16` | Number of polls between scheduler fairness samples. Must be greater than zero. |
| `http` | -- | -- | The HTTP server options. |
| `http.addr` | String | `127.0.0.1:4000` | The address to bind the HTTP server. |
| `http.timeout` | String | `0s` | HTTP request timeout. Set to 0 to disable timeout.<br/>When synchronous Prometheus or shared table batching is enabled, a nonzero timeout is<br/>raised to at least the largest active flush interval plus 1 second. The intervals come from<br/>`prom_store.pending_rows_flush_interval` and `pending_rows_batcher.pending_rows_flush_interval`. |
| `http.timeout` | String | `0s` | HTTP request timeout. Set to 0 to disable timeout.<br/>When synchronous Prometheus, OTLP metrics, or ordinary-table batching is enabled, a nonzero timeout is<br/>raised to at least the largest active flush interval plus 1 second. The intervals come from<br/>`pending_rows_batcher.logical_table.pending_rows_flush_interval` and `pending_rows_batcher.pending_rows_flush_interval`. |
| `http.body_limit` | String | `64MB` | HTTP request body limit.<br/>The following units are supported: `B`, `KB`, `KiB`, `MB`, `MiB`, `GB`, `GiB`, `TB`, `TiB`, `PB`, `PiB`.<br/>Set to 0 to disable limit. |
| `http.enable_cors` | Bool | `true` | HTTP CORS support, it's turned on by default<br/>This allows browser to access http APIs without CORS restrictions |
| `http.cors_allowed_origins` | Array | Unset | Customize allowed origins for HTTP CORS. |
@@ -75,13 +75,20 @@
| `influxdb` | -- | -- | InfluxDB protocol options. |
| `influxdb.enable` | Bool | `true` | Whether to enable InfluxDB protocol in HTTP API. |
| `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". |
| `pending_rows_batcher` | -- | -- | Shared experimental ordinary-table batching for opted-in ingestion protocols.<br/>Legacy Prometheus batching settings under prom_store remain supported.<br/>HTTP write protocols sharing this batcher. Omitted or empty disables all entrances.<br/>Supported: influxdb, opentsdb, otlp, logs, loki, splunk, elasticsearch, http_sql, prom.<br/>Prom uses ordinary-table batching without metric engine, otherwise its dedicated batcher.<br/>Effective shared Prom settings take precedence; existing prom_store settings remain compatible. |
| `pending_rows_batcher` | -- | -- | Ordinary-table batching for opted-in HTTP ingestion protocols.<br/>PENDING_ROWS_BATCH_SYNC defaults to true for both batchers. Set it to false to acknowledge<br/>queue admission without waiting for storage; later failures cannot be returned to the client.<br/>Omitted or empty protocols disables batching. Prom without metric engine uses this batcher.<br/>OTLP logs, traces and ordinary metrics use this batcher. |
| `pending_rows_batcher.pending_rows_flush_interval` | String | `0s` | Flush interval measured from the first pending submission. Zero disables batching. |
| `pending_rows_batcher.max_batch_rows` | Integer | `100000` | Flush after a complete submission reaches this row threshold. |
| `pending_rows_batcher.max_concurrent_flushes` | Integer | `256` | Maximum concurrent flushes shared by the frontend batcher. |
| `pending_rows_batcher.worker_channel_capacity` | Integer | `65526` | Maximum queued submissions per table worker. |
| `pending_rows_batcher.worker_channel_capacity` | Integer | `65536` | Maximum queued submissions per table worker. |
| `pending_rows_batcher.max_inflight_requests` | Integer | `3000` | Maximum admitted original requests awaiting completion. |
| `pending_rows_batcher.flow_notification_queue_capacity` | Integer | `1024` | Maximum number of queued table Flow notifications. |
| `pending_rows_batcher.logical_table` | -- | -- | Metric-engine logical-table batching for Prom remote write and non-legacy OTLP metrics.<br/>Requires prom_store.with_metric_engine. Logs, traces and legacy metrics are not eligible.<br/>Enable independently with protocols and a nonzero flush interval.<br/>Omitted fields use independent defaults, not parent settings.<br/>Empty protocols or a zero interval disables logical batching without fallback.<br/>Omitting this entire section preserves legacy Prom batching; it does not enable OTLP batching. |
| `pending_rows_batcher.logical_table.pending_rows_flush_interval` | String | `0s` | Flush interval measured from the first pending submission. Zero disables batching. |
| `pending_rows_batcher.logical_table.max_batch_rows` | Integer | `100000` | Flush after a complete submission reaches this row threshold. |
| `pending_rows_batcher.logical_table.max_concurrent_flushes` | Integer | `256` | Maximum concurrent flushes shared by Prom and OTLP metrics. |
| `pending_rows_batcher.logical_table.worker_channel_capacity` | Integer | `65536` | Maximum queued submissions per physical-table worker. |
| `pending_rows_batcher.logical_table.max_inflight_requests` | Integer | `3000` | Maximum admitted original requests awaiting completion. |
| `pending_rows_batcher.logical_table.flow_notification_queue_capacity` | Integer | `1024` | Maximum number of queued logical-table Flow notifications. |
| `jaeger` | -- | -- | Jaeger protocol options. |
| `jaeger.enable` | Bool | `true` | Whether to enable Jaeger protocol in HTTP API. |
| `otlp` | -- | -- | OpenTelemetry protocol options. |
@@ -94,12 +101,6 @@
| `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. |
| `prom_store.pending_rows_flush_interval` | String | `0s` | Interval to flush pending rows batcher.<br/>Set to "0s" to disable batching mode in Prometheus Remote Write endpoint |
| `prom_store.max_batch_rows` | Integer | `100000` | Max rows per pending batch before triggering a flush. |
| `prom_store.max_concurrent_flushes` | Integer | `256` | Max number of concurrent batch flushes. |
| `prom_store.worker_channel_capacity` | Integer | `65526` | Capacity of the pending batch worker channel. |
| `prom_store.max_inflight_requests` | Integer | `3000` | Max inflight write requests before backpressure. |
| `prom_store.flow_notification_queue_capacity` | Integer | `1024` | Maximum number of logical-table flow notifications waiting in the shared queue. |
| `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`**. |
@@ -289,7 +290,7 @@
| `runtime.compact_rt_max_blocking_threads` | Integer | `4` | The maximum number of blocking threads for compact operations.<br/>Defaults to max(num_cpus / 2, 2). |
| `http` | -- | -- | The HTTP server options. |
| `http.addr` | String | `127.0.0.1:4000` | The address to bind the HTTP server. |
| `http.timeout` | String | `0s` | HTTP request timeout. Set to 0 to disable timeout.<br/>When synchronous Prometheus or shared table batching is enabled, a nonzero timeout is<br/>raised to at least the largest active flush interval plus 1 second. The intervals come from<br/>`prom_store.pending_rows_flush_interval` and `pending_rows_batcher.pending_rows_flush_interval`. |
| `http.timeout` | String | `0s` | HTTP request timeout. Set to 0 to disable timeout.<br/>When synchronous Prometheus, OTLP metrics, or ordinary-table batching is enabled, a nonzero timeout is<br/>raised to at least the largest active flush interval plus 1 second. The intervals come from<br/>`pending_rows_batcher.logical_table.pending_rows_flush_interval` and `pending_rows_batcher.pending_rows_flush_interval`. |
| `http.body_limit` | String | `64MB` | HTTP request body limit.<br/>The following units are supported: `B`, `KB`, `KiB`, `MB`, `MiB`, `GB`, `GiB`, `TB`, `TiB`, `PB`, `PiB`.<br/>Set to 0 to disable limit. |
| `http.enable_cors` | Bool | `true` | HTTP CORS support, it's turned on by default<br/>This allows browser to access http APIs without CORS restrictions |
| `http.cors_allowed_origins` | Array | Unset | Customize allowed origins for HTTP CORS. |
@@ -342,13 +343,20 @@
| `influxdb` | -- | -- | InfluxDB protocol options. |
| `influxdb.enable` | Bool | `true` | Whether to enable InfluxDB protocol in HTTP API. |
| `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". |
| `pending_rows_batcher` | -- | -- | Shared experimental ordinary-table batching for opted-in ingestion protocols.<br/>Legacy Prometheus batching settings under prom_store remain supported.<br/>HTTP write protocols sharing this batcher. Omitted or empty disables all entrances.<br/>Supported: influxdb, opentsdb, otlp, logs, loki, splunk, elasticsearch, http_sql, prom.<br/>Prom uses ordinary-table batching without metric engine, otherwise its dedicated batcher.<br/>Effective shared Prom settings take precedence; existing prom_store settings remain compatible. |
| `pending_rows_batcher` | -- | -- | Ordinary-table batching for opted-in HTTP ingestion protocols.<br/>PENDING_ROWS_BATCH_SYNC defaults to true for both batchers. Set it to false to acknowledge<br/>queue admission without waiting for storage; later failures cannot be returned to the client.<br/>Omitted or empty protocols disables batching. Prom without metric engine uses this batcher.<br/>OTLP logs, traces and ordinary metrics use this batcher. |
| `pending_rows_batcher.pending_rows_flush_interval` | String | `0s` | Flush interval measured from the first pending submission. Zero disables batching. |
| `pending_rows_batcher.max_batch_rows` | Integer | `100000` | Flush after a complete submission reaches this row threshold. |
| `pending_rows_batcher.max_concurrent_flushes` | Integer | `256` | Maximum concurrent flushes shared by the frontend batcher. |
| `pending_rows_batcher.worker_channel_capacity` | Integer | `65526` | Maximum queued submissions per table worker. |
| `pending_rows_batcher.worker_channel_capacity` | Integer | `65536` | Maximum queued submissions per table worker. |
| `pending_rows_batcher.max_inflight_requests` | Integer | `3000` | Maximum admitted original requests awaiting completion. |
| `pending_rows_batcher.flow_notification_queue_capacity` | Integer | `1024` | Maximum number of queued table Flow notifications. |
| `pending_rows_batcher.logical_table` | -- | -- | Metric-engine logical-table batching for Prom remote write and non-legacy OTLP metrics.<br/>Requires prom_store.with_metric_engine. Logs, traces and legacy metrics are not eligible.<br/>Enable independently with protocols and a nonzero flush interval.<br/>Omitted fields use independent defaults, not parent settings.<br/>Empty protocols or a zero interval disables logical batching without fallback.<br/>Omitting this entire section preserves legacy Prom batching; it does not enable OTLP batching. |
| `pending_rows_batcher.logical_table.pending_rows_flush_interval` | String | `0s` | Flush interval measured from the first pending submission. Zero disables batching. |
| `pending_rows_batcher.logical_table.max_batch_rows` | Integer | `100000` | Flush after a complete submission reaches this row threshold. |
| `pending_rows_batcher.logical_table.max_concurrent_flushes` | Integer | `256` | Maximum concurrent flushes shared by Prom and OTLP metrics. |
| `pending_rows_batcher.logical_table.worker_channel_capacity` | Integer | `65536` | Maximum queued submissions per physical-table worker. |
| `pending_rows_batcher.logical_table.max_inflight_requests` | Integer | `3000` | Maximum admitted original requests awaiting completion. |
| `pending_rows_batcher.logical_table.flow_notification_queue_capacity` | Integer | `1024` | Maximum number of queued logical-table Flow notifications. |
| `jaeger` | -- | -- | Jaeger protocol options. |
| `jaeger.enable` | Bool | `true` | Whether to enable Jaeger protocol in HTTP API. |
| `otlp` | -- | -- | OpenTelemetry protocol options. |
@@ -361,12 +369,6 @@
| `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. |
| `prom_store.pending_rows_flush_interval` | String | `0s` | Interval to flush pending rows batcher.<br/>Set to "0s" to disable batching mode in Prometheus Remote Write endpoint |
| `prom_store.max_batch_rows` | Integer | `100000` | Max rows per pending batch before triggering a flush. |
| `prom_store.max_concurrent_flushes` | Integer | `256` | Max number of concurrent batch flushes. |
| `prom_store.worker_channel_capacity` | Integer | `65526` | Capacity of the pending batch worker channel. |
| `prom_store.max_inflight_requests` | Integer | `3000` | Max inflight write requests before backpressure. |
| `prom_store.flow_notification_queue_capacity` | Integer | `1024` | Maximum number of logical-table flow notifications waiting in the shared queue. |
| `meta_client` | -- | -- | The metasrv client options. |
| `meta_client.metasrv_addrs` | Array | -- | The addresses of the metasrv. |
| `meta_client.timeout` | String | `3s` | Operation timeout. |
+30 -22
View File
@@ -56,9 +56,9 @@ default_column_prefix = "greptime"
## The address to bind the HTTP server.
addr = "127.0.0.1:4000"
## HTTP request timeout. Set to 0 to disable timeout.
## When synchronous Prometheus or shared table batching is enabled, a nonzero timeout is
## When synchronous Prometheus, OTLP metrics, or ordinary-table batching is enabled, a nonzero timeout is
## raised to at least the largest active flush interval plus 1 second. The intervals come from
## `prom_store.pending_rows_flush_interval` and `pending_rows_batcher.pending_rows_flush_interval`.
## `pending_rows_batcher.logical_table.pending_rows_flush_interval` and `pending_rows_batcher.pending_rows_flush_interval`.
timeout = "0s"
## HTTP request body limit.
## The following units are supported: `B`, `KB`, `KiB`, `MB`, `MiB`, `GB`, `GiB`, `TB`, `TiB`, `PB`, `PiB`.
@@ -231,12 +231,11 @@ enable = true
## Available values: "last_non_null", "last_row".
default_merge_mode = "last_non_null"
## Shared experimental ordinary-table batching for opted-in ingestion protocols.
## Legacy Prometheus batching settings under prom_store remain supported.
## HTTP write protocols sharing this batcher. Omitted or empty disables all entrances.
## Supported: influxdb, opentsdb, otlp, logs, loki, splunk, elasticsearch, http_sql, prom.
## Prom uses ordinary-table batching without metric engine, otherwise its dedicated batcher.
## Effective shared Prom settings take precedence; existing prom_store settings remain compatible.
## Ordinary-table batching for opted-in HTTP ingestion protocols.
## PENDING_ROWS_BATCH_SYNC defaults to true for both batchers. Set it to false to acknowledge
## queue admission without waiting for storage; later failures cannot be returned to the client.
## Omitted or empty protocols disables batching. Prom without metric engine uses this batcher.
## OTLP logs, traces and ordinary metrics use this batcher.
[pending_rows_batcher]
# protocols = [
# "influxdb",
@@ -256,12 +255,33 @@ max_batch_rows = 100000
## Maximum concurrent flushes shared by the frontend batcher.
max_concurrent_flushes = 256
## Maximum queued submissions per table worker.
worker_channel_capacity = 65526
worker_channel_capacity = 65536
## Maximum admitted original requests awaiting completion.
max_inflight_requests = 3000
## Maximum number of queued table Flow notifications.
flow_notification_queue_capacity = 1024
## Metric-engine logical-table batching for Prom remote write and non-legacy OTLP metrics.
## Requires prom_store.with_metric_engine. Logs, traces and legacy metrics are not eligible.
## Enable independently with protocols and a nonzero flush interval.
## Omitted fields use independent defaults, not parent settings.
## Empty protocols or a zero interval disables logical batching without fallback.
## Omitting this entire section preserves legacy Prom batching; it does not enable OTLP batching.
[pending_rows_batcher.logical_table]
# protocols = ["prom", "otlp"]
## Flush interval measured from the first pending submission. Zero disables batching.
pending_rows_flush_interval = "0s"
## Flush after a complete submission reaches this row threshold.
max_batch_rows = 100000
## Maximum concurrent flushes shared by Prom and OTLP metrics.
max_concurrent_flushes = 256
## Maximum queued submissions per physical-table worker.
worker_channel_capacity = 65536
## Maximum admitted original requests awaiting completion.
max_inflight_requests = 3000
## Maximum number of queued logical-table Flow notifications.
flow_notification_queue_capacity = 1024
## Jaeger protocol options.
[jaeger]
## Whether to enable Jaeger protocol in HTTP API.
@@ -293,19 +313,7 @@ with_metric_engine = true
prom_validation_mode = "strict"
## Experimental: enable Prometheus remote write v2 native histogram ingestion.
experimental_enable_prometheus_native_histogram = false
## Interval to flush pending rows batcher.
## Set to "0s" to disable batching mode in Prometheus Remote Write endpoint
#+pending_rows_flush_interval = "0s"
## Max rows per pending batch before triggering a flush.
#+max_batch_rows = 100000
## Max number of concurrent batch flushes.
#+max_concurrent_flushes = 256
## Capacity of the pending batch worker channel.
#+worker_channel_capacity = 65526
## Max inflight write requests before backpressure.
#+max_inflight_requests = 3000
## Maximum number of logical-table flow notifications waiting in the shared queue.
#+flow_notification_queue_capacity = 1024
## The metasrv client options.
[meta_client]
+30 -22
View File
@@ -81,9 +81,9 @@ max_concurrent_queries = 0
## The address to bind the HTTP server.
addr = "127.0.0.1:4000"
## HTTP request timeout. Set to 0 to disable timeout.
## When synchronous Prometheus or shared table batching is enabled, a nonzero timeout is
## When synchronous Prometheus, OTLP metrics, or ordinary-table batching is enabled, a nonzero timeout is
## raised to at least the largest active flush interval plus 1 second. The intervals come from
## `prom_store.pending_rows_flush_interval` and `pending_rows_batcher.pending_rows_flush_interval`.
## `pending_rows_batcher.logical_table.pending_rows_flush_interval` and `pending_rows_batcher.pending_rows_flush_interval`.
timeout = "0s"
## HTTP request body limit.
## The following units are supported: `B`, `KB`, `KiB`, `MB`, `MiB`, `GB`, `GiB`, `TB`, `TiB`, `PB`, `PiB`.
@@ -210,12 +210,11 @@ enable = true
## Available values: "last_non_null", "last_row".
default_merge_mode = "last_non_null"
## Shared experimental ordinary-table batching for opted-in ingestion protocols.
## Legacy Prometheus batching settings under prom_store remain supported.
## HTTP write protocols sharing this batcher. Omitted or empty disables all entrances.
## Supported: influxdb, opentsdb, otlp, logs, loki, splunk, elasticsearch, http_sql, prom.
## Prom uses ordinary-table batching without metric engine, otherwise its dedicated batcher.
## Effective shared Prom settings take precedence; existing prom_store settings remain compatible.
## Ordinary-table batching for opted-in HTTP ingestion protocols.
## PENDING_ROWS_BATCH_SYNC defaults to true for both batchers. Set it to false to acknowledge
## queue admission without waiting for storage; later failures cannot be returned to the client.
## Omitted or empty protocols disables batching. Prom without metric engine uses this batcher.
## OTLP logs, traces and ordinary metrics use this batcher.
[pending_rows_batcher]
# protocols = [
# "influxdb",
@@ -235,12 +234,33 @@ max_batch_rows = 100000
## Maximum concurrent flushes shared by the frontend batcher.
max_concurrent_flushes = 256
## Maximum queued submissions per table worker.
worker_channel_capacity = 65526
worker_channel_capacity = 65536
## Maximum admitted original requests awaiting completion.
max_inflight_requests = 3000
## Maximum number of queued table Flow notifications.
flow_notification_queue_capacity = 1024
## Metric-engine logical-table batching for Prom remote write and non-legacy OTLP metrics.
## Requires prom_store.with_metric_engine. Logs, traces and legacy metrics are not eligible.
## Enable independently with protocols and a nonzero flush interval.
## Omitted fields use independent defaults, not parent settings.
## Empty protocols or a zero interval disables logical batching without fallback.
## Omitting this entire section preserves legacy Prom batching; it does not enable OTLP batching.
[pending_rows_batcher.logical_table]
# protocols = ["prom", "otlp"]
## Flush interval measured from the first pending submission. Zero disables batching.
pending_rows_flush_interval = "0s"
## Flush after a complete submission reaches this row threshold.
max_batch_rows = 100000
## Maximum concurrent flushes shared by Prom and OTLP metrics.
max_concurrent_flushes = 256
## Maximum queued submissions per physical-table worker.
worker_channel_capacity = 65536
## Maximum admitted original requests awaiting completion.
max_inflight_requests = 3000
## Maximum number of queued logical-table Flow notifications.
flow_notification_queue_capacity = 1024
## Jaeger protocol options.
[jaeger]
## Whether to enable Jaeger protocol in HTTP API.
@@ -272,19 +292,7 @@ with_metric_engine = true
prom_validation_mode = "strict"
## Experimental: enable Prometheus remote write v2 native histogram ingestion.
experimental_enable_prometheus_native_histogram = false
## Interval to flush pending rows batcher.
## Set to "0s" to disable batching mode in Prometheus Remote Write endpoint
#+pending_rows_flush_interval = "0s"
## Max rows per pending batch before triggering a flush.
#+max_batch_rows = 100000
## Max number of concurrent batch flushes.
#+max_concurrent_flushes = 256
## Capacity of the pending batch worker channel.
#+worker_channel_capacity = 65526
## Max inflight write requests before backpressure.
#+max_inflight_requests = 3000
## Maximum number of logical-table flow notifications waiting in the shared queue.
#+flow_notification_queue_capacity = 1024
## The WAL options.
[wal]
+1 -1
View File
@@ -185,7 +185,7 @@ pub(super) fn default_max_concurrent_flushes() -> u64 {
256
}
pub(super) fn default_worker_channel_capacity() -> u64 {
65526
65_536
}
pub(super) fn default_max_inflight_requests() -> u64 {
3000
+9
View File
@@ -30,6 +30,7 @@ use datanode::config::{DatanodeOptions, RegionEngineConfig, StorageConfig};
use file_engine::config::EngineConfig as FileEngineConfig;
use flow::FlownodeOptions;
use frontend::frontend::FrontendOptions;
use frontend::service_config::PendingRowsBatcherOptions;
use meta_client::MetaClientOptions;
use meta_srv::metasrv::MetasrvOptions;
use meta_srv::selector::SelectorType;
@@ -206,6 +207,10 @@ fn test_load_frontend_example_config() {
);
let expected = GreptimeOptions::<FrontendOptions> {
component: FrontendOptions {
pending_rows_batcher: PendingRowsBatcherOptions {
logical_table: Some(Default::default()),
..Default::default()
},
default_timezone: Some("UTC".to_string()),
default_column_prefix: Some("greptime".to_string()),
auto_create_table: true,
@@ -389,6 +394,10 @@ fn test_load_standalone_example_config() {
);
let expected = GreptimeOptions::<StandaloneOptions> {
component: StandaloneOptions {
pending_rows_batcher: PendingRowsBatcherOptions {
logical_table: Some(Default::default()),
..Default::default()
},
default_timezone: Some("UTC".to_string()),
default_column_prefix: Some("greptime".to_string()),
auto_create_table: true,
+5
View File
@@ -51,6 +51,11 @@ remote datanodes via `operator`/`client`.
- Internal gRPC listeners mark requests with `Channel::Internal` in middleware
(`server.rs`), including requests handled by Enterprise Flight wrappers.
- **Logical-table batching** (`instance/logical_batcher.rs`): `Services` initializes
one shared batcher for opted-in HTTP Prom and nonlegacy OTLP metric-engine
writes. OTLP checks operator eligibility and falls back for incompatible tables.
The schema adapter holds a weak instance reference to avoid an ownership cycle.
- **Table batching** (`instance/builder.rs`): protocol entry points opt in through
`QueryContext`. The primary inserter prepares eligible ordinary-table writes
for `servers::batcher::table::TablePendingRowsBatcher`. A separate execution-only
+44 -5
View File
@@ -67,7 +67,7 @@ pub struct FrontendOptions {
pub postgres: PostgresOptions,
pub opentsdb: OpentsdbOptions,
pub influxdb: InfluxdbOptions,
/// Shared experimental ordinary-table batching; independent of Prom batching.
/// Ordinary-table batching with independent logical-table controls.
pub pending_rows_batcher: PendingRowsBatcherOptions,
pub prom_store: PromStoreOptions,
pub jaeger: JaegerOptions,
@@ -131,6 +131,7 @@ impl Configurable for FrontendOptions {
"meta_client.metasrv_addrs",
"event_recorder.event_types",
"pending_rows_batcher.protocols",
"pending_rows_batcher.logical_table.protocols",
])
}
}
@@ -225,6 +226,38 @@ mod tests {
type GrpcStream<T> =
Pin<Box<dyn Stream<Item = std::result::Result<T, Status>> + Send + Sync + 'static>>;
#[test]
fn test_logical_batcher_protocols_from_env() {
temp_env::with_vars(
[
(
"FRONTEND_LOGICAL_TEST__PENDING_ROWS_BATCHER__PROTOCOLS",
Some("otlp,influxdb"),
),
(
"FRONTEND_LOGICAL_TEST__PENDING_ROWS_BATCHER__LOGICAL_TABLE__PROTOCOLS",
Some("prom,otlp"),
),
],
|| {
let options =
FrontendOptions::load_layered_options(None, "FRONTEND_LOGICAL_TEST").unwrap();
assert_eq!(options.pending_rows_batcher.table.protocols.len(), 2);
assert_eq!(
options
.pending_rows_batcher
.logical_table
.unwrap()
.protocols,
vec![
servers::http::BatchingProtocol::Prom,
servers::http::BatchingProtocol::Otlp
]
);
},
);
}
#[test]
fn test_batcher_protocols_from_env() {
temp_env::with_vars(
@@ -236,7 +269,7 @@ mod tests {
let options =
FrontendOptions::load_layered_options(None, "FRONTEND_BATCHER_TEST").unwrap();
assert_eq!(
options.pending_rows_batcher.protocols,
options.pending_rows_batcher.table.protocols,
vec![
servers::http::BatchingProtocol::Influxdb,
servers::http::BatchingProtocol::HttpSql
@@ -252,6 +285,7 @@ mod tests {
assert!(
!defaults
.pending_rows_batcher
.table
.pending_rows_batching_enabled()
);
let options: FrontendOptions = toml::from_str(
@@ -263,9 +297,14 @@ max_batch_rows = 25
"#,
)
.unwrap();
assert_eq!(options.pending_rows_batcher.max_batch_rows, 25);
assert_eq!(options.pending_rows_batcher.protocols.len(), 2);
assert!(options.pending_rows_batcher.pending_rows_batching_enabled());
assert_eq!(options.pending_rows_batcher.table.max_batch_rows, 25);
assert_eq!(options.pending_rows_batcher.table.protocols.len(), 2);
assert!(
options
.pending_rows_batcher
.table
.pending_rows_batching_enabled()
);
let serialized = toml::to_string(&options).unwrap();
let parsed: FrontendOptions = toml::from_str(&serialized).unwrap();
assert_eq!(options.influxdb, parsed.influxdb);
+4 -1
View File
@@ -21,6 +21,7 @@ mod import_packed;
mod influxdb;
mod jaeger;
mod log_handler;
mod logical_batcher;
mod logs;
mod opentsdb;
mod otlp;
@@ -31,7 +32,7 @@ mod region_query;
use std::collections::HashSet;
use std::pin::Pin;
use std::sync::atomic::AtomicBool;
use std::sync::{Arc, atomic};
use std::sync::{Arc, OnceLock, atomic};
use std::time::{Duration, SystemTime};
use async_stream::stream;
@@ -78,6 +79,7 @@ use query::metrics::OnDone;
use query::parser::{PromQuery, QueryStatement};
use query::query_engine::DescribeResult;
use query::query_engine::options::{QueryOptions, validate_catalog_and_schema};
use servers::batcher::logical_table::LogicalTablePendingRowsBatcher;
use servers::error::{
self as server_error, AuthSnafu, CommonMetaSnafu, ExecuteQuerySnafu,
OtlpMetricModeIncompatibleSnafu, UnexpectedResultSnafu,
@@ -130,6 +132,7 @@ pub struct Instance {
query_engine: QueryEngineRef,
plugins: Plugins,
inserter: InserterRef,
logical_batcher: Arc<OnceLock<Option<Arc<LogicalTablePendingRowsBatcher>>>>,
deleter: DeleterRef,
table_metadata_manager: TableMetadataManagerRef,
event_recorder: EventRecorderRef,
+24 -24
View File
@@ -61,7 +61,7 @@ use crate::heartbeat::frontend_peer_addr;
use crate::instance::Instance;
use crate::instance::entity_graph::EntityGraphProviderImpl;
use crate::instance::region_query::FrontendRegionQueryHandler;
use crate::service_config::PendingRowsBatcherOptions;
use crate::service_config::BatcherOptions;
/// The frontend [`Instance`] builder.
pub struct FrontendBuilder {
@@ -237,31 +237,30 @@ impl FrontendBuilder {
};
// The execution-only inserter owns no batchers, avoiding an Arc cycle.
let bulk_inserter = Arc::new(create_inserter());
let build_batcher =
|options: &PendingRowsBatcherOptions| -> Option<Arc<dyn PendingRowsBatcher>> {
if !options.pending_rows_batching_enabled()
|| (self.options.prom_store.with_metric_engine
&& options
.protocols
.iter()
.all(|protocol| *protocol == BatchingProtocol::Prom))
{
return None;
}
TablePendingRowsBatcher::try_new(
options.pending_rows_flush_interval,
options.max_batch_rows,
options.max_concurrent_flushes,
options.worker_channel_capacity,
options.max_inflight_requests,
options.flow_notification_queue_capacity,
bulk_inserter.clone(),
)
.map(|batcher| batcher as Arc<dyn PendingRowsBatcher>)
};
let build_batcher = |options: &BatcherOptions| -> Option<Arc<dyn PendingRowsBatcher>> {
if !options.pending_rows_batching_enabled()
|| (self.options.prom_store.with_metric_engine
&& options
.protocols
.iter()
.all(|protocol| *protocol == BatchingProtocol::Prom))
{
return None;
}
TablePendingRowsBatcher::try_new(
options.pending_rows_flush_interval,
options.max_batch_rows,
options.max_concurrent_flushes,
options.worker_channel_capacity,
options.max_inflight_requests,
options.flow_notification_queue_capacity,
bulk_inserter.clone(),
)
.map(|batcher| batcher as Arc<dyn PendingRowsBatcher>)
};
let inserter = Arc::new(
create_inserter()
.with_pending_rows_batcher(build_batcher(&self.options.pending_rows_batcher)),
.with_pending_rows_batcher(build_batcher(self.options.table_batcher_options())),
);
let deleter = Arc::new(Deleter::new(
self.catalog_manager.clone(),
@@ -386,6 +385,7 @@ impl FrontendBuilder {
admin_event_recorder.install(&event_recorder);
Ok(Instance {
logical_batcher: Default::default(),
frontend_peer_addr,
experimental_metric_export: self.options.experimental_metric_export,
catalog_manager: self.catalog_manager,
@@ -0,0 +1,102 @@
// Copyright 2023 Greptime Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use std::sync::{Arc, Weak};
use api::v1::ColumnSchema;
use async_trait::async_trait;
use servers::batcher::logical_table::{LogicalTablePendingRowsBatcher, PendingRowsSchemaAlterer};
use servers::error::{BatcherChannelClosedSnafu, Result};
use servers::http::BatchingProtocol;
use session::context::QueryContextRef;
use snafu::OptionExt;
use crate::frontend::FrontendOptions;
use crate::instance::Instance;
impl Instance {
pub(crate) fn init_logical_batcher(self: &Arc<Self>, options: &FrontendOptions) {
self.logical_batcher.get_or_init(|| {
let options_batcher = options.logical_batcher_options();
let enabled = options_batcher
.protocols
.iter()
.any(|protocol| match protocol {
BatchingProtocol::Prom => options.prom_store.enable,
BatchingProtocol::Otlp => options.otlp.enable,
_ => false,
});
if !options.prom_store.with_metric_engine
|| !enabled
|| !options_batcher.pending_rows_batching_enabled()
{
return None;
}
LogicalTablePendingRowsBatcher::try_new(
self.partition_manager().clone(),
self.node_manager().clone(),
self.catalog_manager().clone(),
self.table_flownode_set_cache().clone(),
true,
Arc::new(LogicalTables(Arc::downgrade(self))),
options_batcher.pending_rows_flush_interval,
options_batcher.max_batch_rows,
options_batcher.max_concurrent_flushes,
options_batcher.worker_channel_capacity,
options_batcher.max_inflight_requests,
options_batcher.flow_notification_queue_capacity,
)
});
}
pub(crate) fn logical_batcher(&self) -> Option<&Arc<LogicalTablePendingRowsBatcher>> {
self.logical_batcher.get().and_then(Option::as_ref)
}
}
// The instance owns the shared batcher; schema preparation must not keep that
// owner alive through a reference cycle.
struct LogicalTables(Weak<Instance>);
#[async_trait]
impl PendingRowsSchemaAlterer for LogicalTables {
async fn create_tables_if_missing_batch(
&self,
catalog: &str,
schema: &str,
tables: &[(&str, &[ColumnSchema])],
with_metric_engine: bool,
ctx: QueryContextRef,
) -> Result<()> {
self.0
.upgrade()
.context(BatcherChannelClosedSnafu)?
.create_tables_if_missing_batch(catalog, schema, tables, with_metric_engine, ctx)
.await
}
async fn add_missing_prom_tag_columns_batch(
&self,
catalog: &str,
schema: &str,
tables: &[(&str, &[String])],
ctx: QueryContextRef,
) -> Result<()> {
self.0
.upgrade()
.context(BatcherChannelClosedSnafu)?
.add_missing_prom_tag_columns_batch(catalog, schema, tables, ctx)
.await
}
}
+45 -12
View File
@@ -27,10 +27,11 @@ use client::Output;
use common_catalog::consts::{trace_operations_table_name, trace_services_table_name};
use common_error::ext::BoxedError;
use common_query::prelude::GREPTIME_PHYSICAL_TABLE;
use common_query::{OutputData, OutputMeta};
use common_telemetry::{tracing, warn};
use opentelemetry_proto::tonic::collector::logs::v1::ExportLogsServiceRequest;
use opentelemetry_proto::tonic::collector::trace::v1::ExportTraceServiceRequest;
use operator::insert::{admit_row_insert_batches, admit_write};
use operator::insert::{Inserter, admit_row_insert_batches, admit_write};
use otel_arrow_rust::proto::opentelemetry::collector::metrics::v1::ExportMetricsServiceRequest;
use pipeline::{GreptimePipelineParams, PipelineWay};
use servers::error::{self, AuthSnafu, Result as ServerResult};
@@ -165,27 +166,59 @@ impl OpenTelemetryProtocolHandler for Instance {
Arc::new(c)
};
let physical_table = ctx
.extension(PHYSICAL_TABLE_PARAM)
.unwrap_or(GREPTIME_PHYSICAL_TABLE)
.to_string();
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
.inserter
.can_batch_metric_rows(&requests, &ctx, &physical_table)
.await
.map_err(BoxedError::new)
.context(error::ExecuteGrpcQuerySnafu)?
{
batcher
} else {
None
};
// OTLP tables have one sample field in both the legacy and physical paths.
let output = if metric_ctx.is_legacy || !metric_ctx.with_metric_engine {
let output = if let Some(batcher) = batcher {
let (rows, cost) = batcher
.submit_with(requests, ctx.clone(), |mut requests| {
let ctx = ctx.clone();
async move {
Inserter::meter_row_inserts(&mut requests, &ctx)
.await
.map_err(BoxedError::new)
.context(error::ExecuteGrpcQuerySnafu)
}
})
.await?;
Output::new(
OutputData::AffectedRows(rows as usize),
OutputMeta::new_with_cost(cost as _),
)
} else if metric_ctx.is_legacy || !metric_ctx.with_metric_engine {
self.handle_row_inserts(requests, ctx.clone(), false, true)
.await
.map_err(BoxedError::new)
.context(error::ExecuteGrpcQuerySnafu)
.context(error::ExecuteGrpcQuerySnafu)?
} else {
let physical_table = ctx
.extension(PHYSICAL_TABLE_PARAM)
.unwrap_or(GREPTIME_PHYSICAL_TABLE)
.to_string();
self.handle_metric_row_inserts(requests, ctx.clone(), physical_table)
.await
.map_err(BoxedError::new)
.context(error::ExecuteGrpcQuerySnafu)
}?;
.context(error::ExecuteGrpcQuerySnafu)?
};
outcome.write_cost = output.meta.cost;
// Derived enrichment, written after the metric data is committed:
// failing here would make the client retry data the server already
// accepted, so every failure degrades to a warning instead.
// Derived enrichment follows the accepted metric submission, which may
// still be queued in asynchronous mode. Failures remain warning-only
// to avoid retrying metric data the server already accepted.
if let Some(resource_info) = resource_info {
let written = match self.check_row_insert_permission(
&resource_info,
+126 -45
View File
@@ -24,9 +24,7 @@ use common_base::Plugins;
use common_config::Configurable;
use common_telemetry::{info, warn};
use meta_client::MetaClientOptions;
use servers::batcher::logical_table::{
LogicalTablePendingRowsBatcher, pending_rows_batch_sync_enabled,
};
use servers::batcher::pending_rows_batch_sync_enabled;
use servers::error::Error as ServerError;
use servers::grpc::builder::GrpcServerBuilder;
use servers::grpc::flight::FlightCraftRef;
@@ -74,6 +72,7 @@ where
{
pub fn new(opts: T, instance: Arc<Instance>, plugins: Plugins) -> Self {
let feopts = opts.clone().into();
instance.init_logical_batcher(&feopts);
// Create server request memory limiter for all server protocols
let server_memory_limiter = ServerMemoryLimiter::new(
feopts.max_in_flight_write_bytes.as_bytes(),
@@ -110,7 +109,8 @@ where
request_memory_limiter: ServerMemoryLimiter,
) -> HttpServerBuilder {
let mut builder = HttpServerBuilder::new(effective_http_options(opts))
.with_batching_protocols(opts.pending_rows_batcher.protocols.clone())
.with_batching_protocols(opts.table_batcher_options().protocols.clone())
.with_logical_batching_protocols(opts.logical_batcher_options().protocols)
.with_memory_limiter(request_memory_limiter)
.with_sql_handler(self.instance.clone());
@@ -134,24 +134,12 @@ where
let prom_store = effective_prom_store_options(opts);
if prom_store.enable {
let pending_rows_batcher = if prom_store.with_metric_engine {
LogicalTablePendingRowsBatcher::try_new(
self.instance.partition_manager().clone(),
self.instance.node_manager().clone(),
self.instance.catalog_manager().clone(),
self.instance.table_flownode_set_cache().clone(),
prom_store.with_metric_engine,
self.instance.clone(),
prom_store.pending_rows_flush_interval,
prom_store.max_batch_rows,
prom_store.max_concurrent_flushes,
prom_store.worker_channel_capacity,
prom_store.max_inflight_requests,
prom_store.flow_notification_queue_capacity,
)
} else {
None
};
let pending_rows_batcher = opts
.logical_batcher_options()
.protocols
.contains(&BatchingProtocol::Prom)
.then(|| self.instance.logical_batcher().cloned())
.flatten();
builder = builder
.with_prom_handler(
self.instance.clone(),
@@ -439,15 +427,16 @@ where
/// Selected shared controls override legacy Prom batching knobs, not protocol behavior.
fn effective_prom_store_options(opts: &FrontendOptions) -> PromStoreOptions {
let mut prom_store = opts.prom_store.clone();
let shared = &opts.pending_rows_batcher;
if shared.protocols.contains(&BatchingProtocol::Prom) && shared.pending_rows_batching_enabled()
{
let shared = opts.logical_batcher_options();
if shared.protocols.contains(&BatchingProtocol::Prom) {
prom_store.pending_rows_flush_interval = shared.pending_rows_flush_interval;
prom_store.max_batch_rows = shared.max_batch_rows;
prom_store.max_concurrent_flushes = shared.max_concurrent_flushes;
prom_store.worker_channel_capacity = shared.worker_channel_capacity;
prom_store.max_inflight_requests = shared.max_inflight_requests;
prom_store.flow_notification_queue_capacity = shared.flow_notification_queue_capacity;
} else {
prom_store.pending_rows_flush_interval = Duration::ZERO;
}
prom_store
}
@@ -459,10 +448,9 @@ fn effective_http_options(opts: &FrontendOptions) -> HttpOptions {
fn effective_http_options_with_sync(opts: &FrontendOptions, batch_sync: bool) -> HttpOptions {
let mut http = opts.http.clone();
let prom_store = effective_prom_store_options(opts);
let shared = &opts.pending_rows_batcher;
// Ordinary-table batching always waits for its flush, independently of the
// dedicated Prom batcher's asynchronous acknowledgement mode.
let common_enabled = shared.pending_rows_batching_enabled()
let shared = opts.table_batcher_options();
let common_enabled = batch_sync
&& shared.pending_rows_batching_enabled()
&& shared.protocols.iter().any(|protocol| {
*protocol != BatchingProtocol::Prom
|| (prom_store.enable && !prom_store.with_metric_engine)
@@ -470,7 +458,19 @@ fn effective_http_options_with_sync(opts: &FrontendOptions, batch_sync: bool) ->
let common_interval = common_enabled.then_some(shared.pending_rows_flush_interval);
let prom_interval = (prom_store.pending_rows_batching_enabled() && batch_sync)
.then_some(prom_store.pending_rows_flush_interval);
let Some(flush_interval) = common_interval.into_iter().chain(prom_interval).max() else {
let logical = opts.logical_batcher_options();
let otlp_interval = (batch_sync
&& opts.otlp.enable
&& opts.prom_store.with_metric_engine
&& logical.protocols.contains(&BatchingProtocol::Otlp)
&& logical.pending_rows_batching_enabled())
.then_some(logical.pending_rows_flush_interval);
let Some(flush_interval) = common_interval
.into_iter()
.chain(prom_interval)
.chain(otlp_interval)
.max()
else {
return http;
};
let fallback_timeout = flush_interval.saturating_add(Duration::from_secs(1));
@@ -515,6 +515,81 @@ mod tests {
use crate::instance::builder::FrontendBuilder;
use crate::server::*;
#[tokio::test]
async fn test_logical_batcher_shared_without_prom_endpoint() {
use crate::service_config::BatcherOptions;
let mut options = FrontendOptions::default();
options.prom_store.enable = false;
options.pending_rows_batcher.logical_table = Some(BatcherOptions {
protocols: vec![BatchingProtocol::Otlp],
pending_rows_flush_interval: Duration::from_millis(5),
..Default::default()
});
let meta_client = Arc::new(
MetaClientBuilder::new(0, Role::Frontend)
.enable_procedure()
.build(),
);
let instance = Arc::new(
FrontendBuilder::new_test(&options, meta_client)
.try_build()
.await
.unwrap(),
);
let services = Services::new(options.clone(), instance.clone(), Plugins::default());
let batcher = instance.logical_batcher().unwrap().clone();
instance.init_logical_batcher(&options);
assert!(Arc::ptr_eq(&batcher, instance.logical_batcher().unwrap()));
let weak = Arc::downgrade(&instance);
drop(services);
drop(instance);
assert!(
weak.upgrade().is_none(),
"batcher must not retain its schema owner"
);
}
#[test]
fn test_logical_batcher_http_timeout_and_prom_disable() {
use crate::service_config::pending_rows_batcher::BatcherOptions;
let mut opts = FrontendOptions::default();
opts.http.timeout = Duration::from_secs(1);
opts.prom_store.pending_rows_flush_interval = Duration::from_secs(2);
opts.pending_rows_batcher.logical_table = Some(BatcherOptions {
protocols: vec![BatchingProtocol::Otlp],
pending_rows_flush_interval: Duration::from_secs(5),
..Default::default()
});
assert!(!effective_prom_store_options(&opts).pending_rows_batching_enabled());
// Both logical protocols follow the global acknowledgement policy,
// independently of enabling the Prom HTTP endpoint.
opts.prom_store.enable = false;
assert_eq!(
effective_http_options_with_sync(&opts, true).timeout,
Duration::from_secs(6)
);
assert_eq!(
effective_http_options_with_sync(&opts, false).timeout,
opts.http.timeout
);
opts.prom_store.with_metric_engine = false;
assert_eq!(
effective_http_options_with_sync(&opts, false).timeout,
opts.http.timeout
);
opts.prom_store.with_metric_engine = true;
opts.pending_rows_batcher
.logical_table
.as_mut()
.unwrap()
.protocols
.clear();
assert_eq!(
effective_http_options_with_sync(&opts, true).timeout,
opts.http.timeout
);
}
#[test]
fn test_effective_prom_batching_controls() {
// Only an enabled shared Prom selection replaces the legacy controls.
@@ -532,7 +607,7 @@ mod tests {
opts.prom_store.enable = prom_enabled;
opts.prom_store
.experimental_enable_prometheus_native_histogram = true;
let shared = &mut opts.pending_rows_batcher;
let shared = &mut opts.pending_rows_batcher.table;
shared.protocols = vec![if selected {
BatchingProtocol::Prom
} else {
@@ -569,7 +644,7 @@ mod tests {
#[test]
fn test_http_timeout_covers_synchronous_batchers() {
// Shared ordinary writes remain synchronous even when Prom is asynchronous.
// Only synchronous batchers extend the HTTP timeout.
for (
protocols,
metric_engine,
@@ -579,10 +654,10 @@ mod tests {
timeout_secs,
expected_secs,
) in [
(vec![BatchingProtocol::Prom], false, false, 5, 2, 1, 6),
(vec![BatchingProtocol::Prom], false, false, 5, 2, 1, 1),
(vec![BatchingProtocol::Prom], true, false, 5, 2, 1, 1),
(vec![BatchingProtocol::Prom], true, true, 5, 2, 1, 6),
(vec![BatchingProtocol::Influxdb], true, false, 5, 2, 1, 6),
(vec![BatchingProtocol::Influxdb], true, false, 5, 2, 1, 1),
(vec![BatchingProtocol::Influxdb], true, true, 5, 8, 1, 9),
(vec![BatchingProtocol::Influxdb], true, true, 8, 5, 1, 9),
(vec![BatchingProtocol::Influxdb], true, false, 5, 2, 0, 0),
@@ -595,8 +670,8 @@ mod tests {
opts.http.timeout = Duration::from_secs(timeout_secs);
opts.prom_store.with_metric_engine = metric_engine;
opts.prom_store.pending_rows_flush_interval = Duration::from_secs(legacy_secs);
opts.pending_rows_batcher.protocols = protocols;
opts.pending_rows_batcher.pending_rows_flush_interval =
opts.pending_rows_batcher.table.protocols = protocols;
opts.pending_rows_batcher.table.pending_rows_flush_interval =
Duration::from_secs(shared_secs);
assert_eq!(
effective_http_options_with_sync(&opts, batch_sync).timeout,
@@ -697,24 +772,30 @@ mod tests {
fn test_invalid_shared_batching_preserves_prom_store_options() {
type KnobMutator = fn(&mut FrontendOptions);
let cases: [KnobMutator; 5] = [
|opts| opts.pending_rows_batcher.max_concurrent_flushes = usize::MAX,
|opts| opts.pending_rows_batcher.worker_channel_capacity = usize::MAX,
|opts| opts.pending_rows_batcher.max_inflight_requests = usize::MAX,
|opts| opts.pending_rows_batcher.table.max_concurrent_flushes = usize::MAX,
|opts| opts.pending_rows_batcher.table.worker_channel_capacity = usize::MAX,
|opts| opts.pending_rows_batcher.table.max_inflight_requests = usize::MAX,
|opts| {
opts.pending_rows_batcher.flow_notification_queue_capacity =
NonZeroUsize::new(usize::MAX).unwrap()
opts.pending_rows_batcher
.table
.flow_notification_queue_capacity = NonZeroUsize::new(usize::MAX).unwrap()
},
|opts| opts.pending_rows_batcher.pending_rows_flush_interval = Duration::MAX,
|opts| opts.pending_rows_batcher.table.pending_rows_flush_interval = Duration::MAX,
];
for invalidate in cases {
let mut opts = FrontendOptions::default();
opts.http.timeout = Duration::from_secs(1);
opts.prom_store.pending_rows_flush_interval = Duration::from_secs(5);
opts.pending_rows_batcher.protocols =
opts.pending_rows_batcher.table.protocols =
vec![BatchingProtocol::Prom, BatchingProtocol::Influxdb];
opts.pending_rows_batcher.pending_rows_flush_interval = Duration::from_secs(10);
opts.pending_rows_batcher.table.pending_rows_flush_interval = Duration::from_secs(10);
invalidate(&mut opts);
assert!(!opts.pending_rows_batcher.pending_rows_batching_enabled());
assert!(
!opts
.pending_rows_batcher
.table
.pending_rows_batching_enabled()
);
assert_eq!(opts.prom_store, effective_prom_store_options(&opts));
assert_eq!(
Duration::from_secs(6),
+1 -1
View File
@@ -26,6 +26,6 @@ pub use jaeger::JaegerOptions;
pub use mysql::MysqlOptions;
pub use opentsdb::OpentsdbOptions;
pub use otlp::OtlpOptions;
pub use pending_rows_batcher::PendingRowsBatcherOptions;
pub use pending_rows_batcher::{BatcherOptions, PendingRowsBatcherOptions};
pub use postgres::PostgresOptions;
pub use prom_store::PromStoreOptions;
@@ -20,10 +20,28 @@ use serde::{Deserialize, Serialize};
use servers::http::BatchingProtocol;
use tokio::sync::Semaphore;
/// Experimental table write batching options shared by HTTP ingestion protocols.
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
use crate::frontend::FrontendOptions;
/// Independent ordinary-table and logical-table batching configuration.
#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq)]
#[serde(default)]
pub struct PendingRowsBatcherOptions {
/// Existing ordinary-table controls retain their original TOML paths.
#[serde(flatten)]
pub table: BatcherOptions,
/// Independent logical-table controls; absence retains the legacy Prom fallback.
#[serde(
default,
skip_serializing_if = "Option::is_none",
deserialize_with = "deserialize_logical_options"
)]
pub logical_table: Option<BatcherOptions>,
}
/// Write batching controls shared by HTTP ingestion protocols.
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
#[serde(default)]
pub struct BatcherOptions {
/// HTTP write protocols sharing this batcher; empty disables all entrances.
pub protocols: Vec<BatchingProtocol>,
/// Time from the first pending submission to a timed flush. Zero disables batching.
@@ -41,7 +59,7 @@ pub struct PendingRowsBatcherOptions {
pub flow_notification_queue_capacity: NonZeroUsize,
}
impl PendingRowsBatcherOptions {
impl BatcherOptions {
/// Returns whether a protocol opts in and its controls pass construction validation.
pub fn pending_rows_batching_enabled(&self) -> bool {
!self.protocols.is_empty()
@@ -58,79 +76,229 @@ impl PendingRowsBatcherOptions {
}
}
impl Default for PendingRowsBatcherOptions {
impl Default for BatcherOptions {
fn default() -> Self {
Self {
protocols: Vec::new(),
pending_rows_flush_interval: Duration::ZERO,
max_batch_rows: 100_000,
max_concurrent_flushes: 256,
worker_channel_capacity: 65526,
worker_channel_capacity: 65_536,
max_inflight_requests: 3000,
flow_notification_queue_capacity: NonZeroUsize::new(1024).unwrap_or(NonZeroUsize::MIN),
}
}
}
/// Rejects protocols that cannot write metric-engine logical tables.
fn deserialize_logical_options<'de, D>(deserializer: D) -> Result<Option<BatcherOptions>, D::Error>
where
D: serde::Deserializer<'de>,
{
let options = Option::<BatcherOptions>::deserialize(deserializer)?;
if options.as_ref().is_some_and(|options| {
options
.protocols
.iter()
.any(|protocol| !matches!(protocol, BatchingProtocol::Prom | BatchingProtocol::Otlp))
}) {
return Err(serde::de::Error::custom(
"pending_rows_batcher.logical_table only supports prom and otlp",
));
}
Ok(options)
}
impl FrontendOptions {
/// Explicit configuration wins even when it disables batching.
pub(crate) fn table_batcher_options(&self) -> &BatcherOptions {
&self.pending_rows_batcher.table
}
/// Resolve legacy Prom controls once without merging fields across config blocks.
pub(crate) fn logical_batcher_options(&self) -> BatcherOptions {
if let Some(options) = &self.pending_rows_batcher.logical_table {
return options.clone();
}
let shared = &self.pending_rows_batcher.table;
if shared.protocols.contains(&BatchingProtocol::Prom)
&& shared.pending_rows_batching_enabled()
{
let mut options = shared.clone();
options.protocols = vec![BatchingProtocol::Prom];
return options;
}
let prom = &self.prom_store;
BatcherOptions {
protocols: vec![BatchingProtocol::Prom],
pending_rows_flush_interval: prom.pending_rows_flush_interval,
max_batch_rows: prom.max_batch_rows,
max_concurrent_flushes: prom.max_concurrent_flushes,
worker_channel_capacity: prom.worker_channel_capacity,
max_inflight_requests: prom.max_inflight_requests,
flow_notification_queue_capacity: prom.flow_notification_queue_capacity,
}
}
}
#[cfg(test)]
mod tests {
use crate::service_config::pending_rows_batcher::*;
#[test]
fn test_batcher_config_precedence() {
use crate::frontend::FrontendOptions;
let legacy = "[prom_store]\nenable = true\nwith_metric_engine = true\npending_rows_flush_interval = '2s'\nmax_batch_rows = 9\n";
let shared = "[pending_rows_batcher]\nprotocols = ['prom', 'otlp']\npending_rows_flush_interval = '3s'\n";
let options: FrontendOptions = toml::from_str(legacy).unwrap();
assert_eq!(options.logical_batcher_options().max_batch_rows, 9);
assert_eq!(
options.logical_batcher_options().protocols,
vec![BatchingProtocol::Prom]
);
let options: FrontendOptions = toml::from_str(&format!("{legacy}{shared}")).unwrap();
assert_eq!(
options
.logical_batcher_options()
.pending_rows_flush_interval,
Duration::from_secs(3)
);
assert_eq!(options.logical_batcher_options().max_batch_rows, 100_000);
assert_eq!(
options.logical_batcher_options().protocols,
vec![BatchingProtocol::Prom]
);
for explicit in [
"",
"protocols = ['otlp']\npending_rows_flush_interval = '5s'",
] {
let options: FrontendOptions = toml::from_str(&format!(
"{legacy}{shared}[pending_rows_batcher.logical_table]\n{explicit}"
))
.unwrap();
assert_eq!(
&options.logical_batcher_options(),
options.pending_rows_batcher.logical_table.as_ref().unwrap()
);
assert_eq!(
options.table_batcher_options(),
&toml::from_str::<FrontendOptions>(&format!("{legacy}{shared}"))
.unwrap()
.pending_rows_batcher
.table
);
let restored: FrontendOptions =
toml::from_str(&toml::to_string(&options).unwrap()).unwrap();
assert_eq!(
options.pending_rows_batcher.logical_table,
restored.pending_rows_batcher.logical_table
);
assert_eq!(options.pending_rows_batcher, restored.pending_rows_batcher);
}
}
#[test]
fn test_batcher_config_independent_defaults() {
use crate::frontend::FrontendOptions;
let parent = "[pending_rows_batcher]\nprotocols = ['prom']\npending_rows_flush_interval = '3s'\nmax_batch_rows = 7\n";
let options: FrontendOptions = toml::from_str(parent).unwrap();
let serialized = toml::to_string(&options).unwrap();
assert!(!serialized.contains("[pending_rows_batcher.logical_table]"));
assert!(!serialized.contains("[pending_rows_batcher.table]"));
let restored: FrontendOptions = toml::from_str(&serialized).unwrap();
assert_eq!(restored.logical_batcher_options().max_batch_rows, 7);
let options: FrontendOptions =
toml::from_str(&format!("{parent}[pending_rows_batcher.logical_table]")).unwrap();
assert_eq!(options.table_batcher_options().max_batch_rows, 7);
assert_eq!(options.logical_batcher_options(), BatcherOptions::default());
assert!(
!options
.logical_batcher_options()
.pending_rows_batching_enabled()
);
}
#[test]
fn test_logical_batcher_config_protocols() {
use crate::frontend::FrontendOptions;
for protocols in ["[]", "['prom']", "['otlp']", "['prom', 'otlp']"] {
let options: FrontendOptions = toml::from_str(&format!(
"[pending_rows_batcher.logical_table]\nprotocols = {protocols}"
))
.unwrap();
assert!(options.pending_rows_batcher.logical_table.is_some());
}
for protocol in [
"influxdb",
"logs",
"loki",
"http_sql",
"opentsdb",
"elasticsearch",
"splunk",
] {
assert!(
toml::from_str::<FrontendOptions>(&format!(
"[pending_rows_batcher.logical_table]\nprotocols = ['{protocol}']"
))
.is_err()
);
}
}
#[test]
fn test_protocols() {
let options: PendingRowsBatcherOptions = toml::from_str(
let options: BatcherOptions = toml::from_str(
"protocols = ['influxdb', 'opentsdb', 'otlp', 'logs', 'loki', 'splunk', 'elasticsearch', 'http_sql', 'prom']",
).unwrap();
assert_eq!(options.protocols.len(), 9);
assert!(options.protocols.contains(&BatchingProtocol::HttpSql));
assert!(PendingRowsBatcherOptions::default().protocols.is_empty());
assert!(BatcherOptions::default().protocols.is_empty());
for invalid in ["sql", "jaeger", "unknown"] {
assert!(
toml::from_str::<PendingRowsBatcherOptions>(&format!("protocols = ['{invalid}']"))
.is_err()
toml::from_str::<BatcherOptions>(&format!("protocols = ['{invalid}']")).is_err()
);
}
}
#[test]
fn test_notification_capacity() {
let default = PendingRowsBatcherOptions::default();
let default = BatcherOptions::default();
assert_eq!(default.flow_notification_queue_capacity.get(), 1024);
let configured: PendingRowsBatcherOptions =
let configured: BatcherOptions =
toml::from_str("flow_notification_queue_capacity = 8").unwrap();
assert_eq!(configured.flow_notification_queue_capacity.get(), 8);
assert!(
toml::from_str::<PendingRowsBatcherOptions>("flow_notification_queue_capacity = 0")
.is_err()
);
assert!(toml::from_str::<BatcherOptions>("flow_notification_queue_capacity = 0").is_err());
}
#[test]
fn test_defaults_and_roundtrip() {
let options: PendingRowsBatcherOptions = toml::from_str("").unwrap();
assert_eq!(options, PendingRowsBatcherOptions::default());
let options: BatcherOptions = toml::from_str("").unwrap();
assert_eq!(options, BatcherOptions::default());
assert!(!options.pending_rows_batching_enabled());
assert_eq!(options.max_batch_rows, 100_000);
assert_eq!(options.max_concurrent_flushes, 256);
assert_eq!(options.worker_channel_capacity, 65526);
assert_eq!(options.worker_channel_capacity, 65_536);
assert_eq!(options.max_inflight_requests, 3000);
let serialized = toml::to_string(&options).unwrap();
assert_eq!(
options,
toml::from_str::<PendingRowsBatcherOptions>(&serialized).unwrap()
toml::from_str::<BatcherOptions>(&serialized).unwrap()
);
}
#[test]
fn test_partial_options_and_zero_controls() {
let options: PendingRowsBatcherOptions =
let options: BatcherOptions =
toml::from_str("pending_rows_flush_interval = '5ms'").unwrap();
assert_eq!(
options.pending_rows_flush_interval,
Duration::from_millis(5)
);
assert!(!options.pending_rows_batching_enabled());
let enabled: PendingRowsBatcherOptions =
let enabled: BatcherOptions =
toml::from_str("protocols = ['http_sql']\npending_rows_flush_interval = '5ms'")
.unwrap();
assert!(enabled.pending_rows_batching_enabled());
@@ -140,7 +308,7 @@ mod tests {
"worker_channel_capacity",
"max_inflight_requests",
] {
let options: PendingRowsBatcherOptions = toml::from_str(&format!(
let options: BatcherOptions = toml::from_str(&format!(
"protocols = ['influxdb']\npending_rows_flush_interval = '5ms'\n{field} = 0"
))
.unwrap();
@@ -54,7 +54,7 @@ fn default_max_concurrent_flushes() -> usize {
}
fn default_worker_channel_capacity() -> usize {
65526
65_536
}
fn default_max_inflight_requests() -> usize {
+2 -1
View File
@@ -29,7 +29,8 @@ pub trait PendingRowsBatcher: Send + Sync {
/// Acquires one slot per original request, shared by all of its table submissions.
async fn acquire(&self) -> Result<Arc<OwnedSemaphorePermit>>;
/// Waits for the submitted rows to be written, retaining the slot through completion.
/// Submits rows according to the acknowledgement policy, retaining the slot until
/// writing completes even when the response acknowledges queue admission only.
/// Cancelling the response wait does not retract an already enqueued submission.
async fn submit(
&self,
+310 -5
View File
@@ -156,6 +156,127 @@ pub struct InstantAndNormalInsertRequests {
}
impl Inserter {
/// Checks the assumptions of the logical bulk path without changing tables.
/// Unsupported requests retain ordinary insertion, including schema policy,
/// defaults, instant TTL and row-based Flow delivery.
pub async fn can_batch_metric_rows(
&self,
requests: &RowInsertRequests,
ctx: &QueryContextRef,
physical_table: &str,
) -> Result<bool> {
if self.auto_create_disabled_reason(ctx)?.is_some() || ctx.extension(TTL_KEY).is_some() {
return Ok(false);
}
for request in &requests.inserts {
// The logical bulk encoder only supports scalar metric schemas.
// Check new tables too, before catalog lookup or schema changes.
if request.rows.as_ref().is_some_and(|rows| {
rows.schema.iter().any(|column| {
column.datatype_extension.is_some()
|| !matches!(
ColumnDataType::try_from(column.datatype),
Ok(ColumnDataType::TimestampMillisecond
| ColumnDataType::Float64
| ColumnDataType::String)
)
})
}) {
return Ok(false);
}
let Some(table) = self
.get_table(
ctx.current_catalog(),
&ctx.current_schema(),
&request.table_name,
)
.await?
else {
continue;
};
let info = table.table_info();
if info.meta.engine != METRIC_ENGINE_NAME
|| info.is_ttl_instant_table()
|| info
.meta
.options
.extra_options
.get(LOGICAL_TABLE_METADATA_KEY)
.map(String::as_str)
!= Some(physical_table)
|| info
.meta
.schema
.column_schemas()
.iter()
.any(|column| column.default_constraint().is_some())
{
return Ok(false);
}
// Physical metric tags are nullable even when their logical schema
// is not. Arrow alignment is stricter than ordinary metric insertion.
if info
.meta
.primary_key_indices
.iter()
.any(|&index| !info.meta.schema.column_schemas()[index].is_nullable())
{
return Ok(false);
}
// The current Flow cache does not distinguish streaming and batch
// flows. Keep all Flow sources on the row-based delivery path.
match self.table_flownode_set_cache.get(info.table_id()).await {
Ok(None) => {}
Ok(Some(flows)) if flows.is_empty() => {}
_ => return Ok(false),
}
}
Ok(true)
}
/// Meters an original logical-table request before bulk routing, without
/// cloning its rows or changing the request boundary used for accounting.
pub async fn meter_row_inserts(
requests: &mut RowInsertRequests,
ctx: &QueryContextRef,
) -> Result<u64> {
let metered = InstantAndNormalInsertRequests {
normal_requests: RegionInsertRequests {
requests: requests
.inserts
.iter_mut()
.map(|request| RegionInsertRequest {
rows: request.rows.take(),
..Default::default()
})
.collect(),
},
instant_requests: RegionInsertRequests::default(),
};
let cost = write_meter!(
ctx.current_catalog(),
ctx.current_schema(),
metered,
ctx.write_rows_to_admit(
ctx.current_catalog(),
&ctx.current_schema(),
count_insert_rows(&metered)?
),
ctx.channel() as u8
)
.await
.context(WriteRejectedSnafu);
for (request, region) in requests
.inserts
.iter_mut()
.zip(metered.normal_requests.requests)
{
request.rows = region.rows;
}
cost
}
pub fn new(
catalog_manager: CatalogManagerRef,
partition_manager: PartitionRuleManagerRef,
@@ -2313,19 +2434,67 @@ mod tests {
);
assert_eq!(ctx.channel(), Channel::Prometheus);
}
let attempts = meter.attempts.lock().unwrap();
let totals = attempts
{
let attempts = meter.attempts.lock().unwrap();
let totals = attempts
.iter()
.map(|r| (r.schema.as_str(), r.rows, r.value))
.collect::<Vec<_>>();
assert_eq!(
totals,
if enabled {
vec![("a", 4, 0), ("b", 2, 0)]
} else {
vec![]
}
);
}
meter.attempts.lock().unwrap().clear();
meter.reject.store(false, Ordering::Relaxed);
let ctx = Arc::new(QueryContext::with_channel(
CATALOG,
"logical",
Channel::Otlp,
));
let admitted = admit_write(2, &ctx).await.unwrap();
let mut requests = RowInsertRequests {
inserts: vec![RowInsertRequest {
table_name: "metric".to_string(),
rows: Some(Rows {
schema: vec![],
rows: vec![api::v1::Row::default(); 2],
}),
}],
};
let original = requests.clone();
let cost = Inserter::meter_row_inserts(&mut requests, &admitted)
.await
.unwrap();
assert_eq!(cost, if enabled { 17 } else { 0 });
assert_eq!(requests, original);
let records = meter
.attempts
.lock()
.unwrap()
.iter()
.map(|r| (r.schema.as_str(), r.rows, r.value))
.map(|record| (record.rows, record.value))
.collect::<Vec<_>>();
assert_eq!(
totals,
records,
if enabled {
vec![("a", 4, 0), ("b", 2, 0)]
vec![(2, 0), (0, 17)]
} else {
vec![]
}
);
meter.reject.store(true, Ordering::Relaxed);
let result = Inserter::meter_row_inserts(&mut requests, &admitted).await;
if enabled {
assert_eq!(result.unwrap_err().status_code(), StatusCode::RateLimited);
} else {
assert_eq!(result.unwrap(), 0);
}
assert_eq!(requests, original);
}
#[test]
@@ -2570,6 +2739,142 @@ mod tests {
)
}
#[tokio::test]
async fn test_logical_batcher_eligibility() {
use catalog::RegisterTableRequest;
use catalog::memory::MemoryCatalogManager;
use common_meta::instruction::{CacheIdent, CreateFlow};
use common_meta::kv_backend::KvBackendRef;
use common_meta::kv_backend::memory::MemoryKvBackend;
use datatypes::schema::{ColumnDefaultConstraint, SchemaBuilder};
let requests = RowInsertRequests {
inserts: vec![RowInsertRequest {
table_name: "test_table".to_string(),
rows: None,
}],
};
let original =
make_table_ref_with_schema("ts", "value", ConcreteDataType::float64_datatype())
.table_info();
for case in [
"eligible",
"physical",
"ordinary",
"instant",
"disabled",
"hint",
"flow",
"required_tag",
"default",
] {
let mut info = (*original).clone();
info.meta.engine = METRIC_ENGINE_NAME.to_string();
info.meta.options.extra_options.insert(
LOGICAL_TABLE_METADATA_KEY.to_string(),
"physical".to_string(),
);
let mut ctx = QueryContext::arc().fork();
match case {
"physical" => {
info.meta
.options
.extra_options
.insert(LOGICAL_TABLE_METADATA_KEY.to_string(), "other".to_string());
}
"ordinary" => info.meta.engine = "mito".to_string(),
"instant" => info.meta.options.ttl = Some(common_time::ttl::TimeToLive::Instant),
"hint" => ctx.set_extension(AUTO_CREATE_TABLE_KEY, "false"),
"required_tag" | "default" => {
let mut columns = info.meta.schema.column_schemas().to_vec();
let mut tag = ColumnSchema::new(
"tag",
ConcreteDataType::string_datatype(),
case != "required_tag",
);
if case == "default" {
tag = tag
.with_default_constraint(Some(ColumnDefaultConstraint::null_value()))
.unwrap();
}
columns.push(tag);
info.meta.schema = Arc::new(
SchemaBuilder::try_from_columns(columns)
.unwrap()
.build()
.unwrap(),
);
info.meta.primary_key_indices = vec![2];
}
_ => {}
}
let catalog = MemoryCatalogManager::with_default_setup();
let table = Arc::new(table::Table::new(
Arc::new(info),
table::metadata::FilterPushDownType::Unsupported,
Arc::new(DummyDataSource),
));
catalog
.register_table_sync(RegisterTableRequest {
catalog: DEFAULT_CATALOG_NAME.to_string(),
schema: DEFAULT_SCHEMA_NAME.to_string(),
table_name: "test_table".to_string(),
table_id: 1,
table,
})
.unwrap();
let mut inserter = batcher_test_inserter().await;
inserter.catalog_manager = catalog;
inserter.auto_create_table = case != "disabled";
let kv_backend: KvBackendRef = Arc::new(MemoryKvBackend::default());
inserter.table_flownode_set_cache = Arc::new(new_table_flownode_set_cache(
String::new(),
Cache::new(10),
kv_backend,
));
if case == "flow" {
inserter
.table_flownode_set_cache
.invalidate(&[CacheIdent::CreateFlow(CreateFlow {
flow_id: 1,
source_table_ids: vec![1],
partition_to_peer_mapping: vec![(0, Peer::empty(1))],
})])
.await
.unwrap();
}
assert_eq!(
inserter
.can_batch_metric_rows(&requests, &Arc::new(ctx), "physical")
.await
.unwrap(),
case == "eligible",
"{case}"
);
}
}
#[tokio::test]
async fn test_batcher_meter_preserves_request() {
let mut requests = RowInsertRequests {
inserts: vec![RowInsertRequest {
table_name: "sample".to_string(),
rows: Some(Rows {
schema: vec![],
rows: vec![api::v1::Row {
values: vec![Value {
value_data: Some(api::v1::value::ValueData::F64Value(1.5)),
}],
}],
}),
}],
};
let expected = requests.clone();
Inserter::meter_row_inserts(&mut requests, &QueryContext::arc())
.await
.unwrap();
assert_eq!(requests, expected);
}
#[tokio::test]
async fn test_instant_table_bypasses_batcher() {
let batcher: Arc<dyn PendingRowsBatcher> = Arc::new(UnexpectedBatcher);
+18
View File
@@ -21,3 +21,21 @@ pub mod table;
#[cfg(test)]
mod test_util;
/// Controls whether batching waits for storage before replying to the client.
const PENDING_ROWS_BATCH_SYNC_ENV: &str = "PENDING_ROWS_BATCH_SYNC";
/// Returns whether pending-row batch submissions wait for the flush result
/// before replying to the client (synchronous mode), controlled by the
/// `PENDING_ROWS_BATCH_SYNC` environment variable and defaulting to `true`.
///
/// Callers that reason about how long a remote write request may block (e.g.
/// the frontend HTTP timeout fallback) must consult this instead of
/// duplicating the env lookup.
pub fn pending_rows_batch_sync_enabled() -> bool {
std::env::var(PENDING_ROWS_BATCH_SYNC_ENV)
.ok()
.as_deref()
.and_then(|v| v.parse::<bool>().ok())
.unwrap_or(true)
}
+23 -22
View File
@@ -21,6 +21,7 @@ mod tables;
#[cfg(test)]
mod test_util;
use std::future::{Future, ready};
use std::num::NonZeroUsize;
use std::sync::Arc;
use std::time::Duration;
@@ -54,6 +55,7 @@ pub use crate::batcher::logical_table::region_write::{
pub use crate::batcher::logical_table::tables::{
PendingRowsSchemaAlterer, PendingRowsSchemaAltererRef,
};
use crate::batcher::pending_rows_batch_sync_enabled;
use crate::error;
use crate::error::{Error, Result};
use crate::metrics::{
@@ -62,24 +64,6 @@ use crate::metrics::{
const PHYSICAL_TABLE_KEY: &str = "physical_table";
/// Whether wait for ingestion result before reply to client.
const PENDING_ROWS_BATCH_SYNC_ENV: &str = "PENDING_ROWS_BATCH_SYNC";
/// Returns whether pending-row batch submissions wait for the flush result
/// before replying to the client (synchronous mode), controlled by the
/// `PENDING_ROWS_BATCH_SYNC` environment variable and defaulting to `true`.
///
/// Callers that reason about how long a remote write request may block (e.g.
/// the frontend HTTP timeout fallback) must consult this instead of
/// duplicating the env lookup.
pub fn pending_rows_batch_sync_enabled() -> bool {
std::env::var(PENDING_ROWS_BATCH_SYNC_ENV)
.ok()
.as_deref()
.and_then(|v| v.parse::<bool>().ok())
.unwrap_or(true)
}
const WORKER_IDLE_TIMEOUT_MULTIPLIER: u32 = 3;
#[derive(Debug, Clone, Hash, Eq, PartialEq)]
@@ -182,14 +166,31 @@ impl LogicalTablePendingRowsBatcher {
impl LogicalTablePendingRowsBatcher {
pub async fn submit(&self, requests: RowInsertRequests, ctx: QueryContextRef) -> Result<u64> {
self.submit_with(requests, ctx, |_| ready(Ok(())))
.await
.map(|(rows, ())| rows)
}
/// 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>(
&self,
requests: RowInsertRequests,
ctx: QueryContextRef,
after_prepare: impl FnOnce(RowInsertRequests) -> F,
) -> Result<(u64, T)>
where
F: Future<Output = Result<T>> + Send,
{
let (table_batches, total_rows) = {
let _timer = PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED
.with_label_values(&["submit_build_and_align"])
.start_timer();
self.build_and_align_table_batches(requests, &ctx).await?
self.build_and_align_table_batches(&requests, &ctx).await?
};
let prepared = after_prepare(requests).await?;
if total_rows == 0 {
return Ok(0);
return Ok((0, prepared));
}
// Flushes dispatch directly to datanodes, so admit once before enqueueing.
@@ -269,9 +270,9 @@ impl LogicalTablePendingRowsBatcher {
};
result
.context(error::SubmitBatchSnafu)
.map(|()| total_rows as u64)
.map(|()| (total_rows as u64, prepared))
} else {
Ok(total_rows as u64)
Ok((total_rows as u64, prepared))
}
}
}
+22 -22
View File
@@ -75,7 +75,7 @@ impl LogicalTablePendingRowsBatcher {
/// RecordBatches.
pub(in crate::batcher::logical_table) async fn build_and_align_table_batches(
&self,
requests: RowInsertRequests,
requests: &RowInsertRequests,
ctx: &QueryContextRef,
) -> Result<(Vec<(String, u32, RecordBatchWithTsIdx)>, usize)> {
let catalog = ctx.current_catalog().to_string();
@@ -113,13 +113,13 @@ impl LogicalTablePendingRowsBatcher {
/// Extracts non-empty `(table_name, rows)` pairs and computes total row
/// count across the retained entries.
pub(in crate::batcher::logical_table) fn collect_non_empty_table_rows(
requests: RowInsertRequests,
) -> (Vec<(String, Rows)>, usize) {
let mut table_rows: Vec<(String, Rows)> = Vec::with_capacity(requests.inserts.len());
requests: &RowInsertRequests,
) -> (Vec<(&str, &Rows)>, usize) {
let mut table_rows: Vec<(&str, &Rows)> = Vec::with_capacity(requests.inserts.len());
let mut total_rows = 0;
for request in requests.inserts {
let Some(rows) = request.rows else {
for request in &requests.inserts {
let Some(rows) = &request.rows else {
continue;
};
if rows.rows.is_empty() {
@@ -127,7 +127,7 @@ impl LogicalTablePendingRowsBatcher {
}
total_rows += rows.rows.len();
table_rows.push((request.table_name, rows));
table_rows.push((request.table_name.as_str(), rows));
}
(table_rows, total_rows)
@@ -137,15 +137,15 @@ impl LogicalTablePendingRowsBatcher {
impl LogicalTablePendingRowsBatcher {
/// Returns unique `(table_name, proto_schema)` pairs while keeping the
/// first-seen schema for duplicate table names.
pub(in crate::batcher::logical_table) fn collect_unique_table_schemas(
table_rows: &[(String, Rows)],
) -> Result<Vec<(&str, &[ColumnSchema])>> {
pub(in crate::batcher::logical_table) fn collect_unique_table_schemas<'a>(
table_rows: &[(&'a str, &'a Rows)],
) -> Result<Vec<(&'a str, &'a [ColumnSchema])>> {
let mut unique_tables: Vec<(&str, &[ColumnSchema])> = Vec::with_capacity(table_rows.len());
let mut seen = HashSet::new();
for (table_name, rows) in table_rows {
if seen.insert(table_name.as_str()) {
unique_tables.push((table_name.as_str(), &rows.schema));
if seen.insert(*table_name) {
unique_tables.push((*table_name, &rows.schema));
} else {
// table_rows should group rows by table name.
return error::InvalidPromRemoteRequestSnafu {
@@ -233,7 +233,7 @@ impl LogicalTablePendingRowsBatcher {
catalog: &str,
schema: &str,
ctx: &QueryContextRef,
table_rows: &[(String, Rows)],
table_rows: &[(&str, &Rows)],
plan: &mut TableResolutionPlan,
) -> Result<()> {
if plan.tables_to_create.is_empty() {
@@ -298,7 +298,7 @@ impl LogicalTablePendingRowsBatcher {
/// For newly created tables, re-checks all row schemas and appends alter
/// operations when additional tag columns are still missing.
pub(in crate::batcher::logical_table) fn enqueue_alter_for_new_tables(
table_rows: &[(String, Rows)],
table_rows: &[(&str, &Rows)],
plan: &mut TableResolutionPlan,
) -> Result<()> {
let created_tables: HashSet<&str> = plan
@@ -308,11 +308,11 @@ impl LogicalTablePendingRowsBatcher {
.collect();
for (table_name, rows) in table_rows {
if !created_tables.contains(table_name.as_str()) {
if !created_tables.contains(table_name) {
continue;
}
let Some((region_schema, _)) = plan.region_schemas.get(table_name) else {
let Some((region_schema, _)) = plan.region_schemas.get(*table_name) else {
continue;
};
@@ -321,13 +321,13 @@ impl LogicalTablePendingRowsBatcher {
|| plan
.tables_to_alter
.iter()
.any(|(existing_name, _)| existing_name == table_name)
.any(|(existing_name, _)| existing_name == *table_name)
{
continue;
}
plan.tables_to_alter
.push((table_name.clone(), missing_columns));
.push((table_name.to_string(), missing_columns));
}
Ok(())
@@ -397,13 +397,13 @@ impl LogicalTablePendingRowsBatcher {
/// Converts proto rows to `RecordBatch` values aligned to resolved region
/// schemas and returns `(table_name, table_id, batch)` tuples.
pub(in crate::batcher::logical_table) fn build_aligned_batches(
table_rows: &[(String, Rows)],
table_rows: &[(&str, &Rows)],
region_schemas: &HashMap<String, (Arc<ArrowSchema>, u32)>,
) -> Result<Vec<(String, u32, RecordBatchWithTsIdx)>> {
let mut aligned_batches = Vec::with_capacity(table_rows.len());
for (table_name, rows) in table_rows {
let (region_schema, table_id) =
region_schemas.get(table_name).cloned().with_context(|| {
region_schemas.get(*table_name).cloned().with_context(|| {
error::UnexpectedResultSnafu {
reason: format!("Region schema not resolved for table: {}", table_name),
}
@@ -415,7 +415,7 @@ impl LogicalTablePendingRowsBatcher {
.start_timer();
rows_to_aligned_record_batch(rows, region_schema.as_ref())?
};
aligned_batches.push((table_name.clone(), table_id, record_batch));
aligned_batches.push((table_name.to_string(), table_id, record_batch));
}
Ok(aligned_batches)
@@ -531,7 +531,7 @@ mod tests {
};
let (table_rows, total_rows) =
LogicalTablePendingRowsBatcher::collect_non_empty_table_rows(requests);
LogicalTablePendingRowsBatcher::collect_non_empty_table_rows(&requests);
assert_eq!(2, total_rows);
assert_eq!(1, table_rows.len());
+114 -3
View File
@@ -38,6 +38,7 @@ use snafu::ResultExt;
use table::metadata::TableInfoRef;
use tokio::sync::{OwnedSemaphorePermit, Semaphore, broadcast, oneshot};
use crate::batcher::pending_rows_batch_sync_enabled;
use crate::batcher::table::flow_notifier::FlowNotifier;
use crate::batcher::table::metrics::PENDING_WORKERS;
use crate::batcher::table::pending_batch::PendingBatch;
@@ -70,6 +71,7 @@ pub struct TablePendingRowsBatcher {
flush_limiter: FlushLimiter,
request_limiter: RequestLimiter,
worker_channel_capacity: usize,
pending_rows_batch_sync: bool,
worker_idle_timeout: Duration,
inserter: Arc<Inserter>,
flow_notifier: FlowNotifier,
@@ -110,6 +112,7 @@ impl TablePendingRowsBatcher {
flush_limiter,
request_limiter,
worker_channel_capacity,
pending_rows_batch_sync: pending_rows_batch_sync_enabled(),
worker_idle_timeout: flush_interval.checked_mul(3).unwrap_or(flush_interval),
inserter,
flow_notifier,
@@ -160,8 +163,8 @@ impl PendingRowsBatcher for TablePendingRowsBatcher {
})
}
/// Waits for completed bulk writes. Cancellation does not retract an admitted
/// submission. Combined failures affect all waiters and may be partial writes.
/// Acknowledges according to the global batching policy. Cancellation does
/// not retract an admitted submission. Flush failures may be partial writes.
async fn submit(
&self,
table_info: TableInfoRef,
@@ -210,6 +213,9 @@ impl PendingRowsBatcher for TablePendingRowsBatcher {
}
.fail();
}
if !self.pending_rows_batch_sync {
return Ok(total_rows);
}
response_rx
.await
.map_err(|_| {
@@ -236,7 +242,7 @@ mod tests {
use api::v1::region::{RegionRequest, bulk_insert_request, region_request};
use api::v1::value::ValueData;
use api::v1::{ColumnDataType, Row, Rows, Value};
use arrow::array::{Int32Array, TimestampMillisecondArray};
use arrow::array::{ArrayRef, Int32Array, TimestampMillisecondArray};
use arrow::datatypes::Schema as ArrowSchema;
use arrow::record_batch::RecordBatch;
use catalog::memory::MemoryCatalogManager;
@@ -320,6 +326,111 @@ mod tests {
}
}
#[tokio::test]
async fn test_acknowledgement_preserves_admission() {
use operator::error::UnexpectedSnafu;
use crate::batcher::pending_rows_batch_sync_enabled;
use crate::batcher::table::batch_key_from_ctx;
use crate::batcher::table::pending_batch::notify_batches;
for sync in [false, true] {
for fail in [false, true] {
let backend = prepare_mocked_backend().await;
let nodes = Arc::new(MockDatanodeManager::new(BulkHandler {
requests: Arc::new(Mutex::new(Vec::new())),
report_missing_row: false,
expected_skip_wal: false,
expected_schema: "public".to_string(),
}));
let inserter = Arc::new(Inserter::new(
MemoryCatalogManager::new(),
create_partition_rule_manager(backend).await,
nodes,
mock_table_flownode_cache(1, vec![]).await,
true,
));
let mut batcher = TablePendingRowsBatcher::try_new(
Duration::from_secs(3600),
1,
1,
1,
1,
NonZeroUsize::new(1).unwrap(),
inserter,
)
.unwrap();
assert_eq!(
batcher.pending_rows_batch_sync,
pending_rows_batch_sync_enabled()
);
Arc::get_mut(&mut batcher).unwrap().pending_rows_batch_sync = sync;
let table = Arc::new(new_test_table_info(1, "ack", [0].into_iter()));
let ctx = QueryContext::arc();
let key = batch_key_from_ctx(&table.name, &ctx);
// Hold the worker command to control completion independently of scheduling.
let (_, receiver) = batcher.workers.get_or_create(key, 1).await;
let mut receiver = receiver.unwrap();
let batch = RecordBatch::try_from_iter(vec![(
"a",
Arc::new(Int32Array::from(vec![1])) as ArrayRef,
)])
.unwrap();
let permit = batcher.acquire().await.unwrap();
let submitter = batcher.clone();
let submitted =
tokio::spawn(async move { submitter.submit(table, batch, ctx, permit).await });
let WorkerCommand::Submit(pending) =
timeout(Duration::from_secs(5), receiver.recv())
.await
.unwrap()
.unwrap();
let mut submitted = Some(submitted);
if sync {
assert!(!submitted.as_ref().unwrap().is_finished());
} else {
assert_eq!(
timeout(Duration::from_secs(5), submitted.take().unwrap())
.await
.unwrap()
.unwrap()
.unwrap(),
1
);
}
// Early acknowledgement must not release capacity before completion.
let acquire = batcher.acquire();
tokio::pin!(acquire);
assert!(futures::poll!(acquire.as_mut()).is_pending());
let result = if fail {
Err(Arc::new(
UnexpectedSnafu {
violated: "flush failed".to_string(),
}
.build(),
))
} else {
Ok(())
};
notify_batches(vec![pending], result);
if let Some(submitted) = submitted {
let result = timeout(Duration::from_secs(5), submitted)
.await
.unwrap()
.unwrap();
assert_eq!(result.is_err(), fail);
if !fail {
assert_eq!(result.unwrap(), 1);
}
}
timeout(Duration::from_secs(5), acquire)
.await
.unwrap()
.unwrap();
}
}
}
fn rows(value: i32) -> Rows {
Rows {
schema: vec![
+2
View File
@@ -17,6 +17,7 @@ use std::time::Instant;
use arrow::compute::concat_batches;
use arrow::record_batch::RecordBatch;
use common_telemetry::error;
use operator::error::{ComputeArrowSnafu, Result, UnexpectedSnafu};
use operator::insert::Inserter;
use operator::metrics::DIST_INGEST_ROW_COUNT;
@@ -55,6 +56,7 @@ pub(in crate::batcher::table) async fn flush_batch(
notifier.notify(table, &combined);
}
Err(error) => {
error!(error; "Failed to flush table batch, rows: {}", batch.total_rows);
FLUSH_FAILURES.inc();
FLUSH_DROPPED_ROWS.inc_by(batch.total_rows as u64);
notify_batches(batch.submissions, Err(Arc::new(error)));
+62 -1
View File
@@ -215,6 +215,7 @@ pub struct HttpServer {
// server configs
options: HttpOptions,
batching_protocols: Vec<BatchingProtocol>,
logical_batching_protocols: Vec<BatchingProtocol>,
bind_addr: Option<SocketAddr>,
/// What this server instance exposes. See [`HttpServerKind`].
kind: HttpServerKind,
@@ -235,6 +236,9 @@ pub fn is_api_listener_path(path: &str) -> bool {
is_namespace(path, HTTP_API_PREFIX_WITHOUT_TRAILING_SLASH) || is_namespace(path, "/dashboard")
}
#[derive(Clone)]
struct LogicalBatchingProtocols(Vec<BatchingProtocol>);
/// Sets a local-only write selector after authentication creates the context.
async fn set_http_write_batching(
State(protocol): State<BatchingProtocol>,
@@ -245,8 +249,13 @@ async fn set_http_write_batching(
.extensions()
.get::<Arc<Vec<BatchingProtocol>>>()
.is_some_and(|protocols| protocols.contains(&protocol));
let logical_enabled = req
.extensions()
.get::<LogicalBatchingProtocols>()
.is_some_and(|protocols| protocols.0.contains(&protocol));
if let Some(ctx) = req.extensions_mut().get_mut::<QueryContext>() {
ctx.set_batching_enabled(enabled);
ctx.set_logical_batching_enabled(logical_enabled);
}
next.run(req).await
}
@@ -650,6 +659,7 @@ pub struct DashboardState {
pub struct HttpServerBuilder {
options: HttpOptions,
batching_protocols: Vec<BatchingProtocol>,
logical_batching_protocols: Vec<BatchingProtocol>,
user_provider: Option<UserProviderRef>,
router: Router,
memory_limiter: ServerMemoryLimiter,
@@ -660,13 +670,20 @@ impl HttpServerBuilder {
Self {
options,
batching_protocols: Vec::new(),
logical_batching_protocols: Vec::new(),
user_provider: None,
router: Router::new(),
memory_limiter: ServerMemoryLimiter::default(),
}
}
/// Selects HTTP write protocols allowed to use the shared batcher.
/// Selects HTTP protocols allowed to use logical-table batching.
pub fn with_logical_batching_protocols(mut self, protocols: Vec<BatchingProtocol>) -> Self {
self.logical_batching_protocols = protocols;
self
}
/// Selects HTTP protocols allowed to use ordinary-table batching.
pub fn with_batching_protocols(mut self, protocols: Vec<BatchingProtocol>) -> Self {
self.batching_protocols = protocols;
self
@@ -894,6 +911,7 @@ impl HttpServerBuilder {
HttpServer {
options: self.options,
batching_protocols: self.batching_protocols.clone(),
logical_batching_protocols: self.logical_batching_protocols.clone(),
user_provider: self.user_provider,
shutdown_tx: Mutex::new(None),
router: StdMutex::new(self.router),
@@ -928,6 +946,7 @@ impl HttpServerBuilder {
let internal = HttpServer {
options: self.options,
batching_protocols: self.batching_protocols.clone(),
logical_batching_protocols: self.logical_batching_protocols.clone(),
user_provider: self.user_provider.clone(),
shutdown_tx: Mutex::new(None),
router: StdMutex::new(self.router.clone()),
@@ -946,6 +965,7 @@ impl HttpServerBuilder {
Some(HttpServer {
options: api_options,
batching_protocols: self.batching_protocols,
logical_batching_protocols: self.logical_batching_protocols,
user_provider: self.user_provider.clone(),
shutdown_tx: Mutex::new(None),
router: StdMutex::new(self.router),
@@ -1104,6 +1124,9 @@ impl HttpServer {
authorize::check_http_auth,
))
.layer(Extension(Arc::new(self.batching_protocols.clone())))
.layer(Extension(LogicalBatchingProtocols(
self.logical_batching_protocols.clone(),
)))
.layer(middleware::from_fn(hints::extract_hints))
.layer(middleware::from_fn(client_ip::log_error_with_client_ip))
.layer(middleware::from_fn(
@@ -2462,6 +2485,44 @@ mod batching_tests {
use crate::opentsdb::codec::DataPoint;
use crate::query_handler::{InfluxdbLineProtocolHandler, OpentsdbProtocolHandler};
#[tokio::test]
async fn test_batching_selectors_are_independent() {
use axum::routing::post;
use axum::{Extension, Json, Router, middleware};
use session::context::{QueryContext, QueryContextBuilder};
use crate::http::{LogicalBatchingProtocols, set_http_write_batching};
for table in [false, true] {
for logical in [false, true] {
let app = Router::new()
.route(
"/",
post(|Extension(ctx): Extension<QueryContext>| async move {
Json([ctx.batching_enabled(), ctx.logical_batching_enabled()])
}),
)
.route_layer(middleware::from_fn_with_state(
BatchingProtocol::Otlp,
set_http_write_batching,
))
.layer(Extension(QueryContextBuilder::default().build()))
.layer(Extension(Arc::new(if table {
vec![BatchingProtocol::Otlp]
} else {
vec![]
})))
.layer(Extension(LogicalBatchingProtocols(if logical {
vec![BatchingProtocol::Otlp]
} else {
vec![]
})));
let client = TestClient::new(app).await;
let actual: [bool; 2] = client.post("/").send().await.json().await;
assert_eq!(actual, [table, logical]);
}
}
}
#[test]
fn test_protocol_names_reject_unknown_values() {
assert_eq!(
+22 -2
View File
@@ -82,6 +82,9 @@ pub struct QueryContext {
/// Local-only write batching selection; never transported in protobuf extensions.
#[builder(default)]
batching_enabled: bool,
/// Local-only opt-in for logical metric writes, independent of ordinary-table batching.
#[builder(default)]
logical_batching_enabled: bool,
/// Track which protocol the query comes from.
#[builder(default)]
channel: Channel,
@@ -458,7 +461,17 @@ impl QueryContext {
&self.configuration_parameter
}
/// Whether the local HTTP entry point selected write batching.
/// Whether the local HTTP entry point selected logical-table batching.
pub fn logical_batching_enabled(&self) -> bool {
self.logical_batching_enabled
}
/// Sets local logical-table batching selection without adding a wire-visible extension.
pub fn set_logical_batching_enabled(&mut self, enabled: bool) {
self.logical_batching_enabled = enabled;
}
/// Whether the local HTTP entry point selected ordinary-table batching.
pub fn batching_enabled(&self) -> bool {
self.batching_enabled
}
@@ -634,6 +647,7 @@ impl QueryContextBuilder {
channel,
batching_enabled: self.batching_enabled.unwrap_or_default(),
admitted_write: None,
logical_batching_enabled: self.logical_batching_enabled.unwrap_or_default(),
process_id: self.process_id.unwrap_or_default(),
conn_info: self.conn_info.unwrap_or_default(),
protocol_ctx: self.protocol_ctx.unwrap_or_default(),
@@ -923,11 +937,17 @@ mod test {
fn test_batching_selection_is_local_only() {
let mut ctx = QueryContextBuilder::default().build();
assert!(!ctx.batching_enabled());
assert!(!ctx.logical_batching_enabled());
ctx.set_logical_batching_enabled(true);
assert!(ctx.clone().logical_batching_enabled());
assert!(ctx.fork().logical_batching_enabled());
ctx.set_batching_enabled(true);
assert!(ctx.clone().batching_enabled());
assert!(ctx.fork().batching_enabled());
let wire: api::v1::QueryContext = ctx.into();
assert!(!QueryContext::from(wire).batching_enabled());
let restored = QueryContext::from(wire);
assert!(!restored.batching_enabled());
assert!(!restored.logical_batching_enabled());
let ctx = QueryContextBuilder::default()
.set_extension("batching_enabled".to_string(), "true".to_string())
.build();
+68 -5
View File
@@ -58,7 +58,7 @@ pub struct StandaloneOptions {
pub postgres: PostgresOptions,
pub opentsdb: OpentsdbOptions,
pub influxdb: InfluxdbOptions,
/// Shared experimental ordinary-table batching; independent of Prom batching.
/// Ordinary-table batching with independent logical-table controls.
pub pending_rows_batcher: PendingRowsBatcherOptions,
pub jaeger: JaegerOptions,
pub otlp: OtlpOptions,
@@ -137,6 +137,7 @@ impl Configurable for StandaloneOptions {
"wal.broker_endpoints",
"event_recorder.event_types",
"pending_rows_batcher.protocols",
"pending_rows_batcher.logical_table.protocols",
])
}
}
@@ -219,6 +220,61 @@ mod tests {
use crate::options::*;
#[test]
fn test_logical_batcher_config_forwarding() {
let opts: StandaloneOptions = toml::from_str("[pending_rows_batcher]\nprotocols = ['influxdb']\n[pending_rows_batcher.logical_table]\nprotocols = ['otlp', 'prom']\npending_rows_flush_interval = '10ms'").unwrap();
let frontend = opts.frontend_options();
assert_eq!(frontend.pending_rows_batcher, opts.pending_rows_batcher);
assert_eq!(
frontend.pending_rows_batcher.logical_table,
opts.pending_rows_batcher.logical_table
);
let restored: StandaloneOptions = toml::from_str(&toml::to_string(&opts).unwrap()).unwrap();
assert_eq!(
restored.pending_rows_batcher.logical_table,
opts.pending_rows_batcher.logical_table
);
assert!(
toml::from_str::<StandaloneOptions>(
"[pending_rows_batcher.logical_table]\nprotocols = ['logs']"
)
.is_err()
);
}
#[test]
fn test_logical_batcher_protocols_from_env() {
temp_env::with_vars(
[
(
"STANDALONE_LOGICAL_TEST__PENDING_ROWS_BATCHER__PROTOCOLS",
Some("otlp,influxdb"),
),
(
"STANDALONE_LOGICAL_TEST__PENDING_ROWS_BATCHER__LOGICAL_TABLE__PROTOCOLS",
Some("prom,otlp"),
),
],
|| {
let options =
StandaloneOptions::load_layered_options(None, "STANDALONE_LOGICAL_TEST")
.unwrap();
assert_eq!(options.pending_rows_batcher.table.protocols.len(), 2);
assert_eq!(
options
.pending_rows_batcher
.logical_table
.unwrap()
.protocols,
vec![
servers::http::BatchingProtocol::Prom,
servers::http::BatchingProtocol::Otlp
]
);
},
);
}
#[test]
fn test_batcher_protocols_from_env() {
temp_env::with_vars(
@@ -231,7 +287,7 @@ mod tests {
StandaloneOptions::load_layered_options(None, "STANDALONE_BATCHER_TEST")
.unwrap();
assert_eq!(
options.pending_rows_batcher.protocols,
options.pending_rows_batcher.table.protocols,
vec![
servers::http::BatchingProtocol::Influxdb,
servers::http::BatchingProtocol::HttpSql
@@ -247,6 +303,7 @@ mod tests {
assert!(
!defaults
.pending_rows_batcher
.table
.pending_rows_batching_enabled()
);
let options: StandaloneOptions = toml::from_str(
@@ -259,16 +316,22 @@ flow_notification_queue_capacity = 17
"#,
)
.unwrap();
assert_eq!(options.pending_rows_batcher.max_batch_rows, 25);
assert_eq!(options.pending_rows_batcher.protocols.len(), 2);
assert_eq!(options.pending_rows_batcher.table.max_batch_rows, 25);
assert_eq!(options.pending_rows_batcher.table.protocols.len(), 2);
assert_eq!(
options
.pending_rows_batcher
.table
.flow_notification_queue_capacity
.get(),
17
);
assert!(options.pending_rows_batcher.pending_rows_batching_enabled());
assert!(
options
.pending_rows_batcher
.table
.pending_rows_batching_enabled()
);
let serialized = toml::to_string(&options).unwrap();
let parsed: StandaloneOptions = toml::from_str(&serialized).unwrap();
assert_eq!(options.influxdb, parsed.influxdb);
+227
View File
@@ -503,6 +503,233 @@ WITH(
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn test_otlp_logical_batcher_alignment() {
use std::time::Duration;
use common_base::Plugins;
use frontend::server::Services;
use frontend::service_config::pending_rows_batcher::BatcherOptions;
use otel_arrow_rust::proto::opentelemetry::metrics::v1::{
ExponentialHistogram, ExponentialHistogramDataPoint, exponential_histogram_data_point,
};
use prost::Message;
use servers::batcher::pending_rows_batch_sync_enabled;
use servers::http::BatchingProtocol;
use servers::http::test_helpers::TestClient;
use session::protocol_ctx::{OtlpMetricCtx, ProtocolCtx};
// Fixed workload, only the batcher opt-in changes. Compare stored rows,
// schema evolution and visibility under the configured acknowledgement policy.
let mut results = Vec::new();
for enabled in [false, true] {
let standalone = GreptimeDbStandaloneBuilder::new(&format!("otlp_logical_{enabled}"))
.with_logical_batcher(BatcherOptions {
protocols: if enabled {
vec![BatchingProtocol::Otlp]
} else {
vec![]
},
pending_rows_flush_interval: Duration::from_millis(5),
..Default::default()
})
.build()
.await;
let instance = standalone.fe_instance();
let mut options = standalone.opts.clone();
options.otlp.experimental_enable_exponential_histogram = true;
let services = Services::new(options.clone(), instance.clone(), Plugins::default());
let server = services
.http_server_builder(
&options.frontend_options(),
services.server_memory_limiter.clone(),
)
.build();
let client = TestClient::new(server.build(server.make_app()).unwrap()).await;
let mut ctx = QueryContext::with(DEFAULT_CATALOG_NAME, "public");
ctx.set_logical_batching_enabled(true);
ctx.set_protocol_ctx(ProtocolCtx::OtlpMetric(OtlpMetricCtx {
with_metric_engine: true,
..Default::default()
}));
let ctx = Arc::new(ctx);
let submissions = servers::metrics::PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED
.with_label_values(&["submit_wait_flush_result"]);
let before = submissions.get_sample_count();
for (ts, value) in [(60, 10), (120, 20)] {
let mut request = build_sum_request(
"batch.alignment",
AggregationTemporality::Cumulative,
&[(ts, value)],
);
if ts == 120
&& let Some(metric::Data::Sum(sum)) =
&mut request.resource_metrics[0].scope_metrics[0].metrics[0].data
{
sum.data_points[0].attributes.push(keyvalue("extra", "new"));
}
let response = client
.post("/v1/otlp/v1/metrics")
.header("content-type", "application/x-protobuf")
.body(request.encode_to_vec())
.send()
.await;
assert_eq!(response.status().as_u16(), 200);
}
if enabled {
assert_eq!(
submissions.get_sample_count() - before,
if pending_rows_batch_sync_enabled() {
2
} else {
0
},
"logical submissions must follow the global acknowledgement policy"
);
}
let sql = "SELECT greptime_timestamp, greptime_value, stream, extra FROM batch_alignment_total ORDER BY greptime_timestamp";
let batches = tokio::time::timeout(Duration::from_secs(5), async {
loop {
let output = instance.do_query(sql, ctx.clone()).await.remove(0).unwrap();
let OutputData::Stream(stream) = output.data else {
panic!("expected stream")
};
let batches = RecordBatches::try_collect(stream).await.unwrap();
if !enabled
|| pending_rows_batch_sync_enabled()
|| batches.iter().map(|batch| batch.num_rows()).sum::<usize>() == 2
{
break batches;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.unwrap();
assert_eq!(
batches.iter().map(|batch| batch.num_rows()).sum::<usize>(),
2
);
results.push(batches.pretty_print().unwrap());
let before_histograms = submissions.get_sample_count();
// First create a histogram table, then reuse it in a mixed export.
// Neither request may enter the scalar-only logical batcher.
for mixed in [false, true] {
let mut request = build_sum_request(
"batch.mixed",
AggregationTemporality::Cumulative,
&[(180, 30)],
);
let metrics = &mut request.resource_metrics[0].scope_metrics[0].metrics;
let mut histogram = metrics[0].clone();
histogram.name = "batch.histogram".to_string();
histogram.data = Some(metric::Data::ExponentialHistogram(ExponentialHistogram {
aggregation_temporality: AggregationTemporality::Cumulative as i32,
data_points: vec![ExponentialHistogramDataPoint {
start_time_unix_nano: 1_000_000_000,
time_unix_nano: if mixed { 4_000_000_000 } else { 3_000_000_000 },
count: 4,
sum: Some(8.0),
zero_count: 1,
positive: Some(exponential_histogram_data_point::Buckets {
offset: -1,
bucket_counts: vec![1, 2],
}),
..Default::default()
}],
}));
if !mixed {
metrics.clear();
}
metrics.push(histogram);
let response = client
.post("/v1/otlp/v1/metrics")
.header("content-type", "application/x-protobuf")
.body(request.encode_to_vec())
.send()
.await;
assert_eq!(
response.status().as_u16(),
200,
"enabled={enabled}, mixed={mixed}"
);
}
assert_eq!(submissions.get_sample_count(), before_histograms);
for (sql, expected_rows) in [
(
"SELECT greptime_timestamp, greptime_native_histogram FROM batch_histogram ORDER BY greptime_timestamp",
2,
),
(
"SELECT greptime_timestamp, greptime_value FROM batch_mixed_total ORDER BY greptime_timestamp",
1,
),
] {
let output = instance.do_query(sql, ctx.clone()).await.remove(0).unwrap();
let OutputData::Stream(stream) = output.data else {
panic!("expected stream")
};
let batches = RecordBatches::try_collect(stream).await.unwrap();
assert_eq!(
batches.iter().map(|batch| batch.num_rows()).sum::<usize>(),
expected_rows
);
results.push(batches.pretty_print().unwrap());
}
// An existing metric attached to another physical table must use
// the original routing, not the request's default physical table.
for sql in [
"CREATE TABLE custom_physical (greptime_timestamp TIMESTAMP TIME INDEX, greptime_value DOUBLE) ENGINE=metric WITH ('physical_metric_table'='')",
"CREATE TABLE custom_total (greptime_timestamp TIMESTAMP(3) TIME INDEX, greptime_value DOUBLE, \"stream\" STRING PRIMARY KEY) ENGINE=metric WITH ('on_physical_table'='custom_physical')",
] {
instance.do_query(sql, ctx.clone()).await.remove(0).unwrap();
}
instance
.metrics(
build_sum_request("custom", AggregationTemporality::Cumulative, &[(60, 7)]),
ctx.clone(),
)
.await
.unwrap();
let output = instance
.do_query("SELECT greptime_value FROM custom_total", ctx.clone())
.await
.remove(0)
.unwrap();
let OutputData::Stream(stream) = output.data else {
panic!("expected stream")
};
assert!(
RecordBatches::try_collect(stream)
.await
.unwrap()
.pretty_print()
.unwrap()
.contains("7.0")
);
// Request-level schema policy cannot be bypassed by batching.
let mut fixed = ctx.fork();
fixed.set_extension("auto_create_table", "false");
assert!(
instance
.metrics(
build_sum_request(
"missing",
AggregationTemporality::Cumulative,
&[(60, 1)]
),
Arc::new(fixed)
)
.await
.is_err()
);
}
assert_eq!(results[..3], results[3..]);
}
#[tokio::test(flavor = "multi_thread")]
pub async fn test_otlp_on_standalone() {
let standalone = GreptimeDbStandaloneBuilder::new("test_standalone_otlp")
+14
View File
@@ -52,6 +52,7 @@ use frontend::frontend::Frontend;
use frontend::instance::Instance;
use frontend::instance::builder::FrontendBuilder;
use frontend::server::Services;
use frontend::service_config::{BatcherOptions, PendingRowsBatcherOptions};
use meta_srv::metasrv::{FLOW_ID_SEQ, TABLE_ID_SEQ};
use servers::grpc::GrpcOptions;
use standalone::options::StandaloneOptions;
@@ -88,6 +89,7 @@ pub struct GreptimeDbStandaloneBuilder {
event_recorder_options: EventRecorderOptions,
auto_create_table: bool,
experimental_metric_export: bool,
logical_batcher: Option<BatcherOptions>,
}
impl GreptimeDbStandaloneBuilder {
@@ -108,6 +110,7 @@ impl GreptimeDbStandaloneBuilder {
event_recorder_options: EventRecorderOptions::default(),
auto_create_table: true,
experimental_metric_export: false,
logical_batcher: None,
}
}
@@ -118,6 +121,13 @@ impl GreptimeDbStandaloneBuilder {
self
}
/// Configures logical-table batching for integration tests.
#[must_use]
pub fn with_logical_batcher(mut self, options: BatcherOptions) -> Self {
self.logical_batcher = Some(options);
self
}
#[must_use]
pub fn with_auto_create_table(mut self, auto_create_table: bool) -> Self {
self.auto_create_table = auto_create_table;
@@ -377,6 +387,10 @@ impl GreptimeDbStandaloneBuilder {
event_recorder: self.event_recorder_options.clone(),
auto_create_table: self.auto_create_table,
experimental_metric_export: self.experimental_metric_export,
pending_rows_batcher: PendingRowsBatcherOptions {
logical_table: self.logical_batcher.clone(),
..Default::default()
},
// Tests cover the descriptor, so they run with it enabled.
otlp: frontend::service_config::OtlpOptions {
experimental_enable_resource_info: true,
+2 -2
View File
@@ -2432,7 +2432,7 @@ protocols = []
pending_rows_flush_interval = "0s"
max_batch_rows = 100000
max_concurrent_flushes = 256
worker_channel_capacity = 65526
worker_channel_capacity = 65536
max_inflight_requests = 3000
flow_notification_queue_capacity = 1024
@@ -2453,7 +2453,7 @@ experimental_enable_prometheus_native_histogram = false
pending_rows_flush_interval = "0s"
max_batch_rows = 100000
max_concurrent_flushes = 256
worker_channel_capacity = 65526
worker_channel_capacity = 65536
max_inflight_requests = 3000
flow_notification_queue_capacity = 1024