From 953d01ac546cf13937b8df8b6770fe9395dd11de Mon Sep 17 00:00:00 2001 From: Weny Xu Date: Wed, 23 Sep 2026 04:41:13 +0000 Subject: [PATCH] feat: support pending rows batching for MySQL and PostgreSQL (#9302) * feat: support pending rows batching for MySQL and PostgreSQL Signed-off-by: WenyXu * style: group batcher imports before item definitions Signed-off-by: WenyXu * fix: use 65536 as the default batcher worker channel capacity Signed-off-by: WenyXu * test: use a distinct custom worker channel capacity Signed-off-by: WenyXu * test: complete Prom config in worker capacity override case Signed-off-by: WenyXu --------- Signed-off-by: WenyXu --- config/config.md | 4 +- config/frontend.example.toml | 6 +- config/standalone.example.toml | 6 +- src/frontend/src/frontend.rs | 11 +- src/frontend/src/instance/builder.rs | 2 +- src/frontend/src/instance/logical_batcher.rs | 2 +- src/frontend/src/server.rs | 62 ++++-- .../service_config/pending_rows_batcher.rs | 44 ++++- src/frontend/src/service_config/prom_store.rs | 6 +- src/servers/src/batcher.rs | 39 ++++ src/servers/src/http.rs | 35 +--- src/servers/src/mysql/handler.rs | 31 ++- src/servers/src/mysql/server.rs | 11 +- src/servers/src/postgres.rs | 9 +- src/servers/src/postgres/handler.rs | 12 +- src/servers/src/postgres/server.rs | 12 +- src/servers/tests/http/prom_store_test.rs | 3 +- src/session/src/context.rs | 2 +- src/standalone/src/options.rs | 11 +- tests-integration/src/otlp.rs | 3 +- tests-integration/src/standalone.rs | 11 +- tests-integration/tests/sql.rs | 181 ++++++++++++++++++ 22 files changed, 406 insertions(+), 97 deletions(-) diff --git a/config/config.md b/config/config.md index d736c24d848..c33918a8573 100644 --- a/config/config.md +++ b/config/config.md @@ -75,7 +75,7 @@ | `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` | -- | -- | 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` | -- | -- | Ordinary-table batching for opted-in 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.
MySQL and PostgreSQL use the same acknowledgement policy. Protocol timeouts are unchanged.
Single-connection writes and INSERT SELECT may incur additional flush waits. | | `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. | @@ -343,7 +343,7 @@ | `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` | -- | -- | 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` | -- | -- | Ordinary-table batching for opted-in 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.
MySQL and PostgreSQL use the same acknowledgement policy. Protocol timeouts are unchanged.
Single-connection writes and INSERT SELECT may incur additional flush waits. | | `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. | diff --git a/config/frontend.example.toml b/config/frontend.example.toml index 275db3105c5..055acbb93bd 100644 --- a/config/frontend.example.toml +++ b/config/frontend.example.toml @@ -231,11 +231,13 @@ enable = true ## Available values: "last_non_null", "last_row". default_merge_mode = "last_non_null" -## Ordinary-table batching for opted-in HTTP ingestion protocols. +## Ordinary-table batching for opted-in 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. +## MySQL and PostgreSQL use the same acknowledgement policy. Protocol timeouts are unchanged. +## Single-connection writes and INSERT SELECT may incur additional flush waits. [pending_rows_batcher] # protocols = [ # "influxdb", @@ -246,6 +248,8 @@ default_merge_mode = "last_non_null" # "splunk", # "elasticsearch", # "http_sql", +# "mysql", +# "postgres", # "prom", # ] ## Flush interval measured from the first pending submission. Zero disables batching. diff --git a/config/standalone.example.toml b/config/standalone.example.toml index 257f48eda2c..357ec5904f7 100644 --- a/config/standalone.example.toml +++ b/config/standalone.example.toml @@ -210,11 +210,13 @@ enable = true ## Available values: "last_non_null", "last_row". default_merge_mode = "last_non_null" -## Ordinary-table batching for opted-in HTTP ingestion protocols. +## Ordinary-table batching for opted-in 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. +## MySQL and PostgreSQL use the same acknowledgement policy. Protocol timeouts are unchanged. +## Single-connection writes and INSERT SELECT may incur additional flush waits. [pending_rows_batcher] # protocols = [ # "influxdb", @@ -225,6 +227,8 @@ default_merge_mode = "last_non_null" # "splunk", # "elasticsearch", # "http_sql", +# "mysql", +# "postgres", # "prom", # ] ## Flush interval measured from the first pending submission. Zero disables batching. diff --git a/src/frontend/src/frontend.rs b/src/frontend/src/frontend.rs index 72bf322b1f1..51e291c7906 100644 --- a/src/frontend/src/frontend.rs +++ b/src/frontend/src/frontend.rs @@ -206,6 +206,7 @@ mod tests { use futures::Stream; use meta_client::MetaClientRef; use meta_client::client::MetaClientBuilder; + use servers::batcher::BatchingProtocol; use servers::grpc::{FlightCompression, GRPC_SERVER}; use servers::http::HTTP_SERVER; use servers::http::result::greptime_result_v1::GreptimedbV1Response; @@ -249,10 +250,7 @@ mod tests { .logical_table .unwrap() .protocols, - vec![ - servers::http::BatchingProtocol::Prom, - servers::http::BatchingProtocol::Otlp - ] + vec![BatchingProtocol::Prom, BatchingProtocol::Otlp] ); }, ); @@ -270,10 +268,7 @@ mod tests { FrontendOptions::load_layered_options(None, "FRONTEND_BATCHER_TEST").unwrap(); assert_eq!( options.pending_rows_batcher.table.protocols, - vec![ - servers::http::BatchingProtocol::Influxdb, - servers::http::BatchingProtocol::HttpSql - ] + vec![BatchingProtocol::Influxdb, BatchingProtocol::HttpSql] ); }, ); diff --git a/src/frontend/src/instance/builder.rs b/src/frontend/src/instance/builder.rs index d9b86ac5a05..d261daed48e 100644 --- a/src/frontend/src/instance/builder.rs +++ b/src/frontend/src/instance/builder.rs @@ -50,8 +50,8 @@ use partition::manager::PartitionRuleManager; use pipeline::pipeline_operator::PipelineOperator; use query::QueryEngineFactory; use query::region_query::RegionQueryHandlerFactoryRef; +use servers::batcher::BatchingProtocol; use servers::batcher::table::TablePendingRowsBatcher; -use servers::http::BatchingProtocol; use snafu::{OptionExt, ResultExt}; use crate::error::{self, DataFusionSnafu, ExternalSnafu, Result}; diff --git a/src/frontend/src/instance/logical_batcher.rs b/src/frontend/src/instance/logical_batcher.rs index 28ed19b3e14..e68dcf644d4 100644 --- a/src/frontend/src/instance/logical_batcher.rs +++ b/src/frontend/src/instance/logical_batcher.rs @@ -16,9 +16,9 @@ use std::sync::{Arc, Weak}; use api::v1::ColumnSchema; use async_trait::async_trait; +use servers::batcher::BatchingProtocol; use servers::batcher::logical_table::{LogicalTablePendingRowsBatcher, PendingRowsSchemaAlterer}; use servers::error::{BatcherChannelClosedSnafu, Result}; -use servers::http::BatchingProtocol; use session::context::QueryContextRef; use snafu::OptionExt; diff --git a/src/frontend/src/server.rs b/src/frontend/src/server.rs index 0a6612ba31c..fb1ab2f2603 100644 --- a/src/frontend/src/server.rs +++ b/src/frontend/src/server.rs @@ -24,7 +24,7 @@ use common_base::Plugins; use common_config::Configurable; use common_telemetry::{info, warn}; use meta_client::MetaClientOptions; -use servers::batcher::pending_rows_batch_sync_enabled; +use servers::batcher::{BatchingProtocol, pending_rows_batch_sync_enabled}; use servers::error::Error as ServerError; use servers::grpc::builder::GrpcServerBuilder; use servers::grpc::flight::FlightCraftRef; @@ -34,7 +34,7 @@ use servers::grpc::{GrpcOptions, GrpcServer}; use servers::http::event::LogValidatorRef; use servers::http::result::error_result::ErrorResponse; use servers::http::utils::router::RouterConfigurator; -use servers::http::{BatchingProtocol, HttpOptions, HttpServer, HttpServerBuilder}; +use servers::http::{HttpOptions, HttpServer, HttpServerBuilder}; use servers::interceptor::LogIngestInterceptorRef; use servers::metrics_handler::MetricsHandler; use servers::mysql::server::{MysqlServer, MysqlSpawnConfig, MysqlSpawnRef}; @@ -369,6 +369,15 @@ where } } + let table_batcher = opts.table_batcher_options(); + let batching_enabled = table_batcher.pending_rows_batching_enabled(); + let mysql_batching = + batching_enabled && table_batcher.protocols.contains(&BatchingProtocol::Mysql); + let postgres_batching = batching_enabled + && table_batcher + .protocols + .contains(&BatchingProtocol::Postgres); + if opts.mysql.enable { // Init MySQL server let opts = &opts.mysql; @@ -384,13 +393,16 @@ where let mysql_server = MysqlServer::create_server( common_runtime::global_runtime(), Arc::new(MysqlSpawnRef::new(instance.clone(), user_provider.clone())), - Arc::new(MysqlSpawnConfig::new( - opts.tls.should_force_tls(), - tls_server_config, - opts.keep_alive.as_secs(), - opts.reject_no_database.unwrap_or(false), - opts.prepared_stmt_cache_size, - )), + Arc::new( + MysqlSpawnConfig::new( + opts.tls.should_force_tls(), + tls_server_config, + opts.keep_alive.as_secs(), + opts.reject_no_database.unwrap_or(false), + opts.prepared_stmt_cache_size, + ) + .with_batching_enabled(mysql_batching), + ), Some(instance.process_manager().clone()), ); handlers.insert((mysql_server, mysql_addr)); @@ -407,15 +419,18 @@ where maybe_watch_server_tls_config(tls_server_config.clone()).context(StartServerSnafu)?; - let pg_server = Box::new(PostgresServer::new( - instance.clone(), - opts.tls.should_force_tls(), - tls_server_config, - opts.keep_alive.as_secs(), - common_runtime::global_runtime(), - user_provider.clone(), - Some(self.instance.process_manager().clone()), - )) as Box; + let pg_server = Box::new( + PostgresServer::new( + instance.clone(), + opts.tls.should_force_tls(), + tls_server_config, + opts.keep_alive.as_secs(), + common_runtime::global_runtime(), + user_provider.clone(), + Some(self.instance.process_manager().clone()), + ) + .with_batching_enabled(postgres_batching), + ) as Box; handlers.insert((pg_server, pg_addr)); } @@ -452,8 +467,11 @@ fn effective_http_options_with_sync(opts: &FrontendOptions, batch_sync: bool) -> 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) + !matches!( + protocol, + BatchingProtocol::Mysql | BatchingProtocol::Postgres + ) && (*protocol != BatchingProtocol::Prom + || (prom_store.enable && !prom_store.with_metric_engine)) }); let common_interval = common_enabled.then_some(shared.pending_rows_flush_interval); let prom_interval = (prom_store.pending_rows_batching_enabled() && batch_sync) @@ -655,6 +673,10 @@ mod tests { expected_secs, ) in [ (vec![BatchingProtocol::Prom], false, false, 5, 2, 1, 1), + (vec![BatchingProtocol::Mysql], false, false, 5, 0, 1, 1), + (vec![BatchingProtocol::Postgres], false, false, 5, 0, 1, 1), + (vec![BatchingProtocol::Mysql], false, true, 5, 0, 1, 1), + (vec![BatchingProtocol::Postgres], false, true, 5, 0, 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, 1), diff --git a/src/frontend/src/service_config/pending_rows_batcher.rs b/src/frontend/src/service_config/pending_rows_batcher.rs index 61572015d7c..b671ee9a272 100644 --- a/src/frontend/src/service_config/pending_rows_batcher.rs +++ b/src/frontend/src/service_config/pending_rows_batcher.rs @@ -17,7 +17,7 @@ use std::time::Duration; use common_batcher::flush_policy::timing::TimingFlushPolicy; use serde::{Deserialize, Serialize}; -use servers::http::BatchingProtocol; +use servers::batcher::BatchingProtocol; use tokio::sync::Semaphore; use crate::frontend::FrontendOptions; @@ -38,11 +38,11 @@ pub struct PendingRowsBatcherOptions { pub logical_table: Option, } -/// Write batching controls shared by HTTP ingestion protocols. +/// Write batching controls shared by ingestion protocols. #[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)] #[serde(default)] pub struct BatcherOptions { - /// HTTP write protocols sharing this batcher; empty disables all entrances. + /// 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. #[serde(with = "humantime_serde")] @@ -238,6 +238,8 @@ mod tests { "opentsdb", "elasticsearch", "splunk", + "mysql", + "postgres", ] { assert!( toml::from_str::(&format!( @@ -251,9 +253,11 @@ mod tests { #[test] fn test_protocols() { let options: BatcherOptions = toml::from_str( - "protocols = ['influxdb', 'opentsdb', 'otlp', 'logs', 'loki', 'splunk', 'elasticsearch', 'http_sql', 'prom']", + "protocols = ['influxdb', 'opentsdb', 'otlp', 'logs', 'loki', 'splunk', 'elasticsearch', 'http_sql', 'prom', 'mysql', 'postgres']", ).unwrap(); - assert_eq!(options.protocols.len(), 9); + assert_eq!(options.protocols.len(), 11); + assert!(options.protocols.contains(&BatchingProtocol::Mysql)); + assert!(options.protocols.contains(&BatchingProtocol::Postgres)); assert!(options.protocols.contains(&BatchingProtocol::HttpSql)); assert!(BatcherOptions::default().protocols.is_empty()); for invalid in ["sql", "jaeger", "unknown"] { @@ -289,6 +293,36 @@ mod tests { ); } + #[test] + fn test_worker_capacity_override() { + let options: FrontendOptions = toml::from_str( + r#" +[pending_rows_batcher] +worker_channel_capacity = 12345 +[pending_rows_batcher.logical_table] +worker_channel_capacity = 12345 +[prom_store] +enable = true +with_metric_engine = true +worker_channel_capacity = 12345 +"#, + ) + .unwrap(); + assert_eq!( + options.pending_rows_batcher.table.worker_channel_capacity, + 12_345 + ); + assert_eq!( + options + .pending_rows_batcher + .logical_table + .unwrap() + .worker_channel_capacity, + 12_345 + ); + assert_eq!(options.prom_store.worker_channel_capacity, 12_345); + } + #[test] fn test_partial_options_and_zero_controls() { let options: BatcherOptions = diff --git a/src/frontend/src/service_config/prom_store.rs b/src/frontend/src/service_config/prom_store.rs index 2fdd05043db..3d00f48cdda 100644 --- a/src/frontend/src/service_config/prom_store.rs +++ b/src/frontend/src/service_config/prom_store.rs @@ -108,7 +108,6 @@ mod tests { use crate::service_config::prom_store::{ default_flow_notification_queue_capacity, default_max_batch_rows, default_max_concurrent_flushes, default_max_inflight_requests, - default_worker_channel_capacity, }; #[test] @@ -160,10 +159,7 @@ mod tests { default.max_concurrent_flushes, default_max_concurrent_flushes() ); - assert_eq!( - default.worker_channel_capacity, - default_worker_channel_capacity() - ); + assert_eq!(default.worker_channel_capacity, 65_536); assert_eq!( default.max_inflight_requests, default_max_inflight_requests() diff --git a/src/servers/src/batcher.rs b/src/servers/src/batcher.rs index 6b6ca8f7ee7..3cd9a831aaa 100644 --- a/src/servers/src/batcher.rs +++ b/src/servers/src/batcher.rs @@ -22,6 +22,8 @@ pub mod table; #[cfg(test)] mod test_util; +use serde::{Deserialize, Serialize}; + /// Controls whether batching waits for storage before replying to the client. const PENDING_ROWS_BATCH_SYNC_ENV: &str = "PENDING_ROWS_BATCH_SYNC"; @@ -39,3 +41,40 @@ pub fn pending_rows_batch_sync_enabled() -> bool { .and_then(|v| v.parse::().ok()) .unwrap_or(true) } + +/// Ingestion protocols that can opt into pending-row batching. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum BatchingProtocol { + Prom, + Influxdb, + Opentsdb, + Otlp, + Logs, + Loki, + Splunk, + Elasticsearch, + HttpSql, + Mysql, + Postgres, +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_protocol_names_reject_unknown_values() { + assert_eq!( + serde_json::from_str::("\"prom\"").unwrap(), + BatchingProtocol::Prom + ); + for name in ["sql", "unknown"] { + assert!(serde_json::from_str::(&format!("\"{name}\"")).is_err()); + } + assert_eq!( + serde_json::from_str::("\"http_sql\"").unwrap(), + BatchingProtocol::HttpSql + ); + } +} diff --git a/src/servers/src/http.rs b/src/servers/src/http.rs index 4d168111ade..f947ec1621b 100644 --- a/src/servers/src/http.rs +++ b/src/servers/src/http.rs @@ -53,6 +53,7 @@ use tower_http::trace::TraceLayer; use self::authorize::AuthState; use self::result::table_result::TableResponse; +use crate::batcher::BatchingProtocol; use crate::batcher::logical_table::LogicalTablePendingRowsBatcher; use crate::elasticsearch; use crate::error::{ @@ -189,22 +190,6 @@ pub(crate) enum HttpServerKind { Api, } -/// HTTP write protocols eligible for the shared pending-row batcher. -/// Prometheus uses this selector only when metric-engine storage is disabled. -#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] -#[serde(rename_all = "snake_case")] -pub enum BatchingProtocol { - Prom, - Influxdb, - Opentsdb, - Otlp, - Logs, - Loki, - Splunk, - Elasticsearch, - HttpSql, -} - #[derive(Default)] pub struct HttpServer { router: StdMutex, @@ -2478,9 +2463,10 @@ mod batching_tests { use common_query::Output; use session::context::QueryContextRef; + use crate::batcher::BatchingProtocol; use crate::error::Result as ServerResult; use crate::http::test_helpers::TestClient; - use crate::http::{BatchingProtocol, HttpOptions, HttpServerBuilder}; + use crate::http::{HttpOptions, HttpServerBuilder}; use crate::influxdb::InfluxdbRequest; use crate::opentsdb::codec::DataPoint; use crate::query_handler::{InfluxdbLineProtocolHandler, OpentsdbProtocolHandler}; @@ -2523,21 +2509,6 @@ mod batching_tests { } } - #[test] - fn test_protocol_names_reject_unknown_values() { - assert_eq!( - serde_json::from_str::("\"prom\"").unwrap(), - BatchingProtocol::Prom - ); - for name in ["sql", "unknown"] { - assert!(serde_json::from_str::(&format!("\"{name}\"")).is_err()); - } - assert_eq!( - serde_json::from_str::("\"http_sql\"").unwrap(), - BatchingProtocol::HttpSql - ); - } - #[derive(Default)] struct RecordingWriteHandler { selections: Mutex>, diff --git a/src/servers/src/mysql/handler.rs b/src/servers/src/mysql/handler.rs index 3764591beaa..d0dd1f25279 100644 --- a/src/servers/src/mysql/handler.rs +++ b/src/servers/src/mysql/handler.rs @@ -85,6 +85,7 @@ pub struct MysqlInstanceShim { prepared_stmts_counter: AtomicU32, process_id: u32, prepared_stmt_cache_size: usize, + batching_enabled: bool, } impl MysqlInstanceShim { @@ -122,9 +123,22 @@ impl MysqlInstanceShim { prepared_stmts_counter: AtomicU32::new(1), process_id, prepared_stmt_cache_size, + batching_enabled: false, } } + /// Enables ordinary-table batching for this connection. + pub fn with_batching_enabled(mut self, enabled: bool) -> Self { + self.batching_enabled = enabled; + self + } + + fn new_query_context(&self) -> QueryContextRef { + let mut ctx = self.session.new_query_context(); + Arc::make_mut(&mut ctx).set_batching_enabled(self.batching_enabled); + ctx + } + #[tracing::instrument(skip_all, name = "mysql::do_query")] async fn do_query(&self, query: &str, query_ctx: QueryContextRef) -> Vec> { if let Some(output) = @@ -495,7 +509,7 @@ impl AsyncMysqlShim for MysqlInstanceShi raw_query: &'a str, w: StatementMetaWriter<'a, W>, ) -> Result<()> { - let query_ctx = self.session.new_query_context(); + let query_ctx = self.new_query_context(); let stmt_id = self.prepared_stmts_counter.fetch_add(1, Ordering::Relaxed); let stmt_key = uuid::Uuid::from_u128(stmt_id as u128).to_string(); let (params, columns) = match self @@ -525,7 +539,7 @@ impl AsyncMysqlShim for MysqlInstanceShi ) -> Result<()> { self.session.clear_warnings(); - let query_ctx = self.session.new_query_context(); + let query_ctx = self.new_query_context(); let db = query_ctx.get_db_string(); let _timer = crate::metrics::METRIC_MYSQL_QUERY_TIMER .with_label_values(&[crate::metrics::METRIC_MYSQL_BINQUERY, db.as_str()]) @@ -569,7 +583,7 @@ impl AsyncMysqlShim for MysqlInstanceShi query: &'a str, writer: QueryResultWriter<'a, W>, ) -> Result<()> { - let query_ctx = self.session.new_query_context(); + let query_ctx = self.new_query_context(); let db = query_ctx.get_db_string(); let _timer = crate::metrics::METRIC_MYSQL_QUERY_TIMER .with_label_values(&[crate::metrics::METRIC_MYSQL_TEXTQUERY, db.as_str()]) @@ -1039,6 +1053,17 @@ mod tests { ) } + #[test] + fn test_batching_context() { + for enabled in [false, true] { + let shim = create_shim().with_batching_enabled(enabled); + let ctx = shim.new_query_context(); + assert_eq!(ctx.batching_enabled(), enabled); + assert!(!ctx.logical_batching_enabled()); + assert_eq!(ctx.channel(), Channel::Mysql); + } + } + fn statement_with_transformed_placeholders(query: &str) -> Statement { let mut statements = ParserContext::create_with_dialect(query, &MySqlDialect {}, ParseOptions::default()) diff --git a/src/servers/src/mysql/server.rs b/src/servers/src/mysql/server.rs index 0f48cb0d598..811fc9827ec 100644 --- a/src/servers/src/mysql/server.rs +++ b/src/servers/src/mysql/server.rs @@ -86,6 +86,7 @@ pub struct MysqlSpawnConfig { reject_no_database: bool, // prepared statement cache capacity prepared_stmt_cache_size: usize, + batching_enabled: bool, } impl MysqlSpawnConfig { @@ -102,9 +103,16 @@ impl MysqlSpawnConfig { keep_alive_secs, reject_no_database, prepared_stmt_cache_size, + batching_enabled: false, } } + /// Enables ordinary-table batching for connections accepted by this server. + pub fn with_batching_enabled(mut self, enabled: bool) -> Self { + self.batching_enabled = enabled; + self + } + fn tls(&self) -> Option> { self.tls.get_config() } @@ -213,7 +221,8 @@ impl MysqlServer { stream.peer_addr()?, process_id, spawn_config.prepared_stmt_cache_size, - ); + ) + .with_batching_enabled(spawn_config.batching_enabled); let (mut r, w) = stream.into_split(); let mut w = BufWriter::with_capacity(DEFAULT_RESULT_SET_WRITE_BUFFER_SIZE, w); diff --git a/src/servers/src/postgres.rs b/src/servers/src/postgres.rs index 58ef4fdd7b7..1cd5907157b 100644 --- a/src/servers/src/postgres.rs +++ b/src/servers/src/postgres.rs @@ -80,6 +80,7 @@ pub struct PostgresServerHandlerInner { force_tls: bool, param_provider: Arc, + batching_enabled: bool, session: Arc, query_parser: Arc, } @@ -114,7 +115,12 @@ impl PgWireServerHandlers for PostgresServerHandler { } impl MakePostgresServerHandler { - fn make(&self, addr: Option, process_id: u32) -> PostgresServerHandler { + fn make( + &self, + addr: Option, + process_id: u32, + batching_enabled: bool, + ) -> PostgresServerHandler { let session = Arc::new(Session::new( addr, Channel::Postgres, @@ -127,6 +133,7 @@ impl MakePostgresServerHandler { force_tls: self.force_tls, param_provider: self.param_provider.clone(), + batching_enabled, session: session.clone(), query_parser: Arc::new(DefaultQueryParser::new(self.query_handler.clone(), session)), }; diff --git a/src/servers/src/postgres/handler.rs b/src/servers/src/postgres/handler.rs index e618c55614b..7bf1edca705 100644 --- a/src/servers/src/postgres/handler.rs +++ b/src/servers/src/postgres/handler.rs @@ -59,6 +59,14 @@ use crate::postgres::utils::convert_err; use crate::postgres::{PostgresServerHandlerInner, fixtures}; use crate::query_handler::sql::ServerSqlQueryHandlerRef; +impl PostgresServerHandlerInner { + fn new_query_context(&self) -> QueryContextRef { + let mut ctx = self.session.new_query_context(); + Arc::make_mut(&mut ctx).set_batching_enabled(self.batching_enabled); + ctx + } +} + #[async_trait] impl SimpleQueryHandler for PostgresServerHandlerInner { #[tracing::instrument(skip_all, fields(protocol = "postgres"))] @@ -68,7 +76,7 @@ impl SimpleQueryHandler for PostgresServerHandlerInner { C::Error: Debug, PgWireError: From<>::Error>, { - let query_ctx = self.session.new_query_context(); + let query_ctx = self.new_query_context(); let db = query_ctx.get_db_string(); let _timer = crate::metrics::METRIC_POSTGRES_QUERY_TIMER .with_label_values(&[crate::metrics::METRIC_POSTGRES_SIMPLE_QUERY, db.as_str()]) @@ -425,7 +433,7 @@ impl ExtendedQueryHandler for PostgresServerHandlerInner { C::Error: Debug, PgWireError: From<>::Error>, { - let query_ctx = self.session.new_query_context(); + let query_ctx = self.new_query_context(); let db = query_ctx.get_db_string(); let _timer = crate::metrics::METRIC_POSTGRES_QUERY_TIMER .with_label_values(&[crate::metrics::METRIC_POSTGRES_EXTENDED_QUERY, db.as_str()]) diff --git a/src/servers/src/postgres/server.rs b/src/servers/src/postgres/server.rs index 09bc5015a0b..6e4aa942f9d 100644 --- a/src/servers/src/postgres/server.rs +++ b/src/servers/src/postgres/server.rs @@ -38,6 +38,7 @@ pub struct PostgresServer { make_handler: Arc, tls_server_config: Arc, keep_alive_secs: u64, + batching_enabled: bool, bind_addr: Option, process_manager: Option, } @@ -66,17 +67,25 @@ impl PostgresServer { make_handler, tls_server_config, keep_alive_secs, + batching_enabled: false, bind_addr: None, process_manager, } } + /// Enables ordinary-table batching for connections accepted by this server. + pub fn with_batching_enabled(mut self, enabled: bool) -> Self { + self.batching_enabled = enabled; + self + } + fn accept( &self, io_runtime: Runtime, accepting_stream: AbortableStream, ) -> impl Future + use<> { let handler_maker = self.make_handler.clone(); + let batching_enabled = self.batching_enabled; let tls_server_config = self.tls_server_config.clone(); let process_manager = self.process_manager.clone(); accepting_stream.for_each(move |tcp_stream| { @@ -102,7 +111,8 @@ impl PostgresServer { let _handle = io_runtime.spawn(async move { crate::metrics::METRIC_POSTGRES_CONNECTIONS.inc(); - let pg_handler = Arc::new(handler_maker.make(addr, process_id)); + let pg_handler = + Arc::new(handler_maker.make(addr, process_id, batching_enabled)); let r = process_socket(io_stream, tls_acceptor.clone(), pg_handler).await; crate::metrics::METRIC_POSTGRES_CONNECTIONS.dec(); diff --git a/src/servers/tests/http/prom_store_test.rs b/src/servers/tests/http/prom_store_test.rs index 7c07e3b0344..ffd2c9dcf7f 100644 --- a/src/servers/tests/http/prom_store_test.rs +++ b/src/servers/tests/http/prom_store_test.rs @@ -36,11 +36,12 @@ use datafusion_expr::LogicalPlan; use prost::Message; use query::parser::PromQuery; use query::query_engine::DescribeResult; +use servers::batcher::BatchingProtocol; use servers::error::{self, Result}; use servers::http::header::{CONTENT_ENCODING_SNAPPY, CONTENT_TYPE_PROTOBUF}; use servers::http::prom_store::PHYSICAL_TABLE_PARAM; use servers::http::test_helpers::{TestClient, TestResponse}; -use servers::http::{BatchingProtocol, HttpOptions, HttpServerBuilder}; +use servers::http::{HttpOptions, HttpServerBuilder}; use servers::prom_remote_write::v2::test_util as remote_write_v2; use servers::prom_remote_write::validation::PromValidationMode; use servers::prom_store; diff --git a/src/session/src/context.rs b/src/session/src/context.rs index 5ebf8b36d21..ec10dc6cf3a 100644 --- a/src/session/src/context.rs +++ b/src/session/src/context.rs @@ -471,7 +471,7 @@ impl QueryContext { self.logical_batching_enabled = enabled; } - /// Whether the local HTTP entry point selected ordinary-table batching. + /// Whether the local protocol entry point selected ordinary-table batching. pub fn batching_enabled(&self) -> bool { self.batching_enabled } diff --git a/src/standalone/src/options.rs b/src/standalone/src/options.rs index d0de7408789..fbfcc2ccc1d 100644 --- a/src/standalone/src/options.rs +++ b/src/standalone/src/options.rs @@ -217,6 +217,7 @@ mod tests { use std::sync::Arc; use common_event_recorder::EventTypeFilter; + use servers::batcher::BatchingProtocol; use crate::options::*; @@ -266,10 +267,7 @@ mod tests { .logical_table .unwrap() .protocols, - vec![ - servers::http::BatchingProtocol::Prom, - servers::http::BatchingProtocol::Otlp - ] + vec![BatchingProtocol::Prom, BatchingProtocol::Otlp] ); }, ); @@ -288,10 +286,7 @@ mod tests { .unwrap(); assert_eq!( options.pending_rows_batcher.table.protocols, - vec![ - servers::http::BatchingProtocol::Influxdb, - servers::http::BatchingProtocol::HttpSql - ] + vec![BatchingProtocol::Influxdb, BatchingProtocol::HttpSql] ); }, ); diff --git a/tests-integration/src/otlp.rs b/tests-integration/src/otlp.rs index 338b03ff68f..9fbf80fbeea 100644 --- a/tests-integration/src/otlp.rs +++ b/tests-integration/src/otlp.rs @@ -514,8 +514,7 @@ WITH( ExponentialHistogram, ExponentialHistogramDataPoint, exponential_histogram_data_point, }; use prost::Message; - use servers::batcher::pending_rows_batch_sync_enabled; - use servers::http::BatchingProtocol; + use servers::batcher::{BatchingProtocol, pending_rows_batch_sync_enabled}; use servers::http::test_helpers::TestClient; use session::protocol_ctx::{OtlpMetricCtx, ProtocolCtx}; diff --git a/tests-integration/src/standalone.rs b/tests-integration/src/standalone.rs index 38f546e9ccb..6fc8ce231af 100644 --- a/tests-integration/src/standalone.rs +++ b/tests-integration/src/standalone.rs @@ -90,6 +90,7 @@ pub struct GreptimeDbStandaloneBuilder { auto_create_table: bool, experimental_metric_export: bool, logical_batcher: Option, + table_batcher: BatcherOptions, } impl GreptimeDbStandaloneBuilder { @@ -111,6 +112,7 @@ impl GreptimeDbStandaloneBuilder { auto_create_table: true, experimental_metric_export: false, logical_batcher: None, + table_batcher: BatcherOptions::default(), } } @@ -121,6 +123,13 @@ impl GreptimeDbStandaloneBuilder { self } + /// Configures ordinary-table batching for protocol integration tests. + #[must_use] + pub fn with_table_batcher(mut self, options: BatcherOptions) -> Self { + self.table_batcher = options; + self + } + /// Configures logical-table batching for integration tests. #[must_use] pub fn with_logical_batcher(mut self, options: BatcherOptions) -> Self { @@ -389,7 +398,7 @@ impl GreptimeDbStandaloneBuilder { experimental_metric_export: self.experimental_metric_export, pending_rows_batcher: PendingRowsBatcherOptions { logical_table: self.logical_batcher.clone(), - ..Default::default() + table: self.table_batcher.clone(), }, // Tests cover the descriptor, so they run with it enabled. otlp: frontend::service_config::OtlpOptions { diff --git a/tests-integration/tests/sql.rs b/tests-integration/tests/sql.rs index e3b66fc4447..cb3d151fa85 100644 --- a/tests-integration/tests/sql.rs +++ b/tests-integration/tests/sql.rs @@ -2672,3 +2672,184 @@ pub async fn test_declare_fetch_close_cursor(store_type: StorageType) { let _ = fe_pg_server.shutdown().await; guard.remove_all().await; } + +/// Keeps SQL and input rows fixed while changing only the selected batching protocols. +#[tokio::test(flavor = "multi_thread")] +async fn test_sql_batcher_alignment() { + use common_telemetry::{dump_metrics, info, init_default_ut_logging}; + use frontend::server::Services; + use frontend::service_config::BatcherOptions; + use servers::batcher::BatchingProtocol; + use servers::mysql::server::MYSQL_SERVER; + use tests_integration::standalone::GreptimeDbStandaloneBuilder; + + fn flushes() -> u64 { + dump_metrics() + .unwrap() + .lines() + .find_map(|line| { + line.strip_prefix("greptime_table_batcher_flush_total ") + .map(|value| value.parse().unwrap()) + }) + .unwrap_or(0) + } + + init_default_ut_logging(); + for protocols in [ + vec![], + vec![BatchingProtocol::Mysql], + vec![BatchingProtocol::Postgres], + vec![BatchingProtocol::Mysql, BatchingProtocol::Postgres], + ] { + let mysql_enabled = protocols.contains(&BatchingProtocol::Mysql); + let pg_enabled = protocols.contains(&BatchingProtocol::Postgres); + let mut instance = GreptimeDbStandaloneBuilder::new("sql_batcher_alignment") + .with_table_batcher(BatcherOptions { + protocols, + pending_rows_flush_interval: Duration::from_millis(10), + ..Default::default() + }) + .build() + .await; + let mut opts = instance.opts.clone(); + opts.http.addr = "127.0.0.1:0".into(); + opts.grpc.bind_addr = "127.0.0.1:0".into(); + opts.mysql.addr = "127.0.0.1:0".into(); + opts.postgres.addr = "127.0.0.1:0".into(); + let mut servers = Services::new(opts, instance.fe_instance().clone(), Default::default()) + .build() + .unwrap(); + servers.start_all().await.unwrap(); + let mysql_url = format!("mysql://{}/public", servers.addr(MYSQL_SERVER).unwrap()); + let mut mysql = MySqlConnection::connect(&mysql_url).await.unwrap(); + let (pg, connection) = tokio_postgres::connect( + &format!( + "postgres://{}/public", + servers.addr("POSTGRES_SERVER").unwrap() + ), + NoTls, + ) + .await + .unwrap(); + let pg_task = tokio::spawn(async move { connection.await.unwrap() }); + + for round in 0..2 { + mysql.execute("CREATE TABLE batch_sql (ts TIMESTAMP TIME INDEX, host STRING PRIMARY KEY, val INT DEFAULT 7)").await.unwrap(); + let started = std::time::Instant::now(); + let before = flushes(); + // COM_QUERY and COM_STMT_EXECUTE must both opt in. + assert_eq!( + mysql + .execute("INSERT INTO batch_sql (ts, host) VALUES (0, 'a')") + .await + .unwrap() + .rows_affected(), + 1 + ); + assert_eq!( + sqlx::query("INSERT INTO batch_sql VALUES (0, ?, ?)") + .persistent(false) + .bind("b") + .bind(8_i32) + .execute(&mut mysql) + .await + .unwrap() + .rows_affected(), + 1 + ); + assert_eq!(flushes() - before, if mysql_enabled { 2 } else { 0 }); + + let before = flushes(); + let result = pg + .simple_query("INSERT INTO batch_sql (ts, host) VALUES (0, 'c')") + .await + .unwrap(); + assert!(matches!( + result.as_slice(), + [SimpleQueryMessage::CommandComplete(1)] + )); + let statement = pg + .prepare("INSERT INTO batch_sql VALUES (0, $1, $2)") + .await + .unwrap(); + assert_eq!(pg.execute(&statement, &[&"d", &9_i32]).await.unwrap(), 1); + assert_eq!(flushes() - before, if pg_enabled { 2 } else { 0 }); + + // Independent connections can submit to the same worker concurrently. + let (mysql_result, pg_result) = tokio::join!( + mysql.execute("INSERT INTO batch_sql VALUES (0, 'e', 10)"), + pg.simple_query("INSERT INTO batch_sql VALUES (0, 'f', 11)") + ); + assert_eq!(mysql_result.unwrap().rows_affected(), 1); + assert!(matches!( + pg_result.unwrap().as_slice(), + [SimpleQueryMessage::CommandComplete(1)] + )); + + // Errors must reach the client without affecting later requests. + assert!( + mysql + .execute("INSERT INTO batch_sql (missing) VALUES (1)") + .await + .is_err() + ); + assert!( + pg.simple_query("INSERT INTO batch_sql (missing) VALUES (1)") + .await + .is_err() + ); + let expected = vec![ + ("a".to_string(), 7_i32), + ("b".to_string(), 8), + ("c".to_string(), 7), + ("d".to_string(), 9), + ("e".to_string(), 10), + ("f".to_string(), 11), + ]; + // No polling: successful protocol responses guarantee visibility. + let rows: Vec<(String, i32)> = + sqlx::query_as("SELECT host, val FROM batch_sql ORDER BY host") + .persistent(false) + .fetch_all(&mut mysql) + .await + .unwrap(); + assert_eq!(rows, expected); + let rows: Vec<(String, i32)> = pg + .query("SELECT host, val FROM batch_sql ORDER BY host", &[]) + .await + .unwrap() + .iter() + .map(|row| (row.get(0), row.get(1))) + .collect(); + assert_eq!(rows, expected); + + mysql.execute("CREATE TABLE batch_copy (ts TIMESTAMP TIME INDEX, host STRING PRIMARY KEY, val INT)").await.unwrap(); + assert_eq!( + mysql + .execute("INSERT INTO batch_copy SELECT * FROM batch_sql") + .await + .unwrap() + .rows_affected(), + 6 + ); + let rows: Vec<(String, i32)> = + sqlx::query_as("SELECT host, val FROM batch_copy ORDER BY host") + .persistent(false) + .fetch_all(&mut mysql) + .await + .unwrap(); + assert_eq!(rows, expected); + info!( + "sql_batcher_alignment mysql={mysql_enabled} postgres={pg_enabled} round={round} elapsed={:?}", + started.elapsed() + ); + mysql.execute("DROP TABLE batch_sql").await.unwrap(); + mysql.execute("DROP TABLE batch_copy").await.unwrap(); + } + mysql.close().await.unwrap(); + drop(pg); + pg_task.await.unwrap(); + servers.shutdown_all().await.unwrap(); + instance.guard.remove_all().await; + } +}