From 13bc733531b9eb438b69b95063a4839f3e30e452 Mon Sep 17 00:00:00 2001 From: charlesdong1991 Date: Sat, 28 Feb 2026 21:52:03 +0100 Subject: [PATCH 1/3] Add setter and getter for writer_bucket_no_key_assigner --- bindings/python/fluss/__init__.pyi | 4 ++++ bindings/python/src/config.rs | 21 +++++++++++++++++++++ 2 files changed, 25 insertions(+) diff --git a/bindings/python/fluss/__init__.pyi b/bindings/python/fluss/__init__.pyi index 4c2142d7..dee1923a 100644 --- a/bindings/python/fluss/__init__.pyi +++ b/bindings/python/fluss/__init__.pyi @@ -154,6 +154,10 @@ class Config: @writer_batch_size.setter def writer_batch_size(self, size: int) -> None: ... @property + def writer_bucket_no_key_assigner(self) -> str: ... + @writer_bucket_no_key_assigner.setter + def writer_bucket_no_key_assigner(self, value: str) -> None: ... + @property def scanner_remote_log_prefetch_num(self) -> int: ... @scanner_remote_log_prefetch_num.setter def scanner_remote_log_prefetch_num(self, num: int) -> None: ... diff --git a/bindings/python/src/config.rs b/bindings/python/src/config.rs index 4582d43d..145793c0 100644 --- a/bindings/python/src/config.rs +++ b/bindings/python/src/config.rs @@ -255,6 +255,27 @@ impl Config { self.inner.writer_batch_timeout_ms = timeout; } + /// Get the bucket assignment strategy for tables without bucket keys + #[getter] + fn writer_bucket_no_key_assigner(&self) -> String { + self.inner.writer_bucket_no_key_assigner.to_string() + } + + /// Set the bucket assignment strategy for tables without bucket keys + #[setter] + fn set_writer_bucket_no_key_assigner(&mut self, value: String) -> PyResult<()> { + self.inner.writer_bucket_no_key_assigner = match value.as_str() { + "round_robin" => fcore::config::NoKeyAssigner::RoundRobin, + "sticky" => fcore::config::NoKeyAssigner::Sticky, + other => { + return Err(FlussError::new_err(format!( + "Unknown bucket assigner type: {other}, expected 'sticky' or 'round_robin'" + ))); + } + }; + Ok(()) + } + /// Get the connect timeout in milliseconds #[getter] fn connect_timeout_ms(&self) -> u64 { From 6a25a7d89e8f2cfb312924d8f4d5a020bec6c267 Mon Sep 17 00:00:00 2001 From: charlesdong1991 Date: Tue, 3 Mar 2026 20:40:43 +0100 Subject: [PATCH 2/3] Avoid code duplication --- bindings/cpp/src/lib.rs | 13 ++++--------- bindings/python/src/config.rs | 30 ++++++++++++------------------ crates/fluss/src/config.rs | 18 +++++++----------- 3 files changed, 23 insertions(+), 38 deletions(-) diff --git a/bindings/cpp/src/lib.rs b/bindings/cpp/src/lib.rs index c310fc83..0ebd375d 100644 --- a/bindings/cpp/src/lib.rs +++ b/bindings/cpp/src/lib.rs @@ -632,15 +632,10 @@ fn err_from_core_error(e: &fcore::error::Error) -> ffi::FfiResult { // Connection implementation fn new_connection(config: &ffi::FfiConfig) -> Result<*mut Connection, String> { - let assigner_type = match config.writer_bucket_no_key_assigner.as_str() { - "round_robin" => fluss::config::NoKeyAssigner::RoundRobin, - "sticky" => fluss::config::NoKeyAssigner::Sticky, - other => { - return Err(format!( - "Unknown bucket assigner type: '{other}', expected 'sticky' or 'round_robin'" - )); - } - }; + let assigner_type = config + .writer_bucket_no_key_assigner + .parse::() + .map_err(|e| format!("Invalid bucket assigner type: {e}"))?; let config_core = fluss::config::Config { bootstrap_servers: config.bootstrap_servers.to_string(), writer_request_max_size: config.writer_request_max_size, diff --git a/bindings/python/src/config.rs b/bindings/python/src/config.rs index 145793c0..f99f9c63 100644 --- a/bindings/python/src/config.rs +++ b/bindings/python/src/config.rs @@ -98,15 +98,12 @@ impl Config { })?; } "writer.bucket.no-key-assigner" => { - config.writer_bucket_no_key_assigner = match value.as_str() { - "round_robin" => fcore::config::NoKeyAssigner::RoundRobin, - "sticky" => fcore::config::NoKeyAssigner::Sticky, - other => { - return Err(FlussError::new_err(format!( - "Unknown bucket assigner type: {other}, expected 'sticky' or 'round_robin'" - ))); - } - }; + config.writer_bucket_no_key_assigner = + value.parse::().map_err(|e| { + FlussError::new_err(format!( + "Invalid value '{value}' for '{key}': {e}" + )) + })?; } "connect-timeout" => { config.connect_timeout_ms = value.parse::().map_err(|e| { @@ -264,15 +261,12 @@ impl Config { /// Set the bucket assignment strategy for tables without bucket keys #[setter] fn set_writer_bucket_no_key_assigner(&mut self, value: String) -> PyResult<()> { - self.inner.writer_bucket_no_key_assigner = match value.as_str() { - "round_robin" => fcore::config::NoKeyAssigner::RoundRobin, - "sticky" => fcore::config::NoKeyAssigner::Sticky, - other => { - return Err(FlussError::new_err(format!( - "Unknown bucket assigner type: {other}, expected 'sticky' or 'round_robin'" - ))); - } - }; + self.inner.writer_bucket_no_key_assigner = + value.parse::().map_err(|e| { + FlussError::new_err(format!( + "Invalid value '{value}' for 'writer.bucket.no-key-assigner': {e}" + )) + })?; Ok(()) } diff --git a/crates/fluss/src/config.rs b/crates/fluss/src/config.rs index 438c9483..08ffbfae 100644 --- a/crates/fluss/src/config.rs +++ b/crates/fluss/src/config.rs @@ -17,7 +17,7 @@ use clap::{Parser, ValueEnum}; use serde::{Deserialize, Serialize}; -use std::fmt; +use strum_macros::{Display, EnumString}; const DEFAULT_BOOTSTRAP_SERVER: &str = "127.0.0.1:9123"; const DEFAULT_REQUEST_MAX_SIZE: i32 = 10 * 1024 * 1024; @@ -36,24 +36,20 @@ const DEFAULT_SASL_MECHANISM: &str = "PLAIN"; /// Bucket assigner strategy for tables without bucket keys. /// Matches Java `client.writer.bucket.no-key-assigner`. -#[derive(Debug, Clone, Copy, PartialEq, Eq, ValueEnum, Deserialize, Serialize)] +#[derive( + Debug, Clone, Copy, PartialEq, Eq, ValueEnum, Deserialize, Serialize, EnumString, Display, +)] #[serde(rename_all = "snake_case")] +#[strum(ascii_case_insensitive)] pub enum NoKeyAssigner { /// Sticks to one bucket until the batch is full, then switches. + #[strum(serialize = "sticky")] Sticky, /// Assigns each record to the next bucket in a rotating sequence. + #[strum(serialize = "round_robin")] RoundRobin, } -impl fmt::Display for NoKeyAssigner { - fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - match self { - NoKeyAssigner::Sticky => write!(f, "sticky"), - NoKeyAssigner::RoundRobin => write!(f, "round_robin"), - } - } -} - #[derive(Parser, Clone, Deserialize, Serialize)] #[command(author, version, about, long_about = None)] pub struct Config { From 2449192989c7f34c99d620bb3da8b19d431ef0e7 Mon Sep 17 00:00:00 2001 From: charlesdong1991 Date: Tue, 3 Mar 2026 20:45:56 +0100 Subject: [PATCH 3/3] fix conflicts issue --- bindings/cpp/src/lib.rs | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/bindings/cpp/src/lib.rs b/bindings/cpp/src/lib.rs index 2146e1e4..36b9c516 100644 --- a/bindings/cpp/src/lib.rs +++ b/bindings/cpp/src/lib.rs @@ -648,11 +648,14 @@ fn err_ptr_from_core(e: &fcore::error::Error) -> ffi::FfiPtrResult { } // Connection implementation -fn new_connection(config: &ffi::FfiConfig) -> Result<*mut Connection, String> { - let assigner_type = config +fn new_connection(config: &ffi::FfiConfig) -> ffi::FfiPtrResult { + let assigner_type = match config .writer_bucket_no_key_assigner .parse::() - .map_err(|e| format!("Invalid bucket assigner type: {e}"))?; + { + Ok(v) => v, + Err(e) => return client_err_ptr(format!("Invalid bucket assigner type: {e}")), + }; let config_core = fluss::config::Config { bootstrap_servers: config.bootstrap_servers.to_string(), writer_request_max_size: config.writer_request_max_size,