suxiaogang223 commented on code in PR #67677:
URL: https://github.com/apache/doris/pull/67677#discussion_r4011241824


##########
fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonCppWriteSupport.java:
##########
@@ -0,0 +1,182 @@
+// 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.
+
+package org.apache.doris.datasource.paimon;
+
+import org.apache.doris.common.util.LocationPath;
+import org.apache.doris.datasource.property.storage.StorageProperties;
+import org.apache.doris.thrift.TFileType;
+import org.apache.doris.thrift.TPaimonCppColumn;
+import org.apache.doris.thrift.TPaimonCppStorageDescriptor;
+import org.apache.doris.thrift.TPaimonCppWriteDescriptor;
+import org.apache.doris.thrift.TPaimonWriteMode;
+
+import com.google.common.collect.ImmutableSet;
+import lombok.AccessLevel;
+import lombok.Getter;
+import lombok.RequiredArgsConstructor;
+import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.types.DataField;
+
+import java.net.URI;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+
+/** Pure, pre-writer capability decision. Never retry a failed native writer 
through JNI. */
+public final class PaimonCppWriteSupport {
+    private static final Set<String> OPTIONS = ImmutableSet.of(
+            "bucket", "file.format", "manifest.format", "write-only", "path", 
"owner",
+            "file.compression", "target-file-size", "write-buffer-size",
+            "page-size", "commit.force-create-snapshot");
+    private static final Set<String> TYPES = ImmutableSet.of(

Review Comment:
   Addressed in acf31a17a28, with the supported boundary now documented in the 
PR description. The native path obtains and pins the complete SDK Arrow schema 
and validates the exact `struct<value: binary not null, metadata: binary not 
null>` layout plus Paimon metadata. Ordinary unshredded Parquet VARIANT is 
enabled and has Doris/Spark coverage for top-level, multiple columns, 
ROW/ARRAY/MAP and deep nesting. During integration testing, configured native 
shredding wrote successfully but the Java reader failed because SDK-generated 
physical children lacked field IDs. Therefore explicit and inferred shredding 
remain planning-time JNI fallbacks rather than being enabled prematurely; 
`inferShreddingSchema=false` is accepted for ordinary native VARIANT. ADAPTIVE 
remains SDK-internal and requires no separate Doris state or gate once the 
common shredded-file interoperability issue is fixed. No third-party source is 
changed in this PR.



##########
be/src/exec/sink/writer/paimon/cpp_paimon_write_backend.cpp:
##########
@@ -0,0 +1,442 @@
+// 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 "exec/sink/writer/paimon/cpp_paimon_write_backend.h"
+
+#include <limits>
+
+#include "common/exception.h"
+
+#ifdef USE_PAIMON_CPP
+#include <arrow/api.h>
+#include <arrow/c/bridge.h>
+#include <paimon/commit_message.h>
+#include <paimon/file_store_write.h>
+#include <paimon/memory/memory_pool.h>
+#include <paimon/record_batch.h>
+#include <paimon/write_context.h>
+
+#include <algorithm>
+#include <atomic>
+#include <cstring>
+
+#include "common/config.h"
+#include "common/logging.h"
+#include "core/allocator.h"
+#include "exec/sink/writer/paimon/doris_paimon_file_system.h"
+#include "exec/sink/writer/paimon/paimon_resource_context.h"
+#include "format/arrow/arrow_block_convertor.h"
+#include "format/parquet/arrow_memory_pool.h"
+#include "io/file_factory.h"
+#include "runtime/memory/mem_tracker_limiter.h"
+#include "runtime/query_context.h"
+#include "runtime/runtime_state.h"
+#include "runtime/thread_context.h"
+#include "util/debug_points.h"
+#include "util/defer_op.h"
+#endif
+
+namespace doris {
+
+Status frame_paimon_cpp_commit(const std::string& data, int32_t version,
+                               TPaimonCommitMessage* message) {
+    // Same per-frame bound as the JNI codec. This does NOT bound SDK metadata 
accumulation.
+    constexpr size_t max_payload = 8 * 1024 * 1024;
+    if (message == nullptr || version < 0 || data.size() > max_payload - 12) {
+        return Status::InvalidArgument("Invalid or oversized Paimon commit 
payload");
+    }
+    std::string framed("DPCM");
+    auto append_int = [&](uint32_t value) {
+        for (int shift = 24; shift >= 0; shift -= 8) {
+            framed.push_back(static_cast<char>((value >> shift) & 0xff));
+        }
+    };
+    append_int(static_cast<uint32_t>(version));
+    append_int(static_cast<uint32_t>(data.size()));
+    framed.append(data);
+    message->__set_payload(std::move(framed));
+    return Status::OK();
+}
+
+#ifdef USE_PAIMON_CPP
+namespace {
+
+Status sdk_status(const paimon::Status& status) {
+    if (status.ok()) {
+        return Status::OK();
+    }
+    if (status.IsOutOfMemory()) {
+        return Status::MemoryLimitExceeded(status.ToString());
+    }
+    return Status::InternalError("Paimon native: {}", status.ToString());
+}
+
+class QueryMemoryPool final : public paimon::MemoryPool {
+public:
+    QueryMemoryPool(std::shared_ptr<ResourceContext> context, uint64_t limit)
+            : _context(std::move(context)), _limit(limit) {}
+
+    void* Malloc(uint64_t size, uint64_t alignment = 0) override {
+        if (size > static_cast<uint64_t>(std::numeric_limits<int64_t>::max()) 
||
+            (alignment != 0 && (alignment & (alignment - 1)) != 0)) {
+            throw std::bad_alloc();
+        }
+        const uint64_t charged = std::max<uint64_t>(size, 1);
+        uint64_t used = _used.load();
+        do {
+            if (used > _limit || charged > _limit - used) {
+                throw std::bad_alloc();
+            }
+        } while (!_used.compare_exchange_weak(used, used + charged));
+        void* ptr = nullptr;
+        Defer rollback {[&] {
+            if (!ptr) _used.fetch_sub(charged);
+        }};
+        try {
+            ptr = with_paimon_resource_context(_context, [&] {
+                enable_thread_catch_bad_alloc++;
+                Defer restore {[&] { enable_thread_catch_bad_alloc--; }};
+                return _allocator.alloc(charged, std::max<uint64_t>(alignment, 
64));
+            });
+            if (ptr == nullptr) {
+                throw std::bad_alloc();
+            }
+            auto peak = _peak.load();
+            while (peak < used + charged && !_peak.compare_exchange_weak(peak, 
used + charged)) {
+            }
+            return ptr;
+        } catch (const doris::Exception& e) {
+            // Paimon's ArrowMemPoolAdaptor catches std::bad_alloc, not Doris 
exceptions.
+            // Keep allocation failure inside the SDK's Status-based 
error/cleanup path.
+            if (e.code() == ErrorCode::MEM_ALLOC_FAILED ||
+                e.code() == ErrorCode::MEM_LIMIT_EXCEEDED ||
+                e.code() == ErrorCode::BUFFER_ALLOCATION_FAILED ||
+                e.code() == ErrorCode::QUERY_MEMORY_EXCEEDED ||
+                e.code() == ErrorCode::WORKLOAD_GROUP_MEMORY_EXCEEDED ||
+                e.code() == ErrorCode::PROCESS_MEMORY_EXCEEDED) {
+                throw std::bad_alloc();
+            }
+            throw;
+        }
+    }
+
+    void* Realloc(void* ptr, size_t old_size, size_t new_size, uint64_t 
alignment = 0) override {
+        // Allocate-copy-free deliberately accounts for both live buffers at 
the expansion peak.
+        // A failed allocation leaves ptr, its data and its accounting 
unchanged.
+        void* replacement = Malloc(new_size, alignment);
+        if (ptr != nullptr) {
+            std::memcpy(replacement, ptr, std::min(old_size, new_size));
+            Free(ptr, old_size);
+        }
+        return replacement;
+    }
+
+    void Free(void* ptr, uint64_t size) override {
+        if (ptr != nullptr) {
+            with_paimon_resource_context(
+                    _context, [&] { _allocator.free(ptr, 
std::max<uint64_t>(size, 1)); });
+            _used.fetch_sub(std::max<uint64_t>(size, 1));
+        }
+    }
+    uint64_t CurrentUsage() const override { return _used.load(); }
+    uint64_t MaxMemoryUsage() const override { return _peak.load(); }
+
+private:
+    std::shared_ptr<ResourceContext> _context;
+    uint64_t _limit;
+    Allocator<false> _allocator;
+    std::atomic<uint64_t> _used {0};
+    std::atomic<uint64_t> _peak {0};
+};
+
+// Conversion buffers can outlive Write() and be freed by an SDK worker thread.
+class QueryArrowPool final : public ArrowMemoryPool<> {
+public:
+    explicit QueryArrowPool(std::shared_ptr<ResourceContext> context)
+            : _context(std::move(context)) {}
+    arrow::Status Allocate(int64_t size, int64_t alignment, uint8_t** out) 
override {
+        return with_paimon_resource_context(
+                _context, [&] { return ArrowMemoryPool<>::Allocate(size, 
alignment, out); });
+    }
+    arrow::Status Reallocate(int64_t old_size, int64_t new_size, int64_t 
alignment,
+                             uint8_t** ptr) override {
+        return with_paimon_resource_context(_context, [&] {
+            return ArrowMemoryPool<>::Reallocate(old_size, new_size, 
alignment, ptr);
+        });
+    }
+    void Free(uint8_t* ptr, int64_t size, int64_t alignment) override {
+        with_paimon_resource_context(_context,
+                                     [&] { ArrowMemoryPool<>::Free(ptr, size, 
alignment); });
+    }
+    std::string backend_name() const override { return 
"DorisPaimonConversion"; }
+
+private:
+    std::shared_ptr<ResourceContext> _context;
+};
+
+struct ExportOwner {
+    ArrowArray array;
+    std::shared_ptr<QueryArrowPool> pool;
+    static void release(ArrowArray* exported) {
+        auto* owner = static_cast<ExportOwner*>(exported->private_data);
+        exported->release = nullptr;
+        if (owner->array.release != nullptr) {
+            owner->array.release(&owner->array);
+        }
+        delete owner; // pool outlives all buffer destructors invoked by 
original callback
+    }
+};
+
+std::shared_ptr<arrow::DataType> arrow_type(const std::string& type) {
+    if (type == "BOOLEAN") return arrow::boolean();
+    if (type == "TINYINT") return arrow::int8();
+    if (type == "SMALLINT") return arrow::int16();
+    if (type == "INTEGER") return arrow::int32();
+    if (type == "BIGINT") return arrow::int64();
+    if (type == "FLOAT") return arrow::float32();
+    if (type == "DOUBLE") return arrow::float64();
+    if (type == "VARCHAR") return arrow::utf8();
+    if (type == "VARBINARY") return arrow::binary();
+    return nullptr;
+}
+
+} // namespace
+
+std::shared_ptr<paimon::MemoryPool> make_paimon_query_memory_pool(
+        std::shared_ptr<ResourceContext> context, uint64_t limit) {
+    return std::make_shared<QueryMemoryPool>(std::move(context), limit);
+}
+
+class CppPaimonWriteBackend::Impl {
+public:
+    Status open(const TPaimonTableSink& sink, RuntimeState* state, 
RuntimeProfile* profile) {
+        if (!sink.__isset.cpp_descriptor || !sink.__isset.write_mode ||
+            sink.write_mode != TPaimonWriteMode::APPEND || 
!sink.__isset.commit_user ||
+            sink.commit_user.empty() || state->get_query_ctx() == nullptr ||
+            state->query_mem_tracker() == nullptr) {
+            return Status::InvalidArgument("Incomplete native Paimon write 
description");
+        }
+        const auto& desc = sink.cpp_descriptor;
+        if (desc.columns.empty() || !sink.__isset.column_names ||
+            desc.columns.size() != sink.column_names.size()) {
+            return Status::InvalidArgument("Native Paimon column order is 
missing");
+        }
+        // Reject configuration mismatch rather than silently falling back 
after dispatch.
+        if (!desc.__isset.storage || desc.root_path.empty() || 
desc.storage.root_path.empty() ||
+            (desc.storage.file_type != TFileType::FILE_LOCAL &&
+             desc.storage.file_type != TFileType::FILE_S3)) {
+            return Status::NotSupported("Missing or unsupported Doris Paimon 
storage descriptor");
+        }
+        const auto& root_path = desc.root_path;
+        const auto& storage = desc.storage;
+        // FE selects supported write options and normalizes storage 
locations. The filesystem
+        // adapter owns logical-to-storage path mapping; the backend only 
consumes the descriptor.
+        arrow::FieldVector fields;
+        for (size_t i = 0; i < desc.columns.size(); ++i) {
+            const auto& column = desc.columns[i];
+            auto type = arrow_type(column.type);
+            if (!type) {
+                return Status::NotSupported("Unsupported native Paimon type 
{}", column.type);
+            }
+            fields.push_back(arrow::field(sink.column_names[i], type, 
column.nullable));
+        }
+        _schema = arrow::schema(fields);
+        auto context = state->get_query_ctx()->resource_ctx();
+        int64_t limit = config::paimon_cpp_writer_memory_limit_bytes;
+        const auto query_limit = state->query_mem_tracker()->limit();
+        if (query_limit > 0) {
+            limit = std::min(limit, query_limit / std::max(1, 
state->task_num()));
+        }
+        if (limit <= 0) return Status::MemoryLimitExceeded("No Paimon native 
writer memory budget");
+        _pool = make_paimon_query_memory_pool(context, limit);
+        _arrow_pool = std::make_shared<QueryArrowPool>(context);
+        COUNTER_SET(ADD_COUNTER(profile, "PaimonSdkPoolLimit", TUnit::BYTES), 
limit);
+        _sdk_pool_peak = ADD_COUNTER(profile, "PaimonSdkPoolPeak", 
TUnit::BYTES);
+        _conversion_peak = ADD_COUNTER(profile, "PaimonArrowConversionPeak", 
TUnit::BYTES);
+        profile->add_info_string("PaimonMemoryScope",
+                                 "SDK pool limit excludes conversion, IO and 
non-pool allocations");
+        io::FSPropertiesRef fs_properties(storage.file_type);
+        fs_properties.properties = &storage.properties;
+        io::FileDescription file_description {.path = storage.root_path};
+        auto fs = with_paimon_resource_context(
+                context, [&] { return FileFactory::create_fs(fs_properties, 
file_description); });
+        if (!fs.has_value()) return fs.error();
+        _filesystem = 
std::make_shared<DorisPaimonFileSystem>(std::move(fs.value()), root_path,
+                                                              
storage.root_path, context);
+        auto options = desc.options;
+        options["doris.expected-schema-id"] = std::to_string(desc.schema_id);
+        paimon::WriteContextBuilder builder(root_path, sink.commit_user);
+        auto ctx = builder.SetOptions(options)
+                           .WithFileSystem(_filesystem)
+                           .WithMemoryPool(_pool)
+                           .WithWriteSchema(sink.column_names)
+                           .Finish();
+        if (!ctx.ok()) return sdk_status(ctx.status());
+        auto writer = paimon::FileStoreWrite::Create(std::move(ctx).value());
+        if (!writer.ok()) return sdk_status(writer.status());
+        _sdk = std::move(writer).value();
+        profile->add_info_string("PaimonWriteBackend", "CPP (experimental 
v1)");
+        return Status::OK();
+    }
+
+    Status write(RuntimeState* state, Block& block) {
+        if (block.rows() == 0) return Status::OK();
+        DCHECK(_sdk);
+        if (state->is_cancelled()) return Status::Cancelled("Paimon native 
write cancelled");
+        std::shared_ptr<arrow::RecordBatch> batch;
+        RETURN_IF_ERROR(convert_to_arrow_batch(block, _schema, 
_arrow_pool.get(), &batch,

Review Comment:
   Fixed in acf31a17a28. Before the schema-less ArrowArray export, the backend 
now validates both the pinned batch schema (including metadata) and every 
actual column array type against the SDK schema, then runs RecordBatch 
validation. This rejects utf8/large_utf8 and binary/large_binary offset-width 
mismatches, including nested layouts, instead of allowing 64-bit offsets to be 
interpreted as 32-bit offsets. `RejectsLargeOffsetsBeforeCDataExport` covers 
top-level string/binary and nested Variant-like struct mismatches without 
requiring a 2 GiB allocation.



-- 
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]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to