Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 6 additions & 8 deletions bindings/cpp/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -649,14 +649,12 @@ fn err_ptr_from_core(e: &fcore::error::Error) -> ffi::FfiPtrResult {

// Connection implementation
fn new_connection(config: &ffi::FfiConfig) -> ffi::FfiPtrResult {
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 client_err_ptr(format!(
"Unknown bucket assigner type: '{other}', expected 'sticky' or 'round_robin'"
));
}
let assigner_type = match config
.writer_bucket_no_key_assigner
.parse::<fluss::config::NoKeyAssigner>()
{
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(),
Expand Down
4 changes: 4 additions & 0 deletions bindings/python/fluss/__init__.pyi
Original file line number Diff line number Diff line change
Expand Up @@ -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: ...
Expand Down
33 changes: 24 additions & 9 deletions bindings/python/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::<fcore::config::NoKeyAssigner>().map_err(|e| {
FlussError::new_err(format!(
"Invalid value '{value}' for '{key}': {e}"
))
})?;
}
"connect-timeout" => {
config.connect_timeout_ms = value.parse::<u64>().map_err(|e| {
Expand Down Expand Up @@ -255,6 +252,24 @@ 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 =
value.parse::<fcore::config::NoKeyAssigner>().map_err(|e| {
FlussError::new_err(format!(
"Invalid value '{value}' for 'writer.bucket.no-key-assigner': {e}"
))
})?;
Ok(())
}

/// Get the connect timeout in milliseconds
#[getter]
fn connect_timeout_ms(&self) -> u64 {
Expand Down
18 changes: 7 additions & 11 deletions crates/fluss/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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 {
Expand Down
Loading