diff --git a/config/config.md b/config/config.md index 16384b406f4..d736c24d848 100644 --- a/config/config.md +++ b/config/config.md @@ -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.
When synchronous Prometheus or shared 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`. | +| `http.timeout` | String | `0s` | HTTP request timeout. Set to 0 to disable timeout.
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
`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.
The following units are supported: `B`, `KB`, `KiB`, `MB`, `MiB`, `GB`, `GiB`, `TB`, `TiB`, `PB`, `PiB`.
Set to 0 to disable limit. | | `http.enable_cors` | Bool | `true` | HTTP CORS support, it's turned on by default
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.
Available values: "last_non_null", "last_row". | -| `pending_rows_batcher` | -- | -- | 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. | +| `pending_rows_batcher` | -- | -- | 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.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.
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.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.
Available options:
- strict: deny invalid UTF-8 strings (default).
- lossy: allow invalid UTF-8 strings, replace invalid characters with REPLACEMENT_CHARACTER(U+FFFD).
- 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.
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.
- `raft_engine`: the wal is stored in the local file system by raft-engine.
- `kafka`: it's remote wal that data is stored in Kafka.
- `experimental_object_store`: the wal is stored as objects in an object store.
**Notes: experimental and not supported yet.** | | `wal.dir` | String | Unset | The directory to store the WAL files.
**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.
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.
When synchronous Prometheus or shared 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`. | +| `http.timeout` | String | `0s` | HTTP request timeout. Set to 0 to disable timeout.
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
`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.
The following units are supported: `B`, `KB`, `KiB`, `MB`, `MiB`, `GB`, `GiB`, `TB`, `TiB`, `PB`, `PiB`.
Set to 0 to disable limit. | | `http.enable_cors` | Bool | `true` | HTTP CORS support, it's turned on by default
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.
Available values: "last_non_null", "last_row". | -| `pending_rows_batcher` | -- | -- | 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. | +| `pending_rows_batcher` | -- | -- | 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.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.
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.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.
Available options:
- strict: deny invalid UTF-8 strings (default).
- lossy: allow invalid UTF-8 strings, replace invalid characters with REPLACEMENT_CHARACTER(U+FFFD).
- 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.
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. | diff --git a/config/frontend.example.toml b/config/frontend.example.toml index 043da57db74..275db3105c5 100644 --- a/config/frontend.example.toml +++ b/config/frontend.example.toml @@ -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] diff --git a/config/standalone.example.toml b/config/standalone.example.toml index e59e7ba1b4c..257f48eda2c 100644 --- a/config/standalone.example.toml +++ b/config/standalone.example.toml @@ -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] diff --git a/src/cmd/src/bin/query_perf_fixture/case.rs b/src/cmd/src/bin/query_perf_fixture/case.rs index 0f5239a6b51..271c3928e35 100644 --- a/src/cmd/src/bin/query_perf_fixture/case.rs +++ b/src/cmd/src/bin/query_perf_fixture/case.rs @@ -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 diff --git a/src/cmd/tests/load_config_test.rs b/src/cmd/tests/load_config_test.rs index 67e2e573f38..1c4a83e1397 100644 --- a/src/cmd/tests/load_config_test.rs +++ b/src/cmd/tests/load_config_test.rs @@ -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:: { 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:: { 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, diff --git a/src/frontend/AGENTS.md b/src/frontend/AGENTS.md index d552d73cf39..49a4b8fcd0f 100644 --- a/src/frontend/AGENTS.md +++ b/src/frontend/AGENTS.md @@ -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 diff --git a/src/frontend/src/frontend.rs b/src/frontend/src/frontend.rs index 86ba2cb9bfb..72bf322b1f1 100644 --- a/src/frontend/src/frontend.rs +++ b/src/frontend/src/frontend.rs @@ -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 = Pin> + 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); diff --git a/src/frontend/src/instance.rs b/src/frontend/src/instance.rs index 22d2afaa7bc..92bc74ca10b 100644 --- a/src/frontend/src/instance.rs +++ b/src/frontend/src/instance.rs @@ -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>>>, deleter: DeleterRef, table_metadata_manager: TableMetadataManagerRef, event_recorder: EventRecorderRef, diff --git a/src/frontend/src/instance/builder.rs b/src/frontend/src/instance/builder.rs index 01e9ef034e3..a3cfbcdaffe 100644 --- a/src/frontend/src/instance/builder.rs +++ b/src/frontend/src/instance/builder.rs @@ -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> { - 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) - }; + let build_batcher = |options: &BatcherOptions| -> Option> { + 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) + }; 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, diff --git a/src/frontend/src/instance/logical_batcher.rs b/src/frontend/src/instance/logical_batcher.rs new file mode 100644 index 00000000000..28ed19b3e14 --- /dev/null +++ b/src/frontend/src/instance/logical_batcher.rs @@ -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, 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> { + 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); + +#[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 + } +} diff --git a/src/frontend/src/instance/otlp.rs b/src/frontend/src/instance/otlp.rs index fa40c7eb59a..7a0f706c58d 100644 --- a/src/frontend/src/instance/otlp.rs +++ b/src/frontend/src/instance/otlp.rs @@ -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, diff --git a/src/frontend/src/server.rs b/src/frontend/src/server.rs index 017bba53fba..0a6612ba31c 100644 --- a/src/frontend/src/server.rs +++ b/src/frontend/src/server.rs @@ -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, 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), diff --git a/src/frontend/src/service_config.rs b/src/frontend/src/service_config.rs index bcd1097cfcf..2c761c4242a 100644 --- a/src/frontend/src/service_config.rs +++ b/src/frontend/src/service_config.rs @@ -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; diff --git a/src/frontend/src/service_config/pending_rows_batcher.rs b/src/frontend/src/service_config/pending_rows_batcher.rs index ef7ce14ef02..61572015d7c 100644 --- a/src/frontend/src/service_config/pending_rows_batcher.rs +++ b/src/frontend/src/service_config/pending_rows_batcher.rs @@ -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, +} + +/// 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, /// 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, D::Error> +where + D: serde::Deserializer<'de>, +{ + let options = Option::::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::(&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::(&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::(&format!("protocols = ['{invalid}']")) - .is_err() + toml::from_str::(&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::("flow_notification_queue_capacity = 0") - .is_err() - ); + assert!(toml::from_str::("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::(&serialized).unwrap() + toml::from_str::(&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(); diff --git a/src/frontend/src/service_config/prom_store.rs b/src/frontend/src/service_config/prom_store.rs index f8429a86925..2fdd05043db 100644 --- a/src/frontend/src/service_config/prom_store.rs +++ b/src/frontend/src/service_config/prom_store.rs @@ -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 { diff --git a/src/operator/src/batcher.rs b/src/operator/src/batcher.rs index 4c5e9892c0b..677524d53f3 100644 --- a/src/operator/src/batcher.rs +++ b/src/operator/src/batcher.rs @@ -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>; - /// 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, diff --git a/src/operator/src/insert.rs b/src/operator/src/insert.rs index afc903621e2..49215f89713 100644 --- a/src/operator/src/insert.rs +++ b/src/operator/src/insert.rs @@ -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 { + 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 { + 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::>(); + 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::>(); 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 = Arc::new(UnexpectedBatcher); diff --git a/src/servers/src/batcher.rs b/src/servers/src/batcher.rs index 10f8f558eb3..6b6ca8f7ee7 100644 --- a/src/servers/src/batcher.rs +++ b/src/servers/src/batcher.rs @@ -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::().ok()) + .unwrap_or(true) +} diff --git a/src/servers/src/batcher/logical_table.rs b/src/servers/src/batcher/logical_table.rs index f4c764ee8a8..bd7a26da613 100644 --- a/src/servers/src/batcher/logical_table.rs +++ b/src/servers/src/batcher/logical_table.rs @@ -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::().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 { + 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( + &self, + requests: RowInsertRequests, + ctx: QueryContextRef, + after_prepare: impl FnOnce(RowInsertRequests) -> F, + ) -> Result<(u64, T)> + where + F: Future> + 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)) } } } diff --git a/src/servers/src/batcher/logical_table/tables.rs b/src/servers/src/batcher/logical_table/tables.rs index 4ed07f73a6e..87110366a8c 100644 --- a/src/servers/src/batcher/logical_table/tables.rs +++ b/src/servers/src/batcher/logical_table/tables.rs @@ -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> { + pub(in crate::batcher::logical_table) fn collect_unique_table_schemas<'a>( + table_rows: &[(&'a str, &'a Rows)], + ) -> Result> { 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, u32)>, ) -> Result> { 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()); diff --git a/src/servers/src/batcher/table.rs b/src/servers/src/batcher/table.rs index 1a46a7f10f3..67b3431a810 100644 --- a/src/servers/src/batcher/table.rs +++ b/src/servers/src/batcher/table.rs @@ -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, 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![ diff --git a/src/servers/src/batcher/table/batch.rs b/src/servers/src/batcher/table/batch.rs index 21ebb94a33a..e9ab339a6ac 100644 --- a/src/servers/src/batcher/table/batch.rs +++ b/src/servers/src/batcher/table/batch.rs @@ -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))); diff --git a/src/servers/src/http.rs b/src/servers/src/http.rs index 7f2688d30cf..4d168111ade 100644 --- a/src/servers/src/http.rs +++ b/src/servers/src/http.rs @@ -215,6 +215,7 @@ pub struct HttpServer { // server configs options: HttpOptions, batching_protocols: Vec, + logical_batching_protocols: Vec, bind_addr: Option, /// 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); + /// Sets a local-only write selector after authentication creates the context. async fn set_http_write_batching( State(protocol): State, @@ -245,8 +249,13 @@ async fn set_http_write_batching( .extensions() .get::>>() .is_some_and(|protocols| protocols.contains(&protocol)); + let logical_enabled = req + .extensions() + .get::() + .is_some_and(|protocols| protocols.0.contains(&protocol)); if let Some(ctx) = req.extensions_mut().get_mut::() { 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, + logical_batching_protocols: Vec, user_provider: Option, 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) -> 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) -> 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| 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!( diff --git a/src/session/src/context.rs b/src/session/src/context.rs index 0bb3f0f6a31..5ebf8b36d21 100644 --- a/src/session/src/context.rs +++ b/src/session/src/context.rs @@ -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(); diff --git a/src/standalone/src/options.rs b/src/standalone/src/options.rs index 80ba4cf11a5..d0de7408789 100644 --- a/src/standalone/src/options.rs +++ b/src/standalone/src/options.rs @@ -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::( + "[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); diff --git a/tests-integration/src/otlp.rs b/tests-integration/src/otlp.rs index dacab7ad6bd..338b03ff68f 100644 --- a/tests-integration/src/otlp.rs +++ b/tests-integration/src/otlp.rs @@ -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::() == 2 + { + break batches; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .unwrap(); + assert_eq!( + batches.iter().map(|batch| batch.num_rows()).sum::(), + 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::(), + 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") diff --git a/tests-integration/src/standalone.rs b/tests-integration/src/standalone.rs index 307a7f6daa2..38f546e9ccb 100644 --- a/tests-integration/src/standalone.rs +++ b/tests-integration/src/standalone.rs @@ -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, } 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, diff --git a/tests-integration/tests/http.rs b/tests-integration/tests/http.rs index e4e30baeba8..020e9c023c6 100644 --- a/tests-integration/tests/http.rs +++ b/tests-integration/tests/http.rs @@ -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