Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
76c9eb9
[c++] Add asynchronous write callbacks
naivedogger Sep 14, 2026
587ba94
[rust] Batch write callback completion notifications
naivedogger Sep 14, 2026
1e9137e
[docs] Clarify write callback compatibility and lifecycle
naivedogger Sep 14, 2026
edc2522
[c++] Demonstrate bounded write callbacks in the existing example
naivedogger Sep 14, 2026
91ecc1e
[c++] Bound pending write callbacks per writer
naivedogger Sep 14, 2026
4b78f44
[c++] Flush waits for pending callbacks to finish
naivedogger Sep 15, 2026
dfd3d1c
[c++] Clarify callback guarantees and recovery guidance
naivedogger Sep 16, 2026
1c987f5
[c++] Tune callback defaults and document buffer sizing
naivedogger Sep 16, 2026
81174ad
[c++] Use notify_all in write callback capacity Release to avoid stal…
naivedogger Sep 19, 2026
582124b
[c++] Bound write callback submission by enqueue_timeout
naivedogger Sep 20, 2026
18af731
[c++] Make callback worker count configurable via FLUSS_CALLBACK_WORKERS
naivedogger Sep 20, 2026
f332fe3
[c++] Document callback execution context and recovery guidance
naivedogger Sep 20, 2026
b942134
nit
naivedogger Sep 21, 2026
8c35446
[c++] Bound callback submission by client.writer.buffer.wait-timeout
naivedogger Sep 21, 2026
81eb38a
[c++] Show writer_buffer_wait_timeout_ms and test bounded callback su…
naivedogger Sep 21, 2026
019fde8
[c++] Pass callback capacity limit directly
naivedogger Sep 21, 2026
a4a3648
[c++] Polish callback timeout documentation
naivedogger Sep 21, 2026
6d0fd62
[c++] Address write callback review feedback
naivedogger Sep 22, 2026
cb7258a
[c++] Improve write callback ordering and backpressure
naivedogger Sep 22, 2026
6709ce7
nit
naivedogger Sep 23, 2026
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: 7 additions & 0 deletions fluss-rust/bindings/cpp/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,9 @@ genrule(
name = "cargo_build_debug",
srcs = glob([
"src/**/*.rs",
"src/**/*.hpp",
"include/**/*.hpp",
"build.rs",
"Cargo.toml",
]),
outs = [
Expand Down Expand Up @@ -121,6 +124,9 @@ genrule(
name = "cargo_build_release",
srcs = glob([
"src/**/*.rs",
"src/**/*.hpp",
"include/**/*.hpp",
"build.rs",
"Cargo.toml",
]),
outs = [
Expand Down Expand Up @@ -274,6 +280,7 @@ cc_library(
textual_hdrs = [
"src/ffi_converter.hpp",
"src/type_lowering.hpp",
"src/write_callback.hpp",
":rust_bridge_h_unified",
":lib_rs_h_unified",
":cxx_h_unified",
Expand Down
36 changes: 35 additions & 1 deletion fluss-rust/bindings/cpp/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,7 @@ not apply to `CreateBucketBatchScanner()`.

## Examples and Documentation

- [examples/example.cpp](examples/example.cpp) demonstrates log-table writes, continuous scans,
- [examples/example.cpp](examples/example.cpp) demonstrates log-table writes with Wait and bounded callbacks, continuous scans,
bounded Arrow record-batch scans, projections, and offset queries.
- [examples/admin_example.cpp](examples/admin_example.cpp) demonstrates database, table,
partition, and cluster administration.
Expand All @@ -76,6 +76,40 @@ not apply to `CreateBucketBatchScanner()`.
[C++ API reference](../../website/docs/user-guide/cpp/api-reference.md) and
[log-table examples](../../website/docs/user-guide/cpp/example/log-tables.md).

The SDK executes `WriteCallback` (`void(const WriteCompletion&)`) on one shared
worker, serially in dispatch order and off the I/O threads. `WriteCompletion.result`
is the write outcome; copy it before passing it to another worker.
`CreateWriter(writer)` uses the default `WriteCallbackOptions`; the overload
`CreateWriter(writer, options)` accepts a positive `max_pending_operations` limit
(default 262144) per writer. This operation-count budget is independent of the
Connection's byte-counted write buffer. Slow callbacks can fill it even when the
write buffer has room. The Rust write-buffer permit is released when the batch
completes, before the user callback returns, but the callback object and its
captures remain retained until callback completion. Once the per-writer callback
limit is full, callback-based submissions wait or fail according to
`client.writer.buffer.wait-timeout`; its default is unbounded.

Callback-capacity and buffer waits share `client.writer.buffer.wait-timeout`.
Zero makes those waits fail fast; this is not a deadline for the entire API call,
ACKs, retries, or callback execution. See the
[buffer sizing guidance](../../website/docs/user-guide/cpp/api-reference.md#sizing-callback-capacity-and-write-buffers)
for independent capacity and byte budgets. Callback worker initialization failure
rejects the submission before any data is accepted; there is no parallel fallback.

A failed callback does not prove that the record was not written. Application
resubmission can duplicate it, even with SDK idempotence enabled. The example only
counts and logs outcomes; it does not implement durable recovery. Keep callbacks
short, protect shared state, and handle retries outside the callback with an
application recovery policy.

After submissions stop, `Flush()` first flushes writes and, on success, blocks until
pending callbacks finish, acting as a barrier. A callback that never returns hangs it.
Calling it from a write callback is rejected before flushing. If a write flush
returns an error, keep callback state alive; if
it succeeds, still check individual write results.
See the [callback guarantees and recovery guidance](../../website/docs/user-guide/cpp/api-reference.md#write-guarantees-and-recovery)
for result semantics, callback implementation, and shutdown requirements.

For a bounded log scan, pass the per-bucket offset ranges directly to `TableScan`. The returned
reader yields one Arrow batch at a time until every `[starting_offset, stopping_offset)` range
is complete:
Expand Down
4 changes: 4 additions & 0 deletions fluss-rust/bindings/cpp/build.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,8 +17,12 @@

fn main() {
cxx_build::bridge("src/lib.rs")
.include("include")
.include("src")
.std("c++17")
.compile("fluss-cpp-bridge");

println!("cargo:rerun-if-changed=src/lib.rs");
println!("cargo:rerun-if-changed=src/write_callback.hpp");
println!("cargo:rerun-if-changed=include/fluss.hpp");
}
81 changes: 81 additions & 0 deletions fluss-rust/bindings/cpp/examples/example.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -20,9 +20,13 @@
#include <arrow/record_batch.h>
#include <arrow/type.h>

#include <atomic>
#include <chrono>
#include <iostream>
#include <memory>
#include <mutex>
#include <unordered_map>
#include <utility>
#include <vector>

#include "fluss.hpp"
Expand All @@ -39,6 +43,10 @@ int main() {
// 1) Connect
fluss::Configuration config;
config.bootstrap_servers = "127.0.0.1:9123";
// Callback capacity and buffer waits share this budget. The default UINT64_MAX
// waits indefinitely; a finite value allows handling overload instead. Zero fails
// fast when either resource is unavailable. This is not a whole-call deadline.
config.writer_buffer_wait_timeout_ms = 30000;

fluss::Connection conn;
check("create", fluss::Connection::Create(config, conn));
Expand Down Expand Up @@ -85,6 +93,10 @@ int main() {

// 5) Write rows with scalar and temporal values
fluss::AppendWriter writer;
// Use the default per-writer limit of 262144 outstanding callback operations.
// This is independent of config.writer_buffer_memory_size (per Connection).
// To tune it, pass WriteCallbackOptions based on measured completion latency
// and capture memory; a smaller limit can throttle writes before the buffer fills.
check("new_append_writer", table.NewAppend().CreateWriter(writer));

struct RowData {
Expand Down Expand Up @@ -143,6 +155,75 @@ int main() {
std::cout << "Row acknowledged by server" << std::endl;
}

// Callback acknowledgment
{
// The SDK runs callbacks; no application waiting thread is required.
// Callbacks run on one shared worker, so keep them short and
// non-blocking: do not Flush/Wait or retry synchronously inside a callback.
// This example counts outcomes only. It does not implement durable recovery.
struct CallbackState {
std::atomic<size_t> succeeded{0};
std::atomic<size_t> failed{0};
std::mutex mutex;
int32_t first_failed_id{0};
fluss::Result first_failure;
};
// Shared ownership also keeps state alive if submission throws or flushing fails.
auto state = std::make_shared<CallbackState>();
bool submission_failed = false;
for (const auto& r : rows) {
const int32_t id = 1000 + r.id;
fluss::GenericRow row;
row.SetInt32(0, id);
row.SetString(1, r.name);
row.SetFloat32(2, r.score);
row.SetInt32(3, r.age);
row.SetDate(4, r.date);
row.SetTime(5, r.time);
row.SetTimestampNtz(6, r.ts_ntz);
row.SetTimestampLtz(7, r.ts_ltz);
auto submitted = writer.Append(
row, [id, state](const fluss::WriteCompletion& notification) {
const auto& result = notification.result;
if (result.Ok()) {
++state->succeeded;
} else {
// Copy the invocation-scoped result. Do not log, perform I/O,
// or retry here: one slow callback delays every writer.
if (state->failed.fetch_add(1) == 0) {
std::lock_guard<std::mutex> lock(state->mutex);
state->first_failed_id = id;
state->first_failure = result;
}
}
});
if (!submitted.Ok()) {
// No callback will run for this submission; handle this path too.
submission_failed = true;
std::cerr << "Submission failed for id=" << id << ": " << submitted.error_message
<< '\n';
break;
}
}
// Stop submissions, then drain accepted callbacks. Flush success alone does
// not mean every write succeeded. Writer/Connection destruction is not a drain.
check("flush", writer.Flush());
std::cout << "Callback writes: succeeded=" << state->succeeded.load()
<< " failed=" << state->failed.load() << '\n';
if (state->failed.load() != 0) {
std::lock_guard<std::mutex> lock(state->mutex);
std::cerr << "First failed id=" << state->first_failed_id
<< ": " << state->first_failure.error_message << '\n';
}
// A failed write may have reached the server. Recover outside the callback
// using retained input or a replayable source and application-level deduplication.
// Advance a source position only after Flush and all relevant writes succeed.
// This example reports failure and exits; it does not implement durable recovery.
if (submission_failed || state->failed.load() != 0) {
return 1;
}
}

// Append a row with all fields null (matches Rust log_table.rs all_supported_datatypes)
{
fluss::GenericRow row;
Expand Down
110 changes: 105 additions & 5 deletions fluss-rust/bindings/cpp/include/fluss.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@

#include <chrono>
#include <cstdint>
#include <functional>
#include <limits>
#include <memory>
#include <optional>
Expand All @@ -47,6 +48,8 @@ struct Admin;
struct Table;
struct AppendWriter;
struct WriteResult;
class WriteCallback;
class WriteCallbackCapacity;
struct LogScanner;
struct RecordBatchLogReader;
struct BatchScanner;
Expand Down Expand Up @@ -527,10 +530,62 @@ struct Result {

bool Ok() const { return error_code == 0; }

/// Returns true if retrying the request may succeed. Client-side errors always return false.
/// Returns true if retrying the request may succeed. Does not guarantee that a failed
/// write had no effect or that application resubmission is duplicate-safe.
/// Client-side errors always return false.
bool IsRetriable() const { return ErrorCode::IsRetriable(error_code); }
};

/// Write-specific completion metadata. The reference passed to a callback is valid
/// only for that invocation; copy the result when retaining it for later work.
struct WriteCompletion {
Result result;
};

/// Per-writer admission control for callback operations, independent of buffer bytes.
struct WriteCallbackOptions {
// Includes accepted writes awaiting completion and callbacks queued or executing.
// Must be greater than zero. This is not a byte limit on callback captures.
size_t max_pending_operations{262144};
};

/// Receives the final outcome of an accepted write. Function pointers and lambdas
/// are supported. An empty callback is rejected before submitting the write.
///
/// During normal operation, the SDK owns the callback until completion and invokes
/// it exactly once on a shared SDK worker; no caller polling or waiting thread
/// is needed. Process exit or a crash can prevent delivery. Callbacks never run
/// inline in the submitting call, but may start before the call returns.
/// Keep callbacks short and synchronize access to shared state, including writers.
/// Callback overloads do not make writers safe for concurrent access. Captured
/// references must remain valid until the callback finishes; capturing shared
/// ownership is recommended. Keep the connection alive until completion.
///
/// Success follows the configured acknowledgment policy. Errors are reported after
/// internal retry handling, but do not guarantee that no data was written.
/// Application resubmission is a new operation and can produce duplicates even
/// with SDK idempotence enabled. Retain input identifiers and recovery state as
/// needed; one callback invocation is not an exactly-once delivery guarantee.
///
/// Callbacks run serially in dispatch order on a single shared worker. There is
/// no cross-bucket submission-order guarantee; late registrations are queued
/// when registered. Do not wait for another callback from a callback: it stalls that
/// worker. Synchronous SDK calls require exclusive writer access. Callback
/// submissions to a full writer fail immediately when called from a callback,
/// instead of blocking the worker. Each writer bounds its outstanding callback
/// operations using WriteCallbackOptions, independently of the write buffer size.
/// Hand off retries or expensive work without blocking; bound application queues
/// and handle overflow without silently discarding failed operations.
///
/// Exceptions thrown by callbacks are caught and reported to stderr; they do not
/// change the write outcome or retry the callback. Stop submissions before Flush().
/// After a successful Rust write flush, Flush() blocks until pending callbacks finish,
/// acting as a barrier, so a callback that never returns hangs it. On error, referenced
/// state may still be in use.
/// Flush() called inside any write callback returns a client error without flushing.
/// Flush() does not wait for work handed to application workers or retry queues.
using WriteCallback = std::function<void(const WriteCompletion&)>;

struct TablePath {
std::string database_name;
std::string table_name;
Expand Down Expand Up @@ -1553,9 +1608,12 @@ struct Configuration {
bool writer_enable_idempotence{true};
// Maximum number of in-flight requests per bucket for idempotent writes
size_t writer_max_inflight_requests_per_bucket{5};
// Total memory available for buffering write batches (default 64MB)
// Shared write-batch memory budget per Connection, across its tables and writers
// (default 64 MiB). Not a process RSS limit or a callback-capture memory budget.
size_t writer_buffer_memory_size{64 * 1024 * 1024};
// Maximum time in milliseconds to block waiting for buffer memory
// Shared wait budget in milliseconds for buffer memory and callback capacity.
// Does not bound data conversion, scheduling, ACKs, or callback execution.
// UINT64_MAX waits indefinitely; zero fails fast when capacity or memory is unavailable.
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};
Expand Down Expand Up @@ -1740,6 +1798,8 @@ class TableAppend {
TableAppend& operator=(TableAppend&&) noexcept = default;

Result CreateWriter(AppendWriter& out);
/// Create a writer with an independent, positive callback operation limit.
Result CreateWriter(AppendWriter& out, const WriteCallbackOptions& options);

private:
friend class Table;
Expand All @@ -1759,6 +1819,8 @@ class TableUpsert {
TableUpsert& PartialUpdateByName(std::vector<std::string> column_names);

Result CreateWriter(UpsertWriter& out);
/// Create a writer sharing one callback operation limit across upserts and deletes.
Result CreateWriter(UpsertWriter& out, const WriteCallbackOptions& options);

private:
friend class Table;
Expand Down Expand Up @@ -1877,6 +1939,7 @@ class WriteResult {
friend class UpsertWriter;
WriteResult(ffi::WriteResult* inner) noexcept;

Result Notify(std::unique_ptr<ffi::WriteCallback> callback);
void Destroy() noexcept;
ffi::WriteResult* inner_{nullptr};
};
Expand All @@ -1895,17 +1958,37 @@ class AppendWriter {

Result Append(const GenericRow& row);
Result Append(const GenericRow& row, WriteResult& out);
/// Submit a row and notify callback of its final outcome without waiting for
/// acknowledgment. Callback capacity and buffer waits share the budget from
/// client.writer.buffer.wait-timeout. Ok means the write was accepted and the
/// callback fires exactly once during normal operation; an error means
/// submission failed and no callback runs. A zero timeout makes admission
/// fail fast when callback capacity or buffer memory is unavailable.
Result Append(const GenericRow& row, WriteCallback callback);
Result AppendArrowBatch(const std::shared_ptr<arrow::RecordBatch>& batch);
Result AppendArrowBatch(const std::shared_ptr<arrow::RecordBatch>& batch, WriteResult& out);
/// Like the callback Append overload, but notifies once for the entire batch.
Result AppendArrowBatch(const std::shared_ptr<arrow::RecordBatch>& batch,
WriteCallback callback);
Result Flush();

private:
friend class Table;
friend class TableAppend;
AppendWriter(ffi::AppendWriter* writer) noexcept;
AppendWriter(ffi::AppendWriter* writer,
std::shared_ptr<ffi::WriteCallbackCapacity> callback_capacity) noexcept;

// Submit through FFI bounding the buffer-backpressure wait by submit_budget_ms
// (Kafka max.block.ms style): negative uses the writer's configured buffer wait
// timeout, >= 0 caps the wait at that many ms (0 = fail fast when the buffer is
// full). Only the callback path passes a budget; the public overloads pass -1.
Result AppendWithBudget(const GenericRow& row, WriteResult& out, int64_t submit_budget_ms);
Result AppendArrowBatchWithBudget(const std::shared_ptr<arrow::RecordBatch>& batch,
WriteResult& out, int64_t submit_budget_ms);

void Destroy() noexcept;
ffi::AppendWriter* writer_{nullptr};
std::shared_ptr<ffi::WriteCallbackCapacity> callback_capacity_;
};

class UpsertWriter {
Expand All @@ -1922,16 +2005,33 @@ class UpsertWriter {

Result Upsert(const GenericRow& row);
Result Upsert(const GenericRow& row, WriteResult& out);
/// Submit an upsert and notify callback of its final outcome without waiting
/// for acknowledgment. Callback capacity and buffer waits share the budget from
/// client.writer.buffer.wait-timeout. Ok means the write was accepted and the
/// callback fires exactly once during normal operation; an error means
/// submission failed and no callback runs. A zero timeout makes admission
/// fail fast when callback capacity or buffer memory is unavailable.
Result Upsert(const GenericRow& row, WriteCallback callback);
Result Delete(const GenericRow& row);
Result Delete(const GenericRow& row, WriteResult& out);
/// Like the callback Upsert overload, but deletes a row by primary key.
Result Delete(const GenericRow& row, WriteCallback callback);
Result Flush();

private:
friend class Table;
friend class TableUpsert;
UpsertWriter(ffi::UpsertWriter* writer) noexcept;
UpsertWriter(ffi::UpsertWriter* writer,
std::shared_ptr<ffi::WriteCallbackCapacity> callback_capacity) noexcept;

// See AppendWriter::AppendWithBudget for submit_budget_ms semantics. Only the
// callback path passes a budget; the public overloads pass -1.
Result UpsertWithBudget(const GenericRow& row, WriteResult& out, int64_t submit_budget_ms);
Result DeleteWithBudget(const GenericRow& row, WriteResult& out, int64_t submit_budget_ms);

void Destroy() noexcept;
ffi::UpsertWriter* writer_{nullptr};
std::shared_ptr<ffi::WriteCallbackCapacity> callback_capacity_;
};

class Lookuper {
Expand Down
Loading
Loading