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