lxy-9602 commented on code in PR #376:
URL: https://github.com/apache/paimon-cpp/pull/376#discussion_r4130122875


##########
src/paimon/format/vortex/vortex_format_writer.cpp:
##########
@@ -0,0 +1,178 @@
+/*
+ * 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.
+ */
+
+#include "paimon/format/vortex/vortex_format_writer.h"
+
+#include <utility>
+
+#include "arrow/c/bridge.h"
+#include "arrow/memory_pool.h"
+#include "arrow/type.h"
+#include "paimon/common/metrics/metrics_impl.h"
+#include "paimon/common/utils/arrow/arrow_utils.h"
+#include "paimon/common/utils/arrow/status_utils.h"
+#include "paimon/format/vortex/vortex_ffi_util.h"
+#include "paimon/fs/file_system.h"
+
+namespace paimon::vortex {
+
+VortexFormatWriter::VortexFormatWriter(std::shared_ptr<OutputStream> output,
+                                       std::shared_ptr<arrow::Schema> schema, 
VxSessionPtr session,
+                                       std::shared_ptr<VortexOutputContext> 
output_context,
+                                       vx_callback_sink* sink,
+                                       std::shared_ptr<arrow::MemoryPool> 
arrow_pool)
+    : output_(std::move(output)),
+      schema_(std::move(schema)),
+      arrow_pool_(std::move(arrow_pool)),
+      session_(std::move(session)),
+      output_context_(std::move(output_context)),
+      sink_(sink),
+      metrics_(std::make_shared<MetricsImpl>()) {}
+
+Result<std::unique_ptr<VortexFormatWriter>> VortexFormatWriter::Create(
+    const std::shared_ptr<OutputStream>& output, const 
std::shared_ptr<arrow::Schema>& schema,
+    const std::shared_ptr<arrow::MemoryPool>& arrow_pool) {
+    if (output == nullptr || schema == nullptr || arrow_pool == nullptr) {
+        return Status::Invalid("Vortex writer requires non-null output, schema 
and arrow pool");
+    }
+    VxSessionPtr session(vx_session_new(), vx_session_free);
+    if (session == nullptr) {
+        return Status::IOError("failed to create Vortex session");
+    }
+    ::ArrowSchema ffi_schema = {};
+    PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportSchema(*schema, &ffi_schema));
+    vx_error* error = nullptr;
+    // vx_dtype_from_arrow_schema consumes ffi_schema on both success and 
failure.
+    VxDtypePtr dtype(vx_dtype_from_arrow_schema(&ffi_schema, &error), 
vx_dtype_free);
+    if (dtype == nullptr) {
+        return VortexFfiError("convert Arrow schema to Vortex dtype", error);
+    }
+    error = nullptr;
+    auto output_context = std::make_shared<VortexOutputContext>(output);
+    vx_callback_sink* sink = vx_callback_sink_open(
+        session.get(), VortexOutputContext::MakeCallbacks(output_context), 
dtype.get(), &error);
+    if (sink == nullptr) {
+        return VortexCallbackError("open Vortex callback sink", error,
+                                   output_context->GetCallbackStatus());
+    }
+    return std::unique_ptr<VortexFormatWriter>(new VortexFormatWriter(
+        output, schema, std::move(session), std::move(output_context), sink, 
arrow_pool));
+}
+
+VortexFormatWriter::~VortexFormatWriter() {
+    if (sink_ != nullptr) {
+        // Not finished: abort so Vortex drops the sink without writing a 
footer. Whatever reached
+        // the output stream is an incomplete file, which the caller discards 
along with it.
+        vx_callback_sink_abort(sink_);
+        sink_ = nullptr;
+    }
+}
+
+Status VortexFormatWriter::AddBatch(::ArrowArray* batch) {
+    if (batch == nullptr) {
+        return Status::Invalid("Vortex writer batch is nullptr");
+    }
+    if (finished_ || sink_ == nullptr) {
+        return Status::Invalid("cannot add a batch after Vortex writer is 
finished");
+    }
+    // vx_array_from_arrow imports through arrow-rs, which cannot represent a 
sliced (offset > 0)
+    // top-level struct: the parent offset gets re-applied to children that 
Arrow C++ already
+    // exported with it, tripping an arrow-data slice assertion (end <= len) 
that aborts the
+    // process. Rebase such a batch to offset 0 first, mirroring the read 
path. Offset-0 batches
+    // (the common case) are handed over untouched, so the hot path adds no 
copy.
+    ::ArrowArray* vortex_batch = batch;
+    ::ArrowArray rebased_batch{};
+    if (batch->offset != 0) {
+        PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(
+            std::shared_ptr<arrow::Array> imported,
+            arrow::ImportArray(batch, arrow::struct_(schema_->fields())));
+        PAIMON_ASSIGN_OR_RAISE(std::shared_ptr<arrow::Array> normalized,
+                               ArrowUtils::NormalizeArrayOffsets(imported, 
arrow_pool_.get()));
+        PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportArray(*normalized, 
&rebased_batch));
+        vortex_batch = &rebased_batch;
+    }
+    ::ArrowSchema ffi_schema = {};
+    PAIMON_RETURN_NOT_OK_FROM_ARROW(arrow::ExportSchema(*schema_, 
&ffi_schema));
+    vx_error* error = nullptr;
+    // vx_array_from_arrow consumes both `vortex_batch` and `ffi_schema` on 
success and on failure.
+    VxArrayPtr array(vx_array_from_arrow(vortex_batch, &ffi_schema, 
/*nullable=*/false, &error),
+                     vx_array_free);
+    if (array == nullptr) {
+        return VortexFfiError("convert Arrow batch to Vortex array", error);
+    }
+    error = nullptr;
+    vx_callback_sink_push(sink_, array.get(), &error);
+    if (error != nullptr) {
+        return VortexCallbackError("push batch to Vortex sink", error,
+                                   output_context_->GetCallbackStatus());
+    }
+    return Status::OK();
+}
+
+Status VortexFormatWriter::Flush() {
+    if (finished_) {
+        return Status::OK();
+    }
+    // Flush through the output context so it is serialized with the 
background writer task's writes
+    // (OutputStream has no concurrent Write/Flush contract). Vortex buffers 
internally, so this
+    // flushes the bytes handed over so far; the rest is drained when the sink 
closes in Finish().
+    return output_context_->FlushStream();

Review Comment:
   Could you please confirm whether this `FlushStream` is really necessary? I 
noticed that `output_context` adds locks around `flush` and `write`. Could 
`Flush` just return `OK` here, and let `Finish` ensure all data is written out? 
If so, could we remove the lock on output as well? Please check this.



##########
crates/vortex_callback_io/callback_io.rs:
##########
@@ -0,0 +1,781 @@
+// 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.
+
+//! Callback-based I/O for the Vortex C API.
+//!
+//! The stock C API can only read a whole file already materialized in memory
+//! (`vx_data_source_new_buffer`) or a path that Vortex itself resolves
+//! (`vx_data_source_new`), and it can only write to a local filesystem path
+//! (`vx_array_sink_open_file`, which uses `async_fs::File::create`). paimon 
owns its own pluggable
+//! filesystem, so none of these fit: reading forces the entire file into 
memory, and writing forces
+//! a local temporary file that then has to be copied back.
+//!
+//! This module adds entry points that take host callbacks instead, so Vortex 
issues positional
+//! reads and sequential writes straight against paimon's `InputStream` / 
`OutputStream`. Vortex's
+//! core I/O traits are the plug points; the shape of the bridge mirrors 
`vortex-jni`'s
+//! `JavaReadable` and `JavaWrite`.
+//!
+//! The write side needs its own handle type rather than reusing 
`vx_array_sink`, whose fields are
+//! private to `crate::sink`.
+//!
+//! The sink entry points below (`vx_callback_sink_push` / `_close` / 
`_abort`) are adapted from
+//! vortex-ffi's `src/sink.rs` (SPDX-License-Identifier: Apache-2.0, Copyright 
the Vortex
+//! contributors), reworked to drive the write pipeline through host callbacks 
instead of a
+//! Vortex-resolved filesystem path.
+//!
+//! This file is maintained in the paimon-cpp tree 
(`crates/vortex_callback_io/`) and copied into
+//! `vortex-ffi/src/` at build time; see `cmake_modules/vortex.diff`.
+//!
+//! The C declarations that mirror the types and entry points below live in 
paimon-cpp's
+//! `src/paimon/format/vortex/vortex_ffi.h` (cbindgen does not see this 
module, so they are written
+//! by hand). The two sides are only checked by human review, and a mismatch 
in field order or a
+//! function signature is silent memory corruption, not a compile error. Any 
change to a
+//! `#[repr(C)]` struct, a callback `type`, or a `#[no_mangle]` entry point 
here MUST be mirrored in
+//! that header in the same change.
+
+use std::ffi::c_void;
+use std::io;
+use std::sync::Arc;
+
+use futures::FutureExt;
+use futures::SinkExt;
+use futures::TryStreamExt;
+use futures::channel::mpsc;
+use futures::channel::mpsc::Sender;
+use futures::future::BoxFuture;
+use vortex::array::ArrayRef;
+use vortex::array::buffer::BufferHandle;
+use vortex::array::stream::ArrayStreamAdapter;
+use vortex::buffer::Alignment;
+use vortex::buffer::ByteBufferMut;
+use vortex::dtype::DType;
+use vortex::error::VortexResult;
+use vortex::error::vortex_bail;
+use vortex::error::vortex_ensure;
+use vortex::error::vortex_err;
+use vortex::file::OpenOptionsSessionExt;
+use vortex::file::WriteOptionsSessionExt;
+use vortex::file::WriteStrategyBuilder;
+use vortex::file::WriteSummary;
+use vortex::io::CoalesceConfig;
+use vortex::io::IoBuf;
+use vortex::io::VortexReadAt;
+use vortex::io::VortexWrite;
+use vortex::io::runtime::BlockingRuntime;
+use vortex::io::runtime::Handle;
+use vortex::io::runtime::Task;
+use vortex::io::session::RuntimeSessionExt;
+
+use crate::RUNTIME;
+use crate::array::vx_array;
+use crate::data_source::vx_data_source;
+use crate::dtype::vx_dtype;
+use crate::error::try_or;
+use crate::error::try_or_default;
+use crate::error::vx_error;
+use crate::session::vx_session;
+
+/// Fill `length` bytes starting at `offset` into `dst`.
+///
+/// Returns 0 on success and non-zero on failure; a partial read must be 
reported as a failure.
+/// The host is expected to keep its own error detail on the context, since 
only a status code
+/// crosses the boundary.
+pub type vx_read_at_fn =
+    unsafe extern "C" fn(ctx: *mut c_void, offset: u64, dst: *mut u8, length: 
usize) -> i32;
+
+/// Release the host context. Called exactly once, when the owning handle is 
freed.
+pub type vx_release_fn = unsafe extern "C" fn(ctx: *mut c_void);
+
+/// Append `length` bytes from `src` to the host sink.
+///
+/// Returns 0 on success and non-zero on failure; a partial write must be 
reported as a failure.
+/// Unlike reads, writes are sequential and never concurrent for one sink.
+pub type vx_write_fn = unsafe extern "C" fn(ctx: *mut c_void, src: *const u8, 
length: usize) -> i32;
+
+/// Flush whatever the host has buffered. Returns 0 on success and non-zero on 
failure.
+pub type vx_flush_fn = unsafe extern "C" fn(ctx: *mut c_void) -> i32;
+
+/// Host callbacks backing a positional reader.
+#[repr(C)]
+#[derive(Clone, Copy)]
+pub struct vx_input_callbacks {
+    /// Opaque host context, passed back to every callback.
+    pub ctx: *mut c_void,
+    /// Positional read. Required.
+    ///
+    /// Must be thread-safe: Vortex issues concurrent positional reads, so 
this is called from
+    /// several blocking threads at once for the same context.
+    pub read_at_fn: Option<vx_read_at_fn>,
+    /// Context destructor. Optional; when set, it is the last callback 
invoked.
+    pub release_fn: Option<vx_release_fn>,
+}
+
+/// Host callbacks backing a sequential writer.
+#[repr(C)]
+#[derive(Clone, Copy)]
+pub struct vx_output_callbacks {
+    /// Opaque host context, passed back to every callback.
+    pub ctx: *mut c_void,
+    /// Sequential write. Required.
+    pub write_fn: Option<vx_write_fn>,
+    /// Flush. Optional; when null, flush requests are ignored.
+    pub flush_fn: Option<vx_flush_fn>,
+    /// Context destructor. Optional; when set, it is the last callback 
invoked.
+    pub release_fn: Option<vx_release_fn>,
+}
+
+/// The host context, owned by the reader.
+///
+/// Held behind an `Arc` so in-flight reads keep the context alive, and so 
`release_fn` runs exactly
+/// once when the last reference goes away. It carries only the context, not 
the read callback, so
+/// that ownership can be taken over before anything else is validated.
+struct CallbackTarget {
+    ctx: *mut c_void,
+    release_fn: Option<vx_release_fn>,
+}
+
+// SAFETY: the host contract for `vx_input_callbacks` requires `read_at_fn` to 
be thread-safe for
+// concurrent calls on the same `ctx`, which is what makes it sound to move 
the context across
+// threads and share it between concurrent reads.
+unsafe impl Send for CallbackTarget {}
+// SAFETY: see the `Send` impl above.
+unsafe impl Sync for CallbackTarget {}
+
+impl Drop for CallbackTarget {
+    fn drop(&mut self) {
+        if let Some(release) = self.release_fn {
+            // SAFETY: `ctx` was supplied by the host together with 
`release_fn`, and an `Arc`
+            // guarantees this runs once, after every in-flight read has 
finished.
+            unsafe { release(self.ctx) };
+        }
+    }
+}
+
+/// A [`VortexReadAt`] that forwards positional reads to host callbacks.
+///
+/// Reads run on the runtime's blocking pool, since the host callback is 
synchronous and must not
+/// occupy an async executor thread. The file size is supplied at construction 
time, so `size()`
+/// never crosses the FFI boundary.
+struct CallbackReadAt {
+    target: Arc<CallbackTarget>,
+    read_at_fn: vx_read_at_fn,
+    len: u64,
+    handle: Handle,
+    concurrency: usize,
+}
+
+/// Number of concurrent host read callbacks Vortex may have in flight for one 
file.
+///
+/// The host filesystem may be remote, where concurrency hides latency; it is 
also the bound on
+/// blocking threads occupied by one reader.
+const DEFAULT_CONCURRENCY: usize = 16;
+
+impl VortexReadAt for CallbackReadAt {
+    fn coalesce_config(&self) -> Option<CoalesceConfig> {
+        // Each read costs an FFI hop plus a blocking-pool hand-off, and the 
host filesystem may be
+        // remote, so favor fewer and larger reads.
+        Some(CoalesceConfig::object_storage())
+    }
+
+    fn concurrency(&self) -> usize {
+        self.concurrency
+    }
+
+    fn size(&self) -> BoxFuture<'static, VortexResult<u64>> {
+        let len = self.len;
+        async move { Ok(len) }.boxed()
+    }
+
+    fn read_at(
+        &self,
+        offset: u64,
+        length: usize,
+        alignment: Alignment,
+    ) -> BoxFuture<'static, VortexResult<BufferHandle>> {
+        let target = Arc::clone(&self.target);
+        let read_at_fn = self.read_at_fn;
+        let len = self.len;
+        let handle = self.handle.clone();
+
+        async move {
+            handle
+                .spawn_blocking(move || {
+                    let end = offset
+                        .checked_add(length as u64)
+                        .ok_or_else(|| vortex_err!("read {offset}+{length} 
overflows u64"))?;
+                    if end > len {
+                        vortex_bail!("read {offset}..{end} out of bounds for 
file of length {len}");
+                    }
+
+                    let mut buffer = 
ByteBufferMut::with_capacity_aligned(length, alignment);
+                    if length > 0 {
+                        // SAFETY: the spare capacity covers `length` bytes of 
live allocation, and
+                        // the host contract requires the callback to fill 
exactly that many bytes
+                        // before returning success.
+                        let status = unsafe {
+                            read_at_fn(
+                                target.ctx,
+                                offset,
+                                
buffer.spare_capacity_mut().as_mut_ptr().cast::<u8>(),
+                                length,
+                            )
+                        };
+                        if status != 0 {
+                            vortex_bail!(
+                                "host read callback failed with status 
{status} for {offset}..{end}"
+                            );
+                        }
+                    }
+                    // SAFETY: the callback reported success, which per its 
contract means all
+                    // `length` bytes were written.
+                    unsafe { buffer.set_len(length) };
+                    Ok(BufferHandle::new_host(buffer.freeze()))
+                })
+                .await
+        }
+        .boxed()
+    }
+}
+
+unsafe fn data_source_new_callback(
+    session: *const vx_session,
+    callbacks: vx_input_callbacks,
+    size: u64,
+) -> VortexResult<*const vx_data_source> {
+    // Take ownership of the host context before anything else, so an early 
error below still
+    // releases it exactly once when this `Arc` is dropped.
+    let target = Arc::new(CallbackTarget {
+        ctx: callbacks.ctx,
+        release_fn: callbacks.release_fn,
+    });
+
+    vortex_ensure!(!session.is_null());
+    let read_at_fn = callbacks
+        .read_at_fn
+        .ok_or_else(|| vortex_err!("vx_input_callbacks.read_at_fn is 
required"))?;
+
+    let session = vx_session::as_ref(session).clone();
+    let reader: Arc<dyn VortexReadAt> = Arc::new(CallbackReadAt {
+        target,
+        read_at_fn,
+        len: size,
+        handle: session.handle(),
+        concurrency: DEFAULT_CONCURRENCY,
+    });
+
+    let file = RUNTIME.block_on(async { 
session.open_options().open(reader).await })?;
+    Ok(vx_data_source::new(file.data_source()?))
+}
+
+/// Create a data source that reads through host callbacks.
+///
+/// `size` is the total length of the file in bytes; the host knows it up 
front, so it is passed
+/// here instead of being fetched through another callback.
+///
+/// The callbacks (and the context they carry) are owned by the returned data 
source and released
+/// when it is freed, including when this call fails.
+///
+/// On error, returns NULL and sets "err".
+///
+/// # Safety
+///
+/// `session` must be a valid `vx_session`. `callbacks.read_at_fn` must be 
non-null, thread-safe,
+/// and valid for as long as the data source lives.
+#[unsafe(no_mangle)]
+pub unsafe extern "C-unwind" fn vx_data_source_new_callback(
+    session: *const vx_session,
+    callbacks: vx_input_callbacks,
+    size: u64,
+    err: *mut *mut vx_error,
+) -> *const vx_data_source {
+    try_or(err, std::ptr::null(), || unsafe {
+        data_source_new_callback(session, callbacks, size)
+    })
+}
+
+/// A [`VortexWrite`] that forwards sequential writes to host callbacks.
+///
+/// Writes run inline on the runtime thread driving the write task, the same 
way `vortex-jni`
+/// forwards to Java sinks: the host callback is synchronous and the layout 
writer hands over one
+/// buffer at a time, so there is nothing to overlap.
+struct CallbackWrite {
+    target: Arc<CallbackTarget>,
+    write_fn: vx_write_fn,
+    flush_fn: Option<vx_flush_fn>,
+}
+
+impl CallbackWrite {
+    fn write_slice(&self, bytes: &[u8]) -> io::Result<()> {
+        if bytes.is_empty() {
+            return Ok(());
+        }
+        // SAFETY: `bytes` points at a live slice for the duration of the 
call, and the host
+        // contract forbids retaining it afterwards.
+        let status = unsafe { (self.write_fn)(self.target.ctx, bytes.as_ptr(), 
bytes.len()) };
+        if status != 0 {
+            return Err(io::Error::other(format!(
+                "host write callback failed with status {status}"
+            )));
+        }
+        Ok(())
+    }
+
+    fn flush_host(&self) -> io::Result<()> {
+        let Some(flush) = self.flush_fn else {
+            return Ok(());
+        };
+        // SAFETY: the context is kept alive by `target`.
+        let status = unsafe { flush(self.target.ctx) };
+        if status != 0 {
+            return Err(io::Error::other(format!(
+                "host flush callback failed with status {status}"
+            )));
+        }
+        Ok(())
+    }
+}
+
+impl VortexWrite for CallbackWrite {
+    async fn write_all<B: IoBuf>(&mut self, buffer: B) -> io::Result<B> {
+        self.write_slice(buffer.as_slice())?;
+        Ok(buffer)
+    }
+
+    async fn flush(&mut self) -> io::Result<()> {
+        self.flush_host()
+    }
+
+    async fn shutdown(&mut self) -> io::Result<()> {
+        // The host owns its stream and closes it itself, so this only flushes.
+        self.flush_host()
+    }
+}
+
+/// A sink that writes a Vortex file through host callbacks.
+///
+/// Mirrors `vx_array_sink`: pushed arrays go through a channel feeding a 
write task on the session's
+/// runtime, and errors surface when the sink is closed. This is a separate 
type because
+/// `vx_array_sink`'s fields are private to its own module.
+pub struct vx_callback_sink {
+    sink: Sender<VortexResult<ArrayRef>>,
+    writer: Task<VortexResult<WriteSummary>>,
+    dtype: DType,
+}
+
+unsafe fn callback_sink_open(
+    session: *const vx_session,
+    callbacks: vx_output_callbacks,
+    dtype: *const vx_dtype,
+) -> VortexResult<*mut vx_callback_sink> {
+    // Take ownership of the host context before anything else, so an early 
error below still
+    // releases it exactly once when this `Arc` is dropped.
+    let target = Arc::new(CallbackTarget {
+        ctx: callbacks.ctx,
+        release_fn: callbacks.release_fn,
+    });
+
+    vortex_ensure!(!session.is_null());
+    vortex_ensure!(!dtype.is_null());
+    let write_fn = callbacks
+        .write_fn
+        .ok_or_else(|| vortex_err!("vx_output_callbacks.write_fn is 
required"))?;
+
+    let session = vx_session::as_ref(session).clone();
+    let dtype = vx_dtype::as_ref(dtype).clone();
+    let write = CallbackWrite {
+        target,
+        write_fn,
+        flush_fn: callbacks.flush_fn,
+    };
+
+    // The channel size matches the stock file sink.
+    let (sink, rx) = mpsc::channel(32);
+    let array_stream = ArrayStreamAdapter::new(dtype.clone(), 
rx.into_stream());

Review Comment:
   Java JNI writer uses `WRITE_CHANNEL_CAPACITY = 4`, and the comment 
explicitly says the queue is kept small so the Java thread producing batches 
can feel backpressure in time. Could you explain why there is this 
inconsistency?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to