feat: support pending rows batching for MySQL and PostgreSQL (#9302)

* feat: support pending rows batching for MySQL and PostgreSQL

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

* style: group batcher imports before item definitions

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

* fix: use 65536 as the default batcher worker channel capacity

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

* test: use a distinct custom worker channel capacity

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

* test: complete Prom config in worker capacity override case

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

---------

Signed-off-by: WenyXu <wenymedia@gmail.com>
This commit is contained in:
Weny Xu
2026-09-23 04:41:13 +00:00
committed by GitHub
parent 9281965bac
commit 953d01ac54
22 changed files with 406 additions and 97 deletions
+2 -2
View File
@@ -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.<br/>Available values: "last_non_null", "last_row". |
| `pending_rows_batcher` | -- | -- | Ordinary-table batching for opted-in HTTP ingestion protocols.<br/>PENDING_ROWS_BATCH_SYNC defaults to true for both batchers. Set it to false to acknowledge<br/>queue admission without waiting for storage; later failures cannot be returned to the client.<br/>Omitted or empty protocols disables batching. Prom without metric engine uses this batcher.<br/>OTLP logs, traces and ordinary metrics use this batcher. |
| `pending_rows_batcher` | -- | -- | Ordinary-table batching for opted-in ingestion protocols.<br/>PENDING_ROWS_BATCH_SYNC defaults to true for both batchers. Set it to false to acknowledge<br/>queue admission without waiting for storage; later failures cannot be returned to the client.<br/>Omitted or empty protocols disables batching. Prom without metric engine uses this batcher.<br/>OTLP logs, traces and ordinary metrics use this batcher.<br/>MySQL and PostgreSQL use the same acknowledgement policy. Protocol timeouts are unchanged.<br/>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.<br/>Available values: "last_non_null", "last_row". |
| `pending_rows_batcher` | -- | -- | Ordinary-table batching for opted-in HTTP ingestion protocols.<br/>PENDING_ROWS_BATCH_SYNC defaults to true for both batchers. Set it to false to acknowledge<br/>queue admission without waiting for storage; later failures cannot be returned to the client.<br/>Omitted or empty protocols disables batching. Prom without metric engine uses this batcher.<br/>OTLP logs, traces and ordinary metrics use this batcher. |
| `pending_rows_batcher` | -- | -- | Ordinary-table batching for opted-in ingestion protocols.<br/>PENDING_ROWS_BATCH_SYNC defaults to true for both batchers. Set it to false to acknowledge<br/>queue admission without waiting for storage; later failures cannot be returned to the client.<br/>Omitted or empty protocols disables batching. Prom without metric engine uses this batcher.<br/>OTLP logs, traces and ordinary metrics use this batcher.<br/>MySQL and PostgreSQL use the same acknowledgement policy. Protocol timeouts are unchanged.<br/>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. |
+5 -1
View File
@@ -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.
+5 -1
View File
@@ -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.
+3 -8
View File
@@ -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]
);
},
);
+1 -1
View File
@@ -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};
+1 -1
View File
@@ -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;
+42 -20
View File
@@ -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<dyn Server>;
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<dyn Server>;
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),
@@ -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<BatcherOptions>,
}
/// 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<BatchingProtocol>,
/// 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::<FrontendOptions>(&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 =
@@ -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()
+39
View File
@@ -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::<bool>().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::<BatchingProtocol>("\"prom\"").unwrap(),
BatchingProtocol::Prom
);
for name in ["sql", "unknown"] {
assert!(serde_json::from_str::<BatchingProtocol>(&format!("\"{name}\"")).is_err());
}
assert_eq!(
serde_json::from_str::<BatchingProtocol>("\"http_sql\"").unwrap(),
BatchingProtocol::HttpSql
);
}
}
+3 -32
View File
@@ -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<Router>,
@@ -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::<BatchingProtocol>("\"prom\"").unwrap(),
BatchingProtocol::Prom
);
for name in ["sql", "unknown"] {
assert!(serde_json::from_str::<BatchingProtocol>(&format!("\"{name}\"")).is_err());
}
assert_eq!(
serde_json::from_str::<BatchingProtocol>("\"http_sql\"").unwrap(),
BatchingProtocol::HttpSql
);
}
#[derive(Default)]
struct RecordingWriteHandler {
selections: Mutex<Vec<bool>>,
+28 -3
View File
@@ -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<Result<Output>> {
if let Some(output) =
@@ -495,7 +509,7 @@ impl<W: AsyncWrite + Send + Sync + Unpin> AsyncMysqlShim<W> 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<W: AsyncWrite + Send + Sync + Unpin> AsyncMysqlShim<W> 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<W: AsyncWrite + Send + Sync + Unpin> AsyncMysqlShim<W> 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())
+10 -1
View File
@@ -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<Arc<ServerConfig>> {
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);
+8 -1
View File
@@ -80,6 +80,7 @@ pub struct PostgresServerHandlerInner {
force_tls: bool,
param_provider: Arc<GreptimeDBStartupParameters>,
batching_enabled: bool,
session: Arc<Session>,
query_parser: Arc<DefaultQueryParser>,
}
@@ -114,7 +115,12 @@ impl PgWireServerHandlers for PostgresServerHandler {
}
impl MakePostgresServerHandler {
fn make(&self, addr: Option<SocketAddr>, process_id: u32) -> PostgresServerHandler {
fn make(
&self,
addr: Option<SocketAddr>,
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)),
};
+10 -2
View File
@@ -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<<C as Sink<PgWireBackendMessage>>::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<<C as Sink<PgWireBackendMessage>>::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()])
+11 -1
View File
@@ -38,6 +38,7 @@ pub struct PostgresServer {
make_handler: Arc<MakePostgresServerHandler>,
tls_server_config: Arc<ReloadableTlsServerConfig>,
keep_alive_secs: u64,
batching_enabled: bool,
bind_addr: Option<SocketAddr>,
process_manager: Option<ProcessManagerRef>,
}
@@ -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<Output = ()> + 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();
+2 -1
View File
@@ -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;
+1 -1
View File
@@ -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
}
+3 -8
View File
@@ -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]
);
},
);
+1 -2
View File
@@ -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};
+10 -1
View File
@@ -90,6 +90,7 @@ pub struct GreptimeDbStandaloneBuilder {
auto_create_table: bool,
experimental_metric_export: bool,
logical_batcher: Option<BatcherOptions>,
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 {
+181
View File
@@ -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;
}
}