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
7 changes: 6 additions & 1 deletion fluss-rust/bindings/cpp/include/fluss.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -188,6 +188,8 @@ struct ErrorCode {
static constexpr int INVALID_ALTER_TABLE_EXCEPTION = 56;
/// Deletion operations are disabled on this table.
static constexpr int DELETION_DISABLED_EXCEPTION = 57;
/// The server rejected a write due to storage backpressure.
static constexpr int STORAGE_BACKPRESSURE_EXCEPTION = 72;

/// Returns true if retrying the request may succeed. Mirrors Java's RetriableException hierarchy.
static constexpr bool IsRetriable(int32_t code) {
Expand All @@ -198,7 +200,8 @@ struct ErrorCode {
code == UNKNOWN_TABLE_OR_BUCKET_EXCEPTION || code == REQUEST_TIME_OUT ||
code == STORAGE_EXCEPTION ||
code == NOT_ENOUGH_REPLICAS_AFTER_APPEND_EXCEPTION ||
code == NOT_ENOUGH_REPLICAS_EXCEPTION || code == LEADER_NOT_AVAILABLE_EXCEPTION;
code == NOT_ENOUGH_REPLICAS_EXCEPTION || code == LEADER_NOT_AVAILABLE_EXCEPTION ||
code == STORAGE_BACKPRESSURE_EXCEPTION;
}
};

Expand Down Expand Up @@ -1388,6 +1391,8 @@ struct Configuration {
size_t writer_buffer_memory_size{64 * 1024 * 1024};
// Maximum time in milliseconds to block waiting for buffer memory
uint64_t writer_buffer_wait_timeout_ms{std::numeric_limits<uint64_t>::max()};
// Maximum KV backpressure throttle in milliseconds
uint64_t writer_kv_backpressure_max_throttle_ms{3000};
// Connect timeout in milliseconds for TCP transport connect
uint64_t connect_timeout_ms{120000};
// Security protocol: "PLAINTEXT" (default, no auth) or "sasl" (SASL auth)
Expand Down
2 changes: 2 additions & 0 deletions fluss-rust/bindings/cpp/src/ffi_converter.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -234,6 +234,8 @@ inline ffi::FfiConfig to_ffi_config(const Configuration& config) {
config.writer_max_inflight_requests_per_bucket;
ffi_config.writer_buffer_memory_size = config.writer_buffer_memory_size;
ffi_config.writer_buffer_wait_timeout_ms = config.writer_buffer_wait_timeout_ms;
ffi_config.writer_kv_backpressure_max_throttle_ms =
config.writer_kv_backpressure_max_throttle_ms;
ffi_config.connect_timeout_ms = config.connect_timeout_ms;
ffi_config.security_protocol = rust::String(config.security_protocol);
ffi_config.security_sasl_mechanism = rust::String(config.security_sasl_mechanism);
Expand Down
2 changes: 2 additions & 0 deletions fluss-rust/bindings/cpp/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,7 @@ mod ffi {
writer_max_inflight_requests_per_bucket: usize,
writer_buffer_memory_size: usize,
writer_buffer_wait_timeout_ms: u64,
writer_kv_backpressure_max_throttle_ms: u64,
connect_timeout_ms: u64,
security_protocol: String,
security_sasl_mechanism: String,
Expand Down Expand Up @@ -974,6 +975,7 @@ fn new_connection(config: &ffi::FfiConfig) -> ffi::FfiPtrResult {
writer_max_inflight_requests_per_bucket: config.writer_max_inflight_requests_per_bucket,
writer_buffer_memory_size: config.writer_buffer_memory_size,
writer_buffer_wait_timeout_ms: config.writer_buffer_wait_timeout_ms,
writer_kv_backpressure_max_throttle_ms: config.writer_kv_backpressure_max_throttle_ms,
connect_timeout_ms: config.connect_timeout_ms,
security_protocol: config.security_protocol.to_string(),
security_sasl_mechanism: config.security_sasl_mechanism.to_string(),
Expand Down
11 changes: 11 additions & 0 deletions fluss-rust/bindings/cpp/test/test_ffi_converter.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,17 @@ fluss::ffi::FfiTypeNode Node(int32_t type_id, uint32_t child_count = 0, bool nul

} // namespace

TEST(FfiConverterTest, KvBackpressureConfiguration) {
fluss::Configuration config;
EXPECT_EQ(config.writer_kv_backpressure_max_throttle_ms, 3000u);

config.writer_kv_backpressure_max_throttle_ms = 1500;
auto ffi_config = fluss::utils::to_ffi_config(config);
EXPECT_EQ(ffi_config.writer_kv_backpressure_max_throttle_ms, 1500u);
EXPECT_EQ(fluss::ErrorCode::STORAGE_BACKPRESSURE_EXCEPTION, 72);
EXPECT_TRUE(fluss::ErrorCode::IsRetriable(72));
}

// --- DataType value semantics ---

TEST(DataTypeTest, DefaultNullable) { EXPECT_TRUE(DataType::Int().nullable()); }
Expand Down
7 changes: 7 additions & 0 deletions fluss-rust/bindings/elixir/lib/fluss/config.ex
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@ defmodule Fluss.Config do
writer_bucket_no_key_assigner: nil,
writer_buffer_memory_size: nil,
writer_buffer_wait_timeout_ms: nil,
writer_kv_backpressure_max_throttle_ms: nil,
writer_dynamic_batch_size_enabled: nil,
writer_dynamic_batch_size_min: nil,
writer_enable_idempotence: nil,
Expand Down Expand Up @@ -80,6 +81,7 @@ defmodule Fluss.Config do
writer_bucket_no_key_assigner: :sticky | :round_robin | nil,
writer_buffer_memory_size: non_neg_integer() | nil,
writer_buffer_wait_timeout_ms: non_neg_integer() | nil,
writer_kv_backpressure_max_throttle_ms: non_neg_integer() | nil,
writer_dynamic_batch_size_enabled: boolean() | nil,
writer_dynamic_batch_size_min: non_neg_integer() | nil,
writer_enable_idempotence: boolean() | nil,
Expand Down Expand Up @@ -186,6 +188,11 @@ defmodule Fluss.Config do
def set_writer_buffer_wait_timeout_ms(%__MODULE__{} = config, ms) when is_non_neg_integer(ms),
do: %{config | writer_buffer_wait_timeout_ms: ms}

@spec set_writer_kv_backpressure_max_throttle_ms(t(), non_neg_integer()) :: t()
def set_writer_kv_backpressure_max_throttle_ms(%__MODULE__{} = config, ms)
when is_non_neg_integer(ms),
do: %{config | writer_kv_backpressure_max_throttle_ms: ms}

@spec set_writer_dynamic_batch_size_enabled(t(), boolean()) :: t()
def set_writer_dynamic_batch_size_enabled(%__MODULE__{} = config, enabled)
when is_boolean(enabled),
Expand Down
6 changes: 4 additions & 2 deletions fluss-rust/bindings/elixir/lib/fluss/error.ex
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ defmodule Fluss.Error do
Fields:

* `:code` — stable atom for pattern matching.
* `:error_code` — raw integer code. Protocol codes `0..57`, `-1` for
* `:error_code` — raw integer code. Protocol codes are non-negative; `-1` for
`:unknown_server_error`, `-2` for `:client_error`.
* `:message` — human-readable description.

Expand Down Expand Up @@ -96,6 +96,7 @@ defmodule Fluss.Error do
| :ineligible_replica_exception
| :invalid_alter_table_exception
| :deletion_disabled_exception
| :storage_backpressure_exception
| :client_error

@type t :: %__MODULE__{code: code(), error_code: integer(), message: String.t()}
Expand All @@ -113,7 +114,8 @@ defmodule Fluss.Error do
:storage_exception,
:not_enough_replicas_after_append_exception,
:not_enough_replicas_exception,
:leader_not_available_exception
:leader_not_available_exception,
:storage_backpressure_exception
]

@impl true
Expand Down
2 changes: 2 additions & 0 deletions fluss-rust/bindings/elixir/native/fluss_nif/src/atoms.rs
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,7 @@ rustler::atoms! {
ineligible_replica_exception,
invalid_alter_table_exception,
deletion_disabled_exception,
storage_backpressure_exception,
client_error,
}

Expand Down Expand Up @@ -212,6 +213,7 @@ fn api_error_atom(code: i32) -> Atom {
FlussError::IneligibleReplicaException => ineligible_replica_exception(),
FlussError::InvalidAlterTableException => invalid_alter_table_exception(),
FlussError::DeletionDisabledException => deletion_disabled_exception(),
FlussError::StorageBackpressureException => storage_backpressure_exception(),
}
}

Expand Down
4 changes: 4 additions & 0 deletions fluss-rust/bindings/elixir/native/fluss_nif/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@ pub struct NifConfig {
pub writer_bucket_no_key_assigner: Option<NifNoKeyAssigner>,
pub writer_buffer_memory_size: Option<u64>,
pub writer_buffer_wait_timeout_ms: Option<u64>,
pub writer_kv_backpressure_max_throttle_ms: Option<u64>,
pub writer_dynamic_batch_size_enabled: Option<bool>,
pub writer_dynamic_batch_size_min: Option<i32>,
pub writer_enable_idempotence: Option<bool>,
Expand Down Expand Up @@ -130,6 +131,9 @@ impl NifConfig {
if let Some(timeout_ms) = self.writer_buffer_wait_timeout_ms {
config.writer_buffer_wait_timeout_ms = timeout_ms;
}
if let Some(timeout_ms) = self.writer_kv_backpressure_max_throttle_ms {
config.writer_kv_backpressure_max_throttle_ms = timeout_ms;
}
if let Some(enabled) = self.writer_enable_idempotence {
config.writer_enable_idempotence = enabled;
}
Expand Down
8 changes: 8 additions & 0 deletions fluss-rust/bindings/elixir/test/config_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -182,6 +182,14 @@ defmodule Fluss.ConfigTest do
assert config.writer_buffer_wait_timeout_ms == 5_000
end

test "set_writer_kv_backpressure_max_throttle_ms/2 sets the throttle" do
config =
Fluss.Config.new("localhost:9123")
|> Fluss.Config.set_writer_kv_backpressure_max_throttle_ms(1_500)

assert config.writer_kv_backpressure_max_throttle_ms == 1_500
end

test "set_writer_enable_idempotence/2 sets the idempotence flag" do
config =
Fluss.Config.new("localhost:9123")
Expand Down
3 changes: 2 additions & 1 deletion fluss-rust/bindings/elixir/test/error_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,8 @@ defmodule Fluss.ErrorTest do
:storage_exception,
:not_enough_replicas_after_append_exception,
:not_enough_replicas_exception,
:leader_not_available_exception
:leader_not_available_exception,
:storage_backpressure_exception
]

@non_retriable_codes [
Expand Down
5 changes: 5 additions & 0 deletions fluss-rust/bindings/python/fluss/__init__.pyi
Original file line number Diff line number Diff line change
Expand Up @@ -232,6 +232,10 @@ class Config:
@writer_buffer_wait_timeout_ms.setter
def writer_buffer_wait_timeout_ms(self, timeout: int) -> None: ...
@property
def writer_kv_backpressure_max_throttle_ms(self) -> int: ...
@writer_kv_backpressure_max_throttle_ms.setter
def writer_kv_backpressure_max_throttle_ms(self, timeout: int) -> None: ...
@property
def connect_timeout_ms(self) -> int: ...
@connect_timeout_ms.setter
def connect_timeout_ms(self, timeout: int) -> None: ...
Expand Down Expand Up @@ -1280,6 +1284,7 @@ class ErrorCode:
INELIGIBLE_REPLICA_EXCEPTION: int
INVALID_ALTER_TABLE_EXCEPTION: int
DELETION_DISABLED_EXCEPTION: int
STORAGE_BACKPRESSURE_EXCEPTION: int

@final
class OffsetSpec:
Expand Down
20 changes: 20 additions & 0 deletions fluss-rust/bindings/python/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -174,6 +174,14 @@ impl Config {
))
})?;
}
"writer.kv-backpressure.max-throttle-ms" => {
config.writer_kv_backpressure_max_throttle_ms =
value.parse::<u64>().map_err(|e| {
FlussError::new_err(format!(
"Invalid value '{value}' for '{key}': {e}"
))
})?;
}
"writer.bucket.no-key-assigner" => {
config.writer_bucket_no_key_assigner =
value.parse::<fcore::config::NoKeyAssigner>().map_err(|e| {
Expand Down Expand Up @@ -419,6 +427,18 @@ impl Config {
self.inner.writer_buffer_wait_timeout_ms = timeout;
}

/// Get the maximum KV backpressure throttle in milliseconds
#[getter]
fn writer_kv_backpressure_max_throttle_ms(&self) -> u64 {
self.inner.writer_kv_backpressure_max_throttle_ms
}

/// Set the maximum KV backpressure throttle in milliseconds
#[setter]
fn set_writer_kv_backpressure_max_throttle_ms(&mut self, timeout: u64) {
self.inner.writer_kv_backpressure_max_throttle_ms = timeout;
}

/// Get the connect timeout in milliseconds
#[getter]
fn connect_timeout_ms(&self) -> u64 {
Expand Down
3 changes: 3 additions & 0 deletions fluss-rust/bindings/python/src/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -273,4 +273,7 @@ impl ErrorCode {
/// Deletion operations are disabled on this table.
#[classattr]
const DELETION_DISABLED_EXCEPTION: i32 = 57;
/// The server rejected a write due to storage backpressure.
#[classattr]
const STORAGE_BACKPRESSURE_EXCEPTION: i32 = 72;
}
32 changes: 32 additions & 0 deletions fluss-rust/bindings/python/test/test_config.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
# Licensed to the Apache Software Foundation (ASF) under one
# or more contributor license agreements. See the NOTICE file
# distributed with this work for additional information
# regarding copyright ownership. The ASF licenses this file
# to you under the Apache License, Version 2.0 (the
# "License"); you may not use this file except in compliance
# with the License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing,
# software distributed under the License is distributed on an
# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
# KIND, either express or implied. See the License for the
# specific language governing permissions and limitations
# under the License.

import fluss


def test_kv_backpressure_configuration():
config = fluss.Config({"writer.kv-backpressure.max-throttle-ms": "1500"})
assert config.writer_kv_backpressure_max_throttle_ms == 1500

config.writer_kv_backpressure_max_throttle_ms = 750
assert config.writer_kv_backpressure_max_throttle_ms == 750


def test_storage_backpressure_error_is_retriable():
assert fluss.ErrorCode.STORAGE_BACKPRESSURE_EXCEPTION == 72
error = fluss.FlussError("backpressure", 72)
assert error.is_retriable
Loading
Loading