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]
