github-actions[bot] commented on code in PR #67395:
URL: https://github.com/apache/doris/pull/67395#discussion_r4213990946


##########
fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonWritePlanProvider.java:
##########
@@ -0,0 +1,535 @@
+// 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.connector.paimon;
+
+import org.apache.doris.connector.spi.ConnectorColumn;
+import org.apache.doris.connector.spi.ConnectorContext;
+import org.apache.doris.connector.spi.ConnectorSession;
+import org.apache.doris.connector.spi.DorisConnectorException;
+import org.apache.doris.connector.spi.handle.ConnectorTableHandle;
+import org.apache.doris.connector.spi.handle.ConnectorTransaction;
+import org.apache.doris.connector.spi.handle.ConnectorWriteHandle;
+import org.apache.doris.connector.spi.handle.WriteOperation;
+import org.apache.doris.connector.spi.write.ConnectorChangelogMode;
+import org.apache.doris.connector.spi.write.ConnectorRowChangeStyle;
+import org.apache.doris.connector.spi.write.ConnectorRowLevelDmlRequest;
+import org.apache.doris.connector.spi.write.ConnectorSinkPlan;
+import org.apache.doris.connector.spi.write.ConnectorWriteDistribution;
+import org.apache.doris.connector.spi.write.ConnectorWritePlanProvider;
+import org.apache.doris.filesystem.properties.StorageProperties;
+import org.apache.doris.thrift.TDataSink;
+import org.apache.doris.thrift.TDataSinkType;
+import org.apache.doris.thrift.TPaimonTableSink;
+import org.apache.doris.thrift.TPaimonWriteBackendType;
+import org.apache.doris.thrift.TPaimonWriteMode;
+
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.catalog.Catalog;
+import org.apache.paimon.catalog.Identifier;
+import org.apache.paimon.options.CatalogOptions;
+import org.apache.paimon.options.Options;
+import org.apache.paimon.schema.TableSchema;
+import org.apache.paimon.table.BucketMode;
+import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.table.Table;
+import org.apache.paimon.types.DataField;
+import org.apache.paimon.types.DataTypeRoot;
+
+import java.net.URI;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.EnumSet;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.Set;
+import java.util.TreeMap;
+import java.util.stream.Collectors;
+
+/** Builds the JNI-backed Paimon sink and binds it to the active connector 
transaction. */
+public class PaimonWritePlanProvider implements ConnectorWritePlanProvider {
+    // Keep in sync with Config.PAIMON_WRITE_MIN_BE_EXEC_VERSION. Connector 
modules cannot depend
+    // on fe-common, so the wire-protocol boundary is repeated at the provider 
admission gate.
+    static final int MIN_BE_EXEC_VERSION = 16;
+
+    static final String ROW_KIND_COLUMN = "__DORIS_PAIMON_ROW_KIND__";
+    static final byte INSERT_OPERATION = 0;
+    static final byte UPDATE_OPERATION = 1;
+    static final byte DELETE_OPERATION = 2;
+
+    private final PaimonCatalogProperties catalogProperties;
+    private final PaimonCatalogOps catalogOps;
+    private final ConnectorContext context;
+
+    PaimonWritePlanProvider(PaimonCatalogProperties catalogProperties,
+            PaimonCatalogOps catalogOps, ConnectorContext context) {
+        this.catalogProperties = catalogProperties;
+        this.catalogOps = catalogOps;
+        this.context = context;
+    }
+
+    @Override
+    public Optional<List<ConnectorColumn>> getWriteColumns(ConnectorSession 
session,
+            ConnectorTableHandle tableHandle, Optional<String> branchName) {
+        FileStoreTable table = resolveTable((PaimonTableHandle) tableHandle);
+        PaimonTypeMapping.Options options = new PaimonTypeMapping.Options(
+                catalogProperties.isEnableMappingVarbinary(),
+                catalogProperties.isEnableMappingTimestampTz());
+        return Optional.of(mapWriteColumns(table.schema().fields(), 
table.primaryKeys(), options));
+    }
+
+    static List<ConnectorColumn> mapWriteColumns(List<DataField> fields, 
List<String> primaryKeys,
+            PaimonTypeMapping.Options options) {
+        Set<String> keyNames = new 
java.util.TreeSet<>(String.CASE_INSENSITIVE_ORDER);
+        keyNames.addAll(primaryKeys);
+        List<ConnectorColumn> columns = new ArrayList<>(fields.size());
+        for (DataField field : fields) {
+            ConnectorColumn column = new ConnectorColumn(
+                    field.name(), 
PaimonTypeMapping.toConnectorType(field.type(), options),
+                    field.description(), field.type().isNullable(), 
field.defaultValue(),
+                    keyNames.contains(field.name())).withUniqueId(field.id());
+            if (field.defaultValue() != null) {
+                column = 
column.withDefaultValueSql(toDorisDefaultValueSql(field));
+            }
+            if (field.type().getTypeRoot() == 
DataTypeRoot.TIMESTAMP_WITH_LOCAL_TIME_ZONE) {
+                column = column.withTimeZone();
+            }
+            columns.add(column);
+        }
+        return Collections.unmodifiableList(columns);
+    }
+
+    private static String toDorisDefaultValueSql(DataField field) {
+        String value = field.defaultValue();
+        switch (field.type().getTypeRoot()) {
+            case CHAR:
+            case VARCHAR:
+            case BINARY:
+            case VARBINARY:
+            case DATE:
+            case TIME_WITHOUT_TIME_ZONE:
+            case TIMESTAMP_WITHOUT_TIME_ZONE:
+            case TIMESTAMP_WITH_LOCAL_TIME_ZONE:
+            case VARIANT:
+            case BLOB:
+                if ((value.startsWith("'") && value.endsWith("'"))

Review Comment:
   [P2] Always quote raw Paimon string defaults when building Doris SQL. 
`field.defaultValue()` is a value, so a VARCHAR default containing the literal 
three characters `'x'` enters this branch and is returned as SQL `'x'`. 
`BindSink` then parses it as `x` when the INSERT omits the column, silently 
changing the stored value. Escape the raw value even when its boundary 
characters are quotes, and cover an omitted-column insert with such a default.



##########
be/src/exec/spill/spill_file_manager.cpp:
##########
@@ -31,14 +32,105 @@
 #include "exec/spill/spill_file.h"
 #include "io/fs/file_system.h"
 #include "io/fs/local_file_system.h"
+#include "runtime/query_context.h"
 #include "storage/olap_define.h"
 #include "util/debug_points.h"
 #include "util/parse_util.h"
 #include "util/pretty_printer.h"
 #include "util/time.h"
+#include "util/uid_util.h"
 
 namespace doris {
 
+ExternalSpillSession::ExternalSpillSession(SpillFileManager* manager, 
QueryContext* query_context,
+                                           std::string relative_path)
+        : _manager(manager),
+          _query_context(query_context->weak_from_this()),
+          _resource_context(query_context->resource_ctx()),
+          _query_id(print_id(query_context->query_id())),
+          _relative_path(std::move(relative_path)) {
+    DCHECK(_manager != nullptr);
+    DCHECK(!_query_context.expired());
+    DCHECK(_resource_context != nullptr);
+}
+
+ExternalSpillSession::~ExternalSpillSession() {
+    _manager->_release_external_spill_session(this);
+}
+
+Status ExternalSpillSession::get_paths(std::vector<std::string>* paths) {
+    if (paths == nullptr) {
+        return Status::InvalidArgument("External spill paths output must not 
be null");
+    }
+    std::lock_guard lock(_mutex);
+    if (_data_dir == nullptr) {
+        RETURN_IF_ERROR(_manager->_initialize_external_spill_session(this));
+    }
+    *paths = {_path};
+    return Status::OK();
+}
+
+bool ExternalSpillSession::_contains(const std::string& path) const {
+    return path == _path ||
+           (path.size() > _path.size() && path.starts_with(_path) && 
path[_path.size()] == '/');
+}
+
+Status ExternalSpillSession::reserve(const std::string& path, int64_t bytes) {
+    if (bytes <= 0) {
+        return Status::InvalidArgument("External spill reservation must be 
positive: {}", bytes);
+    }
+
+    std::lock_guard lock(_mutex);
+    if (_data_dir == nullptr || !_contains(path)) {
+        return Status::InvalidArgument("External spill path is not managed by 
Doris: {}", path);
+    }
+    if (bytes > std::numeric_limits<int64_t>::max() - _accounted_bytes) {
+        return Status::InvalidArgument("External spill reservation overflows: 
bytes={}", bytes);
+    }
+    if (_data_dir->reach_capacity_limit(bytes)) {

Review Comment:
   [P2] Make root spill admission and charging atomic for concurrent Paimon 
sessions. This check and `update_spill_data_usage(bytes)` take the 
`SpillDataDir` mutex in separate calls, while `reserve` locks only its own 
session. With 100 bytes free, two sessions can each pass a 60-byte check before 
either updates usage, then both write and charge 120 bytes. Reserve under one 
root lock so tracked buffer writes obey the spill limit even when sink tasks 
run in parallel.



##########
fe/fe-connector/fe-connector-paimon/src/main/java/org/apache/doris/connector/paimon/PaimonConnectorTransaction.java:
##########
@@ -0,0 +1,326 @@
+// 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.connector.paimon;
+
+import org.apache.doris.connector.spi.ConnectorContext;
+import org.apache.doris.connector.spi.DorisConnectorException;
+import org.apache.doris.connector.spi.handle.ConnectorTransaction;
+import org.apache.doris.thrift.TPaimonCommitMessage;
+
+import com.google.common.base.Preconditions;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+import org.apache.paimon.io.DataInputDeserializer;
+import org.apache.paimon.table.sink.CommitMessage;
+import org.apache.paimon.table.sink.CommitMessageSerializer;
+import org.apache.paimon.table.sink.InnerTableCommit;
+import org.apache.thrift.TDeserializer;
+import org.apache.thrift.TException;
+import org.apache.thrift.protocol.TBinaryProtocol;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+
+/** Connector-SPI transaction for one Paimon write statement. */
+public class PaimonConnectorTransaction implements ConnectorTransaction {
+
+    private static final Logger LOG = 
LogManager.getLogger(PaimonConnectorTransaction.class);
+    private static final int COMMIT_HEADER_SIZE = 12;
+    private static final byte[] COMMIT_MAGIC = new byte[] {'D', 'P', 'C', 'M'};
+
+    enum CommitState {
+        PREPARED,
+        COMMITTING,
+        COMMITTED,
+        OUTCOME_UNKNOWN
+    }
+
+    private final long transactionId;
+    private final String commitUser;
+    private final ConnectorContext context;
+    private final List<byte[]> commitPayloads = new ArrayList<>();
+    private final Set<CommitPayloadKey> commitPayloadSet = new HashSet<>();
+    private PaimonWriteBinding binding;
+    private CommitState state = CommitState.PREPARED;
+
+    PaimonConnectorTransaction(long transactionId, ConnectorContext context) {
+        Preconditions.checkArgument(transactionId > 0, "Paimon transaction id 
must be positive");
+        this.transactionId = transactionId;
+        this.context = Preconditions.checkNotNull(context, "Paimon connector 
context must not be null");
+        this.commitUser = commitUser(context.getClusterId(), transactionId);
+    }
+
+    synchronized void bind(PaimonWriteBinding writeBinding) {
+        Preconditions.checkNotNull(writeBinding, "Paimon write binding must 
not be null");
+        Preconditions.checkState(binding == null, "Paimon transaction is 
already bound");
+        Preconditions.checkState(state == CommitState.PREPARED,
+                "Paimon transaction can only be bound while prepared");
+        binding = writeBinding;
+    }
+
+    String getCommitUser() {
+        return commitUser;
+    }
+
+    static String commitUser(int clusterId, long transactionId) {
+        return "doris_cluster_" + clusterId + "_txn_" + transactionId;
+    }
+
+    @Override
+    public long getTransactionId() {
+        return transactionId;
+    }
+
+    @Override
+    public void addCommitData(byte[] commitFragment) {
+        TPaimonCommitMessage message = new TPaimonCommitMessage();
+        try {
+            new TDeserializer(new 
TBinaryProtocol.Factory()).deserialize(message, commitFragment);
+        } catch (TException e) {
+            throw new DorisConnectorException("Failed to deserialize Paimon 
commit message", e);
+        }
+        Preconditions.checkState(message.isSetPayload(), "Paimon commit 
message payload is missing");
+        byte[] payload = message.getPayload();
+        Preconditions.checkState(payload.length > 0, "Paimon commit message 
payload is empty");
+        synchronized (this) {
+            if (commitPayloadSet.add(new CommitPayloadKey(payload))) {
+                commitPayloads.add(payload);
+            }
+        }
+    }
+
+    @Override
+    public void commit() {
+        PaimonWriteBinding writeBinding = requireBinding();
+        List<byte[]> payloads = snapshotPayloads();
+        if (payloads.isEmpty() && !writeBinding.isOverwrite()) {
+            markPreparedTransactionCommitted();
+            return;
+        }
+        try {
+            List<CommitMessage> messages = deserializePayloads(payloads);
+            doCommitWithReconciliation(writeBinding, messages);
+        } catch (Exception e) {
+            throw new DorisConnectorException("Failed to commit Paimon 
transaction on FE", e);
+        }
+    }
+
+    @Override
+    public void rollback() {
+        CommitState current = getState();
+        if (current == CommitState.COMMITTED) {
+            return;
+        }
+        if (current == CommitState.COMMITTING || current == 
CommitState.OUTCOME_UNKNOWN) {
+            LOG.warn("Skip rollback for Paimon transaction in state {}, 
txnId={}, table={}",
+                    current, transactionId, tableName());
+            return;
+        }
+        List<byte[]> payloads = snapshotPayloads();
+        if (payloads.isEmpty()) {
+            return;
+        }
+        try {
+            PaimonWriteBinding writeBinding = requireBinding();
+            List<CommitMessage> messages = deserializePayloads(payloads);
+            context.executeAuthenticated(() -> {
+                try (InnerTableCommit committer = 
writeBinding.getTable().newCommit(commitUser)) {
+                    committer.abort(messages);
+                }
+                return null;
+            });
+        } catch (Exception e) {
+            LOG.warn("Failed to rollback Paimon transaction, txnId={}, 
table={}",
+                    transactionId, tableName(), e);
+        }
+    }
+
+    @Override
+    public void close() {
+        // The Paimon committer is scoped to commit/rollback and closed there.
+    }
+
+    @Override
+    public String profileLabel() {
+        return "PAIMON";
+    }
+
+    private void doCommitWithReconciliation(PaimonWriteBinding writeBinding,
+            List<CommitMessage> messages) throws Exception {
+        Exception firstFailure;
+        try {
+            doCommit(writeBinding, messages);
+            return;
+        } catch (Exception e) {
+            if (getState() == CommitState.COMMITTED) {
+                return;
+            }
+            if (getState() == CommitState.PREPARED) {

Review Comment:
   [P2] Keep abort ownership when commit fails before `markCommitting()`. If 
`newCommit(commitUser)` or authenticated committer setup throws, this branch 
rethrows in `PREPARED` state without aborting the staged messages. 
`PluginDrivenTransactionManager.commit()` has already removed the transaction, 
so the executor's later `rollback(txnId)` is a no-op and the files remain 
orphaned. Abort them on this definite pre-commit failure, or retain the 
transaction for rollback; test a committer-creation failure after BE prepare.



##########
fe/be-java-extensions/paimon-scanner/src/main/java/org/apache/doris/paimon/PaimonJniWriter.java:
##########
@@ -0,0 +1,791 @@
+// 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.paimon;
+
+import org.apache.doris.jni.spi.JniWriter;
+import org.apache.doris.jni.spi.ThreadContextClassLoader;
+import org.apache.doris.jni.spi.vec.VectorTable;
+import org.apache.doris.kerberos.PreExecutionAuthenticator;
+import org.apache.doris.kerberos.PreExecutionAuthenticatorCache;
+
+import org.apache.arrow.c.ArrowArray;
+import org.apache.arrow.c.ArrowSchema;
+import org.apache.arrow.c.CDataDictionaryProvider;
+import org.apache.arrow.c.Data;
+import org.apache.arrow.memory.BufferAllocator;
+import org.apache.arrow.memory.RootAllocator;
+import org.apache.arrow.vector.VectorSchemaRoot;
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.crosspartition.IndexBootstrap;
+import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.index.BucketAssigner;
+import org.apache.paimon.index.HashBucketAssigner;
+import org.apache.paimon.index.SimpleHashBucketAssigner;
+import org.apache.paimon.memory.MemoryPoolFactory;
+import org.apache.paimon.table.BucketMode;
+import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.table.sink.CommitMessage;
+import org.apache.paimon.table.sink.InnerTableCommit;
+import org.apache.paimon.table.sink.PartitionKeyExtractor;
+import org.apache.paimon.table.sink.RowPartitionKeyExtractor;
+import org.apache.paimon.table.sink.SinkRecord;
+import org.apache.paimon.table.sink.TableWriteImpl;
+import org.apache.paimon.utils.ExecutorThreadFactory;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.io.IOException;
+import java.nio.ByteBuffer;
+import java.time.ZoneId;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashSet;
+import java.util.Iterator;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Set;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
+
+/**
+ * JNI entry point for Paimon write operations.
+ *
+ * <p>Called from C++ ({@code JniPaimonWriter}) via JNI. One instance per BE 
pipeline
+ * fragment (one per {@code PaimonTableWriter}). Data path:
+ *
+ * <pre>
+ *   C++ Block → Arrow RecordBatch → Arrow C Data Interface
+ *   → PaimonJniWriter.writeArrow(arrayAddress, schemaAddress)
+ *   → zero-copy VectorSchemaRoot view
+ *   → PaimonArrowBatchAdapter (Arrow-backed Paimon columnar row)
+ *   → PaimonWriteSchema.tableRow() (canonical table-schema order)
+ *   → Paimon SDK bucket assignment and table write
+ * </pre>
+ *
+ * <p>Commit path:
+ *
+ * <pre>
+ *   PaimonTableWriter::close() → JNI → PaimonJniWriter.prepareCommit()
+ *   → TableWriteImpl.prepareCommit()
+ *   → PaimonCommitCodec.encode() → DPCM-framed byte[][]
+ *   → C++ collects TPaimonCommitMessage[] → RPC to FE → PaimonTransaction
+ * </pre>
+ */
+public class PaimonJniWriter extends JniWriter {
+    private static final Logger LOG = 
LoggerFactory.getLogger(PaimonJniWriter.class);
+    private static final int APPEND_ONLY_WRITER_MIN_PAGES = 1;
+    private static final int MERGE_TREE_WRITER_MIN_PAGES = 3;
+    private static final long COMPACTION_CLOSE_TIMEOUT_SECONDS = 60;
+
+    private final PaimonCommitCodec commitCodec = new PaimonCommitCodec();
+
+    private BufferAllocator allocator;
+    private PreExecutionAuthenticator preExecutionAuthenticator;
+    private PaimonArrowBatchAdapter arrowAdapter;
+
+    private PaimonWriteSchema writeSchema;
+    private FileStoreTable table;
+    private TableWriteImpl<?> writer;
+    private DorisIOManager ioManager;
+    private ExecutorService compactionExecutor;
+    private long commitIdentifier;
+    private String commitUser;
+    private BucketMode bucketMode;
+    private BucketAssigner hashBucketAssigner;
+    private PartitionKeyExtractor<InternalRow> dynamicBucketExtractor;
+    private GlobalIndexAssigner globalIndexAssigner;
+    private boolean fullCompactionChangelog;
+    private final Set<PartitionBucket> fullCompactionBuckets = new HashSet<>();
+    private List<CommitMessage> preparedCommitMessages = 
Collections.emptyList();
+    private boolean sdkCloseFailed;
+
+    public PaimonJniWriter() {
+        this(0, Collections.emptyMap());
+    }
+
+    public PaimonJniWriter(int batchSize, Map<String, String> params) {
+        super(batchSize, params);
+        // Imported C Data vectors reference Doris-owned buffers; this 
allocator owns only Arrow's
+        // Java-side views and metadata. Physical buffers remain charged to 
the C++ Arrow memory
+        // pool until the synchronous writeArrow call returns.
+        this.allocator = new RootAllocator(Long.MAX_VALUE);
+    }
+
+    // ────────────────────────────────────────────────────────────
+    // JNI entry points (called from C++)
+    // ────────────────────────────────────────────────────────────
+
+    /**
+     * Initialize the writer. Called once per BE pipeline fragment via JNI.
+     *
+     * <p>This method:
+     * <ol>
+     *   <li>Deserializes the target Paimon {@link FileStoreTable} selected by 
FE.</li>
+     *   <li>Creates a {@link PaimonWriteSchema} which normalizes Doris input
+     *       columns to the table-schema row layout.</li>
+     *   <li>Opens one Paimon SDK writer session.</li>
+     * </ol>
+     *
+     * @param serializedTable serialized Paimon table selected by FE
+     * @param hadoopConfig   filesystem and authentication configuration
+     * @param columnNames    output column names in the order produced by BE
+     * @param transactionId  Doris external transaction identifier
+     * @param commitUser     Paimon commit user shared with the FE committer
+     * @param overwrite      whether this is an overwrite write
+     * @param changelogWrite whether the first input column contains a row 
change operation
+     * @param timeZone       normalized Doris session timezone used for Paimon 
LTZ values
+     * @param nativePageMemoryLimitBytes maximum Doris-managed Paimon page 
memory
+     * @param nativeMemoryManager opaque BE manager used to allocate tracked 
native pages
+     * @param nativeSpillSession opaque managed spill session used for 
capacity and I/O accounting
+     */
+    public void open(String serializedTable, Map<String, String> hadoopConfig,
+                     String[] columnNames, long transactionId, String 
commitUser,
+                     boolean overwrite, boolean changelogWrite, String 
timeZone,
+                     long nativePageMemoryLimitBytes, long nativeMemoryManager,
+                     long nativeSpillSession) throws Exception {
+        try (ThreadContextClassLoader ignored =
+                new ThreadContextClassLoader(getClass().getClassLoader())) {
+            if (nativePageMemoryLimitBytes <= 0) {
+                throw new IllegalArgumentException(
+                        "PaimonJniWriter requires a positive native page 
memory limit");
+            }
+            if (nativeMemoryManager == 0) {
+                throw new IllegalArgumentException(
+                        "PaimonJniWriter requires a native memory manager");
+            }
+            this.preExecutionAuthenticator = 
PreExecutionAuthenticatorCache.getAuthenticator(hadoopConfig);
+            preExecutionAuthenticator.execute(() -> {
+                try {
+                    FileStoreTable table = 
PaimonUtils.deserialize(serializedTable);
+                    LOG.info("PaimonJniWriter opening: table={}, columns={}",
+                            table.fullName(), columnNames != null ? 
columnNames.length : 0);
+                    this.commitIdentifier = transactionId;
+                    this.table = table;
+                    this.commitUser = commitUser;
+                    this.bucketMode = table.bucketMode();
+
+                    CoreOptions coreOptions = 
CoreOptions.fromMap(table.options());
+                    this.writeSchema = PaimonWriteSchema.create(
+                            table.rowType(), columnNames, changelogWrite);
+                    this.arrowAdapter = new PaimonArrowBatchAdapter(
+                            writeSchema.inputType(), ZoneId.of(timeZone), 
allocator);
+                    validateWriteColumnsForMergeEngine(
+                            columnNames.length - (changelogWrite ? 1 : 0), 
coreOptions);
+                    this.fullCompactionChangelog =
+                            !coreOptions.writeOnly()
+                                    && coreOptions.changelogProducer()
+                                    == 
CoreOptions.ChangelogProducer.FULL_COMPACTION;
+                    openFileStoreWriter(
+                            table,
+                            commitUser,
+                            overwrite,
+                            coreOptions,
+                            nativePageMemoryLimitBytes,
+                            nativeMemoryManager,
+                            nativeSpillSession);
+                    return null;
+                } catch (Throwable t) {
+                    try {
+                        closeResources();
+                    } catch (Throwable closeFailure) {
+                        t.addSuppressed(closeFailure);
+                    }
+                    throw new RuntimeException("PaimonJniWriter open failed", 
t);
+                }
+            });
+        }
+    }
+
+    /** Return the exact Arrow schema derived from the pinned Paimon input 
type. */
+    public byte[] getArrowSchema() {
+        if (arrowAdapter == null) {
+            throw new IllegalStateException("PaimonJniWriter is not open");
+        }
+        return arrowAdapter.serializedArrowSchema();
+    }
+
+    /**
+     * Import and synchronously consume one C++ Arrow RecordBatch through the 
C Data Interface.
+     *
+     * <p>The native structs are valid only for this call. Import transfers 
the ArrowArray release
+     * callbacks into Java. Closing the imported root releases the exported 
C++ buffers after Paimon
+     * has consumed every row; the native caller releases only callbacks left 
by a partial import.
+     */
+    public void writeArrow(long arrayAddress, long schemaAddress) throws 
Exception {
+        try (ThreadContextClassLoader ignored =
+                new ThreadContextClassLoader(getClass().getClassLoader())) {
+            preExecutionAuthenticator.execute(() -> {
+                try (ArrowArray array = ArrowArray.wrap(arrayAddress);
+                        ArrowSchema schema = ArrowSchema.wrap(schemaAddress);
+                        CDataDictionaryProvider dictionaries = new 
CDataDictionaryProvider();
+                        VectorSchemaRoot root = Data.importVectorSchemaRoot(
+                                allocator, array, schema, dictionaries)) {
+                    writeBatch(root);
+                    ioManager.reconcileDirectFileGrowth();
+                    return null;
+                } catch (Throwable t) {
+                    throw new RuntimeException("PaimonJniWriter C Data write 
failed", t);
+                }
+            });
+        }
+    }
+
+    /**
+     * Prepare commit: flush all in-memory data, close files, and serialize 
commit
+     * messages for the FE coordinator.
+     *
+     * <p>Flushes and collects Paimon {@link CommitMessage}s, then encodes 
them via
+     * {@link PaimonCommitCodec} into DPCM-framed byte chunks that are 
forwarded to
+     * FE through the BE.
+     *
+     * @return byte[][]  each element is a DPCM-framed serialized 
CommitMessage chunk
+     */
+    public byte[][] prepareCommit() throws Exception {
+        try (ThreadContextClassLoader ignored =
+                new ThreadContextClassLoader(getClass().getClassLoader())) {
+            return preExecutionAuthenticator.execute(() -> {
+                try {
+                    List<CommitMessage> messages = prepareCommitMessages();
+                    ioManager.reconcileDirectFileGrowth();
+                    if (messages.isEmpty()) {
+                        LOG.info("PaimonJniWriter prepareCommit: empty");
+                        return new byte[0][];
+                    }
+                    LOG.info("PaimonJniWriter prepareCommit: {} messages", 
messages.size());
+                    return commitCodec.encode(messages);
+                } catch (Throwable t) {
+                    throw new RuntimeException("PaimonJniWriter prepareCommit 
failed", t);
+                }
+            });
+        }
+    }
+
+    /**
+     * Abort: discard all written data files and close the SDK writer.
+     * Called from C++ when write or prepareCommit fails.
+     */
+    public void abort() throws Exception {
+        try (ThreadContextClassLoader ignored =
+                new ThreadContextClassLoader(getClass().getClassLoader())) {
+            try {
+                if (preExecutionAuthenticator != null) {
+                    preExecutionAuthenticator.execute(() -> {
+                        abortWriter();
+                        return null;
+                    });
+                } else {
+                    abortWriter();
+                }
+            } catch (Exception e) {
+                LOG.error("PaimonJniWriter abort failed", e);
+                throw e;
+            }
+        }
+    }
+
+    @Override
+    protected void openInternal() {
+        throw new UnsupportedOperationException(
+                "Paimon writer requires its Arrow-specific open entry point");
+    }
+
+    @Override
+    protected void writeInternal(VectorTable inputTable) {
+        throw new UnsupportedOperationException(
+                "Paimon writer requires its Arrow C Data write entry point");
+    }
+
+    /** Close: release all resources through the shared JNI writer lifecycle. 
*/
+    @Override
+    protected void closeInternal() throws IOException {
+        try {
+            if (preExecutionAuthenticator != null) {
+                preExecutionAuthenticator.execute(() -> {
+                    closeResources();
+                    return null;
+                });
+            } else {
+                closeResources();
+            }
+        } catch (Exception e) {
+            LOG.warn("PaimonJniWriter close error", e);
+            if (e instanceof IOException) {
+                throw (IOException) e;
+            }
+            throw new IOException("PaimonJniWriter close failed", e);
+        }
+    }
+
+    /** Whether a failed close may have left an SDK task using Doris native 
callbacks. */
+    public boolean hasUnresolvedNativeCallbacks() {
+        return sdkCloseFailed;
+    }
+
+    // ────────────────────────────────────────────────────────────
+    // Initialization helpers
+    // ────────────────────────────────────────────────────────────
+
+    private void openFileStoreWriter(FileStoreTable table, String commitUser, 
boolean overwrite,
+            CoreOptions coreOptions, long nativePageMemoryLimitBytes,
+            long nativeMemoryManager, long nativeSpillSession) throws 
Exception {
+        writer = table.newWrite(commitUser);
+        compactionExecutor = Executors.newSingleThreadExecutor(
+                new ExecutorThreadFactory("doris-paimon-compaction"));
+        writer.withCompactExecutor(compactionExecutor);
+        if (overwrite) {
+            writer.withIgnorePreviousFiles(true);
+        }
+        openMemoryResources(table, coreOptions, nativePageMemoryLimitBytes,
+                nativeMemoryManager, nativeSpillSession);
+        openDynamicBucketAssigner(table, commitUser, overwrite, coreOptions,
+                nativePageMemoryLimitBytes);
+    }
+
+    private void validateWriteColumnsForMergeEngine(int writeColumnCount, 
CoreOptions coreOptions) {
+        if (writeColumnCount == table.rowType().getFieldCount() || 
table.primaryKeys().isEmpty()) {
+            return;
+        }
+
+        CoreOptions.MergeEngine mergeEngine = coreOptions.mergeEngine();
+        if (mergeEngine != CoreOptions.MergeEngine.PARTIAL_UPDATE) {
+            throw new UnsupportedOperationException(
+                    "Paimon primary-key partial-column write requires "
+                            + "merge-engine=partial-update, but table uses 
merge-engine="
+                            + mergeEngine);
+        }
+    }
+
+    private void openMemoryResources(
+            FileStoreTable table,
+            CoreOptions coreOptions,
+            long nativePageMemoryLimitBytes,
+            long nativeMemoryManager,
+            long nativeSpillSession) throws Exception {
+        int pageSize = coreOptions.pageSize();
+        long writeBufferSize = coreOptions.writeBufferSize();
+        // Paimon creates merge-tree bucket writers lazily on the first write. 
Their
+        // SortBufferWriteBuffer requires three pages at construction time, so 
reject a permanent
+        // per-writer capacity shortage during open instead of failing 
nondeterministically when a
+        // particular bucket first receives a row. Paimon's MemoryPoolFactory 
shares these pages
+        // among bucket owners; the requirement is three pages per Doris 
writer, not per bucket.
+        long effectivePoolLimit = 
validateAndGetMemoryPoolLimit(writeBufferSize,
+                nativePageMemoryLimitBytes, pageSize, 
!table.primaryKeys().isEmpty());
+        DorisMemorySegmentPool memorySegmentPool =
+                new DorisMemorySegmentPool(effectivePoolLimit, pageSize, 
nativeMemoryManager);
+        MemoryPoolFactory memoryPoolFactory = new 
MemoryPoolFactory(memorySegmentPool);
+        writer.withMemoryPoolFactory(memoryPoolFactory);
+        LOG.info("Paimon writer uses Doris-managed memory pool: limit={} 
bytes, pageSize={}",
+                memoryPoolFactory.totalBufferSize(), pageSize);
+
+        // All Paimon temporary files, including lookup and clustering files 
written directly by
+        // Paimon, use the same Doris-managed directory. DorisIOManager 
requests that directory only
+        // on its first actual use, so a memory-only writer does not depend on 
spill storage.
+        ioManager = DorisIOManager.create(nativeSpillSession);
+        writer.withIOManager(ioManager);
+        LOG.info("Paimon writer uses a lazy Doris-managed spill session");
+    }
+
+    static long validateAndGetMemoryPoolLimit(long writeBufferSize,
+            long nativePageMemoryLimitBytes, int pageSize, boolean 
mergeTreeWriter) {
+        long effectivePoolLimit = Math.min(writeBufferSize, 
nativePageMemoryLimitBytes);
+        int requiredPages = mergeTreeWriter
+                ? MERGE_TREE_WRITER_MIN_PAGES
+                : APPEND_ONLY_WRITER_MIN_PAGES;
+        long availablePages = effectivePoolLimit / pageSize;
+        if (availablePages < requiredPages) {
+            String writerType = mergeTreeWriter ? "merge-tree" : "append-only";
+            throw new IllegalArgumentException("Paimon " + writerType
+                    + " writer requires at least " + requiredPages
+                    + " memory pages, but the effective pool contains " + 
availablePages
+                    + " pages: effectivePoolLimit=" + effectivePoolLimit
+                    + ", pageSize=" + pageSize
+                    + ", writeBufferSize=" + writeBufferSize
+                    + ", nativePageMemoryLimitBytes=" + 
nativePageMemoryLimitBytes
+                    + ". Increase the query memory limit, reduce sink 
parallelism, or adjust "
+                    + "paimon_jni_writer_memory_pool_limit_bytes, 
write-buffer-size, or page-size");
+        }
+        return effectivePoolLimit;
+    }
+
+    private void openDynamicBucketAssigner(FileStoreTable table, String 
commitUser,
+            boolean overwrite, CoreOptions coreOptions, long 
writerMemoryLimitBytes) throws Exception {
+        switch (bucketMode) {
+            case HASH_DYNAMIC:
+                openHashDynamicBucketAssigner(table, commitUser, overwrite, 
coreOptions);
+                break;
+            case KEY_DYNAMIC:
+                openKeyDynamicBucketAssigner(table, writerMemoryLimitBytes);
+                break;
+            default:
+                // Fixed, unaware and postpone modes route through 
TableWrite.write(row).
+                break;
+        }
+    }
+
+    private void openHashDynamicBucketAssigner(FileStoreTable table, String 
commitUser,
+            boolean overwrite, CoreOptions coreOptions) {
+        dynamicBucketExtractor = new RowPartitionKeyExtractor(table.schema());
+        if (overwrite) {
+            hashBucketAssigner =
+                    new SimpleHashBucketAssigner(
+                            1,
+                            0,
+                            coreOptions.dynamicBucketTargetRowNum(),
+                            coreOptions.dynamicBucketMaxBuckets());
+            return;
+        }
+
+        hashBucketAssigner =
+                new HashBucketAssigner(
+                        table.snapshotManager(),
+                        commitUser,
+                        table.store().newIndexFileHandler(),
+                        1,
+                        1,
+                        0,
+                        coreOptions.dynamicBucketTargetRowNum(),
+                        coreOptions.dynamicBucketMaxBuckets());
+    }
+
+    private void openKeyDynamicBucketAssigner(FileStoreTable table,
+            long writerMemoryLimitBytes) throws Exception {
+        globalIndexAssigner = new GlobalIndexAssigner(table, 
writerMemoryLimitBytes);
+        globalIndexAssigner.open(1, 0, this::writeAssignedRow);
+        new IndexBootstrap(table).bootstrap(
+                1, 0, this::bootstrapGlobalIndexKey);
+        globalIndexAssigner.finishBootstrap();
+    }
+
+    // ────────────────────────────────────────────────────────────
+    // Data writing
+    // ────────────────────────────────────────────────────────────
+
+    private void writeBatch(VectorSchemaRoot root) throws Exception {
+        int rowCount = root.getRowCount();
+        if (rowCount == 0) {
+            return;
+        }
+        // The adapter exposes imported Arrow vectors directly as a Paimon 
columnar row. Only the
+        // table-layout row is materialized; there is no decoded Arrow copy or 
Object[][] batch.
+        PaimonArrowBatchAdapter.Rows rows = arrowAdapter.rows(root);
+        for (int r = 0; r < rowCount; r++) {
+            InternalRow row = writeSchema.tableRow(rows.row(r));
+            switch (bucketMode) {
+                case HASH_DYNAMIC:
+                    writeHashDynamicRow(row);
+                    break;
+                case KEY_DYNAMIC:
+                    globalIndexAssigner.processInput(row);
+                    break;
+                default:
+                    writeRow(row);
+                    break;
+            }
+        }
+    }
+
+    private void writeHashDynamicRow(InternalRow row) throws Exception {
+        int bucket =
+                hashBucketAssigner.assign(
+                        dynamicBucketExtractor.partition(row),
+                        
dynamicBucketExtractor.trimmedPrimaryKey(row).hashCode());
+        writeRow(row, bucket);
+    }
+
+    private void writeAssignedRow(InternalRow row, Integer bucket) {
+        try {
+            writeRow(row, bucket);
+        } catch (Exception e) {
+            throw new RuntimeException("Failed to write Paimon key-dynamic 
bucket row", e);
+        }
+    }
+
+    private void bootstrapGlobalIndexKey(InternalRow row) {
+        try {
+            globalIndexAssigner.bootstrapKey(row);
+        } catch (Exception e) {
+            throw new RuntimeException("Failed to bootstrap Paimon key-dynamic 
index", e);
+        }
+    }
+
+    private void writeRow(InternalRow row) throws Exception {
+        if (!fullCompactionChangelog) {
+            writer.write(row);
+            return;
+        }
+
+        trackFullCompactionBucket(writer.writeAndReturn(row));
+    }
+
+    private void writeRow(InternalRow row, int bucket) throws Exception {
+        if (!fullCompactionChangelog) {
+            writer.write(row, bucket);
+            return;
+        }
+
+        trackFullCompactionBucket(writer.writeAndReturn(row, bucket));
+    }
+
+    private void trackFullCompactionBucket(SinkRecord sinkRecord) {
+        if (sinkRecord == null) {
+            return;
+        }
+        fullCompactionBuckets.add(
+                new PartitionBucket(
+                        sinkRecord.partition().copy(), sinkRecord.bucket()));
+    }
+
+    // ────────────────────────────────────────────────────────────
+    // Resource management
+    // ────────────────────────────────────────────────────────────
+
+    private void closeResources() throws Exception {
+        FileStoreTable preparedTable = table;
+        String preparedCommitUser = commitUser;
+        List<CommitMessage> messages = preparedCommitMessages;
+        try {
+            closeWriter();
+        } catch (Exception closeFailure) {
+            abortPreparedAfterCompletedClose(preparedTable, 
preparedCommitUser, messages, closeFailure);
+            throw closeFailure;
+        } finally {
+            writeSchema = null;
+            arrowAdapter = null;
+            if (allocator != null) {
+                try {
+                    allocator.close();
+                } catch (Exception allocatorFailure) {
+                    abortPreparedAfterCompletedClose(
+                            preparedTable, preparedCommitUser, messages, 
allocatorFailure);
+                    throw allocatorFailure;
+                } finally {
+                    allocator = null;
+                }
+            }
+        }
+    }
+
+    private void abortPreparedAfterCompletedClose(FileStoreTable preparedTable,
+            String preparedCommitUser, List<CommitMessage> messages, Exception 
closeFailure) {
+        if (sdkCloseFailed || messages.isEmpty()) {
+            // An unresolved SDK task can still write these files. Keep them 
for reconciliation.
+            return;
+        }
+        try (InnerTableCommit committer = 
preparedTable.newCommit(preparedCommitUser)) {
+            committer.abort(messages);
+        } catch (Exception abortFailure) {
+            closeFailure.addSuppressed(abortFailure);
+        }
+    }
+
+    private List<CommitMessage> prepareCommitMessages() throws Exception {
+        if (writer == null) {
+            throw new IllegalStateException("Paimon writer is not open");
+        }
+        prepareDynamicBucketCommit();
+        submitFullCompaction();
+        List<CommitMessage> messages = commitIdentifier > 0

Review Comment:
   [P2] Retain earlier partition messages if a later prepare fails. Paimon 
1.4.2 prepares partition writers sequentially and clears each writer's new-file 
increment as it builds a local result list. If the second partition's flush 
throws, that list never returns here, so `preparedCommitMessages` stays empty; 
`abortWriter()` retries prepare but the first partition's increment has already 
been drained, leaving its staged files without an abort message. Preserve 
partial results or track created files for failure cleanup, and fault-test a 
second-partition prepare failure.



##########
fe/be-java-extensions/paimon-scanner/src/main/java/org/apache/doris/paimon/PaimonJniWriter.java:
##########
@@ -0,0 +1,791 @@
+// 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.paimon;
+
+import org.apache.doris.jni.spi.JniWriter;
+import org.apache.doris.jni.spi.ThreadContextClassLoader;
+import org.apache.doris.jni.spi.vec.VectorTable;
+import org.apache.doris.kerberos.PreExecutionAuthenticator;
+import org.apache.doris.kerberos.PreExecutionAuthenticatorCache;
+
+import org.apache.arrow.c.ArrowArray;
+import org.apache.arrow.c.ArrowSchema;
+import org.apache.arrow.c.CDataDictionaryProvider;
+import org.apache.arrow.c.Data;
+import org.apache.arrow.memory.BufferAllocator;
+import org.apache.arrow.memory.RootAllocator;
+import org.apache.arrow.vector.VectorSchemaRoot;
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.crosspartition.IndexBootstrap;
+import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.index.BucketAssigner;
+import org.apache.paimon.index.HashBucketAssigner;
+import org.apache.paimon.index.SimpleHashBucketAssigner;
+import org.apache.paimon.memory.MemoryPoolFactory;
+import org.apache.paimon.table.BucketMode;
+import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.table.sink.CommitMessage;
+import org.apache.paimon.table.sink.InnerTableCommit;
+import org.apache.paimon.table.sink.PartitionKeyExtractor;
+import org.apache.paimon.table.sink.RowPartitionKeyExtractor;
+import org.apache.paimon.table.sink.SinkRecord;
+import org.apache.paimon.table.sink.TableWriteImpl;
+import org.apache.paimon.utils.ExecutorThreadFactory;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.io.IOException;
+import java.nio.ByteBuffer;
+import java.time.ZoneId;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashSet;
+import java.util.Iterator;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Set;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
+
+/**
+ * JNI entry point for Paimon write operations.
+ *
+ * <p>Called from C++ ({@code JniPaimonWriter}) via JNI. One instance per BE 
pipeline
+ * fragment (one per {@code PaimonTableWriter}). Data path:
+ *
+ * <pre>
+ *   C++ Block → Arrow RecordBatch → Arrow C Data Interface
+ *   → PaimonJniWriter.writeArrow(arrayAddress, schemaAddress)
+ *   → zero-copy VectorSchemaRoot view
+ *   → PaimonArrowBatchAdapter (Arrow-backed Paimon columnar row)
+ *   → PaimonWriteSchema.tableRow() (canonical table-schema order)
+ *   → Paimon SDK bucket assignment and table write
+ * </pre>
+ *
+ * <p>Commit path:
+ *
+ * <pre>
+ *   PaimonTableWriter::close() → JNI → PaimonJniWriter.prepareCommit()
+ *   → TableWriteImpl.prepareCommit()
+ *   → PaimonCommitCodec.encode() → DPCM-framed byte[][]
+ *   → C++ collects TPaimonCommitMessage[] → RPC to FE → PaimonTransaction
+ * </pre>
+ */
+public class PaimonJniWriter extends JniWriter {
+    private static final Logger LOG = 
LoggerFactory.getLogger(PaimonJniWriter.class);
+    private static final int APPEND_ONLY_WRITER_MIN_PAGES = 1;
+    private static final int MERGE_TREE_WRITER_MIN_PAGES = 3;
+    private static final long COMPACTION_CLOSE_TIMEOUT_SECONDS = 60;
+
+    private final PaimonCommitCodec commitCodec = new PaimonCommitCodec();
+
+    private BufferAllocator allocator;
+    private PreExecutionAuthenticator preExecutionAuthenticator;
+    private PaimonArrowBatchAdapter arrowAdapter;
+
+    private PaimonWriteSchema writeSchema;
+    private FileStoreTable table;
+    private TableWriteImpl<?> writer;
+    private DorisIOManager ioManager;
+    private ExecutorService compactionExecutor;
+    private long commitIdentifier;
+    private String commitUser;
+    private BucketMode bucketMode;
+    private BucketAssigner hashBucketAssigner;
+    private PartitionKeyExtractor<InternalRow> dynamicBucketExtractor;
+    private GlobalIndexAssigner globalIndexAssigner;
+    private boolean fullCompactionChangelog;
+    private final Set<PartitionBucket> fullCompactionBuckets = new HashSet<>();
+    private List<CommitMessage> preparedCommitMessages = 
Collections.emptyList();
+    private boolean sdkCloseFailed;
+
+    public PaimonJniWriter() {
+        this(0, Collections.emptyMap());
+    }
+
+    public PaimonJniWriter(int batchSize, Map<String, String> params) {
+        super(batchSize, params);
+        // Imported C Data vectors reference Doris-owned buffers; this 
allocator owns only Arrow's
+        // Java-side views and metadata. Physical buffers remain charged to 
the C++ Arrow memory
+        // pool until the synchronous writeArrow call returns.
+        this.allocator = new RootAllocator(Long.MAX_VALUE);
+    }
+
+    // ────────────────────────────────────────────────────────────
+    // JNI entry points (called from C++)
+    // ────────────────────────────────────────────────────────────
+
+    /**
+     * Initialize the writer. Called once per BE pipeline fragment via JNI.
+     *
+     * <p>This method:
+     * <ol>
+     *   <li>Deserializes the target Paimon {@link FileStoreTable} selected by 
FE.</li>
+     *   <li>Creates a {@link PaimonWriteSchema} which normalizes Doris input
+     *       columns to the table-schema row layout.</li>
+     *   <li>Opens one Paimon SDK writer session.</li>
+     * </ol>
+     *
+     * @param serializedTable serialized Paimon table selected by FE
+     * @param hadoopConfig   filesystem and authentication configuration
+     * @param columnNames    output column names in the order produced by BE
+     * @param transactionId  Doris external transaction identifier
+     * @param commitUser     Paimon commit user shared with the FE committer
+     * @param overwrite      whether this is an overwrite write
+     * @param changelogWrite whether the first input column contains a row 
change operation
+     * @param timeZone       normalized Doris session timezone used for Paimon 
LTZ values
+     * @param nativePageMemoryLimitBytes maximum Doris-managed Paimon page 
memory
+     * @param nativeMemoryManager opaque BE manager used to allocate tracked 
native pages
+     * @param nativeSpillSession opaque managed spill session used for 
capacity and I/O accounting
+     */
+    public void open(String serializedTable, Map<String, String> hadoopConfig,
+                     String[] columnNames, long transactionId, String 
commitUser,
+                     boolean overwrite, boolean changelogWrite, String 
timeZone,
+                     long nativePageMemoryLimitBytes, long nativeMemoryManager,
+                     long nativeSpillSession) throws Exception {
+        try (ThreadContextClassLoader ignored =
+                new ThreadContextClassLoader(getClass().getClassLoader())) {
+            if (nativePageMemoryLimitBytes <= 0) {
+                throw new IllegalArgumentException(
+                        "PaimonJniWriter requires a positive native page 
memory limit");
+            }
+            if (nativeMemoryManager == 0) {
+                throw new IllegalArgumentException(
+                        "PaimonJniWriter requires a native memory manager");
+            }
+            this.preExecutionAuthenticator = 
PreExecutionAuthenticatorCache.getAuthenticator(hadoopConfig);
+            preExecutionAuthenticator.execute(() -> {
+                try {
+                    FileStoreTable table = 
PaimonUtils.deserialize(serializedTable);
+                    LOG.info("PaimonJniWriter opening: table={}, columns={}",
+                            table.fullName(), columnNames != null ? 
columnNames.length : 0);
+                    this.commitIdentifier = transactionId;
+                    this.table = table;
+                    this.commitUser = commitUser;
+                    this.bucketMode = table.bucketMode();
+
+                    CoreOptions coreOptions = 
CoreOptions.fromMap(table.options());
+                    this.writeSchema = PaimonWriteSchema.create(
+                            table.rowType(), columnNames, changelogWrite);
+                    this.arrowAdapter = new PaimonArrowBatchAdapter(
+                            writeSchema.inputType(), ZoneId.of(timeZone), 
allocator);
+                    validateWriteColumnsForMergeEngine(
+                            columnNames.length - (changelogWrite ? 1 : 0), 
coreOptions);
+                    this.fullCompactionChangelog =
+                            !coreOptions.writeOnly()
+                                    && coreOptions.changelogProducer()
+                                    == 
CoreOptions.ChangelogProducer.FULL_COMPACTION;
+                    openFileStoreWriter(
+                            table,
+                            commitUser,
+                            overwrite,
+                            coreOptions,
+                            nativePageMemoryLimitBytes,
+                            nativeMemoryManager,
+                            nativeSpillSession);
+                    return null;
+                } catch (Throwable t) {
+                    try {
+                        closeResources();
+                    } catch (Throwable closeFailure) {
+                        t.addSuppressed(closeFailure);
+                    }
+                    throw new RuntimeException("PaimonJniWriter open failed", 
t);
+                }
+            });
+        }
+    }
+
+    /** Return the exact Arrow schema derived from the pinned Paimon input 
type. */
+    public byte[] getArrowSchema() {
+        if (arrowAdapter == null) {
+            throw new IllegalStateException("PaimonJniWriter is not open");
+        }
+        return arrowAdapter.serializedArrowSchema();
+    }
+
+    /**
+     * Import and synchronously consume one C++ Arrow RecordBatch through the 
C Data Interface.
+     *
+     * <p>The native structs are valid only for this call. Import transfers 
the ArrowArray release
+     * callbacks into Java. Closing the imported root releases the exported 
C++ buffers after Paimon
+     * has consumed every row; the native caller releases only callbacks left 
by a partial import.
+     */
+    public void writeArrow(long arrayAddress, long schemaAddress) throws 
Exception {
+        try (ThreadContextClassLoader ignored =
+                new ThreadContextClassLoader(getClass().getClassLoader())) {
+            preExecutionAuthenticator.execute(() -> {
+                try (ArrowArray array = ArrowArray.wrap(arrayAddress);
+                        ArrowSchema schema = ArrowSchema.wrap(schemaAddress);
+                        CDataDictionaryProvider dictionaries = new 
CDataDictionaryProvider();
+                        VectorSchemaRoot root = Data.importVectorSchemaRoot(
+                                allocator, array, schema, dictionaries)) {
+                    writeBatch(root);
+                    ioManager.reconcileDirectFileGrowth();
+                    return null;
+                } catch (Throwable t) {
+                    throw new RuntimeException("PaimonJniWriter C Data write 
failed", t);
+                }
+            });
+        }
+    }
+
+    /**
+     * Prepare commit: flush all in-memory data, close files, and serialize 
commit
+     * messages for the FE coordinator.
+     *
+     * <p>Flushes and collects Paimon {@link CommitMessage}s, then encodes 
them via
+     * {@link PaimonCommitCodec} into DPCM-framed byte chunks that are 
forwarded to
+     * FE through the BE.
+     *
+     * @return byte[][]  each element is a DPCM-framed serialized 
CommitMessage chunk
+     */
+    public byte[][] prepareCommit() throws Exception {
+        try (ThreadContextClassLoader ignored =
+                new ThreadContextClassLoader(getClass().getClassLoader())) {
+            return preExecutionAuthenticator.execute(() -> {
+                try {
+                    List<CommitMessage> messages = prepareCommitMessages();
+                    ioManager.reconcileDirectFileGrowth();
+                    if (messages.isEmpty()) {
+                        LOG.info("PaimonJniWriter prepareCommit: empty");
+                        return new byte[0][];
+                    }
+                    LOG.info("PaimonJniWriter prepareCommit: {} messages", 
messages.size());
+                    return commitCodec.encode(messages);
+                } catch (Throwable t) {
+                    throw new RuntimeException("PaimonJniWriter prepareCommit 
failed", t);
+                }
+            });
+        }
+    }
+
+    /**
+     * Abort: discard all written data files and close the SDK writer.
+     * Called from C++ when write or prepareCommit fails.
+     */
+    public void abort() throws Exception {
+        try (ThreadContextClassLoader ignored =
+                new ThreadContextClassLoader(getClass().getClassLoader())) {
+            try {
+                if (preExecutionAuthenticator != null) {
+                    preExecutionAuthenticator.execute(() -> {
+                        abortWriter();
+                        return null;
+                    });
+                } else {
+                    abortWriter();
+                }
+            } catch (Exception e) {
+                LOG.error("PaimonJniWriter abort failed", e);
+                throw e;
+            }
+        }
+    }
+
+    @Override
+    protected void openInternal() {
+        throw new UnsupportedOperationException(
+                "Paimon writer requires its Arrow-specific open entry point");
+    }
+
+    @Override
+    protected void writeInternal(VectorTable inputTable) {
+        throw new UnsupportedOperationException(
+                "Paimon writer requires its Arrow C Data write entry point");
+    }
+
+    /** Close: release all resources through the shared JNI writer lifecycle. 
*/
+    @Override
+    protected void closeInternal() throws IOException {
+        try {
+            if (preExecutionAuthenticator != null) {
+                preExecutionAuthenticator.execute(() -> {
+                    closeResources();
+                    return null;
+                });
+            } else {
+                closeResources();
+            }
+        } catch (Exception e) {
+            LOG.warn("PaimonJniWriter close error", e);
+            if (e instanceof IOException) {
+                throw (IOException) e;
+            }
+            throw new IOException("PaimonJniWriter close failed", e);
+        }
+    }
+
+    /** Whether a failed close may have left an SDK task using Doris native 
callbacks. */
+    public boolean hasUnresolvedNativeCallbacks() {
+        return sdkCloseFailed;
+    }
+
+    // ────────────────────────────────────────────────────────────
+    // Initialization helpers
+    // ────────────────────────────────────────────────────────────
+
+    private void openFileStoreWriter(FileStoreTable table, String commitUser, 
boolean overwrite,
+            CoreOptions coreOptions, long nativePageMemoryLimitBytes,
+            long nativeMemoryManager, long nativeSpillSession) throws 
Exception {
+        writer = table.newWrite(commitUser);
+        compactionExecutor = Executors.newSingleThreadExecutor(
+                new ExecutorThreadFactory("doris-paimon-compaction"));
+        writer.withCompactExecutor(compactionExecutor);
+        if (overwrite) {
+            writer.withIgnorePreviousFiles(true);
+        }
+        openMemoryResources(table, coreOptions, nativePageMemoryLimitBytes,
+                nativeMemoryManager, nativeSpillSession);
+        openDynamicBucketAssigner(table, commitUser, overwrite, coreOptions,
+                nativePageMemoryLimitBytes);
+    }
+
+    private void validateWriteColumnsForMergeEngine(int writeColumnCount, 
CoreOptions coreOptions) {
+        if (writeColumnCount == table.rowType().getFieldCount() || 
table.primaryKeys().isEmpty()) {
+            return;
+        }
+
+        CoreOptions.MergeEngine mergeEngine = coreOptions.mergeEngine();
+        if (mergeEngine != CoreOptions.MergeEngine.PARTIAL_UPDATE) {
+            throw new UnsupportedOperationException(
+                    "Paimon primary-key partial-column write requires "
+                            + "merge-engine=partial-update, but table uses 
merge-engine="
+                            + mergeEngine);
+        }
+    }
+
+    private void openMemoryResources(
+            FileStoreTable table,
+            CoreOptions coreOptions,
+            long nativePageMemoryLimitBytes,
+            long nativeMemoryManager,
+            long nativeSpillSession) throws Exception {
+        int pageSize = coreOptions.pageSize();
+        long writeBufferSize = coreOptions.writeBufferSize();
+        // Paimon creates merge-tree bucket writers lazily on the first write. 
Their
+        // SortBufferWriteBuffer requires three pages at construction time, so 
reject a permanent
+        // per-writer capacity shortage during open instead of failing 
nondeterministically when a
+        // particular bucket first receives a row. Paimon's MemoryPoolFactory 
shares these pages
+        // among bucket owners; the requirement is three pages per Doris 
writer, not per bucket.
+        long effectivePoolLimit = 
validateAndGetMemoryPoolLimit(writeBufferSize,
+                nativePageMemoryLimitBytes, pageSize, 
!table.primaryKeys().isEmpty());
+        DorisMemorySegmentPool memorySegmentPool =
+                new DorisMemorySegmentPool(effectivePoolLimit, pageSize, 
nativeMemoryManager);
+        MemoryPoolFactory memoryPoolFactory = new 
MemoryPoolFactory(memorySegmentPool);
+        writer.withMemoryPoolFactory(memoryPoolFactory);
+        LOG.info("Paimon writer uses Doris-managed memory pool: limit={} 
bytes, pageSize={}",
+                memoryPoolFactory.totalBufferSize(), pageSize);
+
+        // All Paimon temporary files, including lookup and clustering files 
written directly by
+        // Paimon, use the same Doris-managed directory. DorisIOManager 
requests that directory only
+        // on its first actual use, so a memory-only writer does not depend on 
spill storage.
+        ioManager = DorisIOManager.create(nativeSpillSession);
+        writer.withIOManager(ioManager);
+        LOG.info("Paimon writer uses a lazy Doris-managed spill session");
+    }
+
+    static long validateAndGetMemoryPoolLimit(long writeBufferSize,
+            long nativePageMemoryLimitBytes, int pageSize, boolean 
mergeTreeWriter) {
+        long effectivePoolLimit = Math.min(writeBufferSize, 
nativePageMemoryLimitBytes);
+        int requiredPages = mergeTreeWriter
+                ? MERGE_TREE_WRITER_MIN_PAGES
+                : APPEND_ONLY_WRITER_MIN_PAGES;
+        long availablePages = effectivePoolLimit / pageSize;
+        if (availablePages < requiredPages) {
+            String writerType = mergeTreeWriter ? "merge-tree" : "append-only";
+            throw new IllegalArgumentException("Paimon " + writerType
+                    + " writer requires at least " + requiredPages
+                    + " memory pages, but the effective pool contains " + 
availablePages
+                    + " pages: effectivePoolLimit=" + effectivePoolLimit
+                    + ", pageSize=" + pageSize
+                    + ", writeBufferSize=" + writeBufferSize
+                    + ", nativePageMemoryLimitBytes=" + 
nativePageMemoryLimitBytes
+                    + ". Increase the query memory limit, reduce sink 
parallelism, or adjust "
+                    + "paimon_jni_writer_memory_pool_limit_bytes, 
write-buffer-size, or page-size");
+        }
+        return effectivePoolLimit;
+    }
+
+    private void openDynamicBucketAssigner(FileStoreTable table, String 
commitUser,
+            boolean overwrite, CoreOptions coreOptions, long 
writerMemoryLimitBytes) throws Exception {
+        switch (bucketMode) {
+            case HASH_DYNAMIC:
+                openHashDynamicBucketAssigner(table, commitUser, overwrite, 
coreOptions);
+                break;
+            case KEY_DYNAMIC:
+                openKeyDynamicBucketAssigner(table, writerMemoryLimitBytes);
+                break;
+            default:
+                // Fixed, unaware and postpone modes route through 
TableWrite.write(row).
+                break;
+        }
+    }
+
+    private void openHashDynamicBucketAssigner(FileStoreTable table, String 
commitUser,
+            boolean overwrite, CoreOptions coreOptions) {
+        dynamicBucketExtractor = new RowPartitionKeyExtractor(table.schema());
+        if (overwrite) {
+            hashBucketAssigner =
+                    new SimpleHashBucketAssigner(
+                            1,
+                            0,
+                            coreOptions.dynamicBucketTargetRowNum(),
+                            coreOptions.dynamicBucketMaxBuckets());
+            return;
+        }
+
+        hashBucketAssigner =
+                new HashBucketAssigner(
+                        table.snapshotManager(),
+                        commitUser,
+                        table.store().newIndexFileHandler(),
+                        1,
+                        1,
+                        0,
+                        coreOptions.dynamicBucketTargetRowNum(),
+                        coreOptions.dynamicBucketMaxBuckets());
+    }
+
+    private void openKeyDynamicBucketAssigner(FileStoreTable table,
+            long writerMemoryLimitBytes) throws Exception {
+        globalIndexAssigner = new GlobalIndexAssigner(table, 
writerMemoryLimitBytes);
+        globalIndexAssigner.open(1, 0, this::writeAssignedRow);
+        new IndexBootstrap(table).bootstrap(
+                1, 0, this::bootstrapGlobalIndexKey);
+        globalIndexAssigner.finishBootstrap();
+    }
+
+    // ────────────────────────────────────────────────────────────
+    // Data writing
+    // ────────────────────────────────────────────────────────────
+
+    private void writeBatch(VectorSchemaRoot root) throws Exception {
+        int rowCount = root.getRowCount();
+        if (rowCount == 0) {
+            return;
+        }
+        // The adapter exposes imported Arrow vectors directly as a Paimon 
columnar row. Only the
+        // table-layout row is materialized; there is no decoded Arrow copy or 
Object[][] batch.
+        PaimonArrowBatchAdapter.Rows rows = arrowAdapter.rows(root);
+        for (int r = 0; r < rowCount; r++) {
+            InternalRow row = writeSchema.tableRow(rows.row(r));
+            switch (bucketMode) {
+                case HASH_DYNAMIC:
+                    writeHashDynamicRow(row);
+                    break;
+                case KEY_DYNAMIC:
+                    globalIndexAssigner.processInput(row);
+                    break;
+                default:
+                    writeRow(row);
+                    break;
+            }
+        }
+    }
+
+    private void writeHashDynamicRow(InternalRow row) throws Exception {
+        int bucket =
+                hashBucketAssigner.assign(
+                        dynamicBucketExtractor.partition(row),
+                        
dynamicBucketExtractor.trimmedPrimaryKey(row).hashCode());
+        writeRow(row, bucket);
+    }
+
+    private void writeAssignedRow(InternalRow row, Integer bucket) {
+        try {
+            writeRow(row, bucket);
+        } catch (Exception e) {
+            throw new RuntimeException("Failed to write Paimon key-dynamic 
bucket row", e);
+        }
+    }
+
+    private void bootstrapGlobalIndexKey(InternalRow row) {
+        try {
+            globalIndexAssigner.bootstrapKey(row);
+        } catch (Exception e) {
+            throw new RuntimeException("Failed to bootstrap Paimon key-dynamic 
index", e);
+        }
+    }
+
+    private void writeRow(InternalRow row) throws Exception {
+        if (!fullCompactionChangelog) {
+            writer.write(row);

Review Comment:
   [P2] Gate binary primary-key lookup writes with the pinned Paimon runtime. 
This new write path accepts BINARY/VARBINARY primary keys and tables with 
`changelog-producer=lookup` or deletion vectors, but Paimon 1.4.2 throws 
`ClassCastException` in `SliceComparator` once lookup compaction probes an 
older level for a byte-array key 
([apache/paimon#8873](https://github.com/apache/paimon/issues/8873), explicitly 
reproduced on 1.4.2). Reject that combination until the dependency includes the 
fix, or upgrade and add a regression that forces the higher-level lookup.



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