This is an automated email from the ASF dual-hosted git repository.

yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/branch-4.1 by this push:
     new 55c8638d05d branch-4.1: [opt](paimon) add table cache in paimon jni 
scanner (#66018) (#66438)
55c8638d05d is described below

commit 55c8638d05d453cdddf546b250a4a61f1c866e43
Author: zhangstar333 <[email protected]>
AuthorDate: Wed Aug 5 07:43:21 2026 +0800

    branch-4.1: [opt](paimon) add table cache in paimon jni scanner (#66018) 
(#66438)
    
    cherr-pick from master #66018
---
 be/src/exec/connector/jni_connector.cpp            |  85 +++++++++------
 be/src/exec/connector/jni_connector.h              |   6 +
 be/src/format/table/paimon_jni_reader.cpp          |   7 ++
 be/src/format/table/paimon_jni_reader.h            |   7 ++
 be/src/format_v2/jni/paimon_jni_reader.cpp         |   7 ++
 be/test/format/table/paimon_jni_reader_test.cpp    |  67 ++++++++++++
 be/test/format_v2/jni/paimon_jni_reader_test.cpp   |  37 +++++++
 be/test/format_v2/table/paimon_reader_test.cpp     |   2 +
 .../org/apache/doris/paimon/PaimonJniScanner.java  |  54 ++++++++-
 .../org/apache/doris/paimon/PaimonTableCache.java  |  82 ++++++++++++++
 .../apache/doris/paimon/PaimonJniScannerTest.java  |  76 ++++++++++++-
 .../apache/doris/paimon/PaimonTableCacheTest.java  | 121 +++++++++++++++++++++
 .../apache/doris/datasource/FileQueryScanNode.java |   6 +
 .../datasource/paimon/source/PaimonScanNode.java   |   7 ++
 .../paimon/source/PaimonScanNodeTest.java          |  11 ++
 gensrc/thrift/PlanNodes.thrift                     |   2 +
 16 files changed, 536 insertions(+), 41 deletions(-)

diff --git a/be/src/exec/connector/jni_connector.cpp 
b/be/src/exec/connector/jni_connector.cpp
index a3e20f600ed..8717341eb92 100644
--- a/be/src/exec/connector/jni_connector.cpp
+++ b/be/src/exec/connector/jni_connector.cpp
@@ -20,6 +20,7 @@
 #include <glog/logging.h>
 
 #include <sstream>
+#include <utility>
 #include <variant>
 
 #include "core/block/block.h"
@@ -166,43 +167,30 @@ Status JniConnector::get_statistics(JNIEnv* env, 
std::map<std::string, std::stri
 }
 
 Status JniConnector::close() {
-    if (!_closed) {
-        JNIEnv* env = nullptr;
-        RETURN_IF_ERROR(Jni::Env::Get(&env));
-        if (_scanner_opened) {
-            COUNTER_UPDATE(_open_scanner_time, _jni_scanner_open_watcher);
-            COUNTER_UPDATE(_fill_block_time, _fill_block_watcher);
-
-            RETURN_ERROR_IF_EXC(env);
-            int64_t _append = 0;
-            RETURN_IF_ERROR(
-                    _jni_scanner_obj.call_long_method(env, 
_jni_scanner_get_append_data_time)
-                            .call(&_append));
-
-            COUNTER_UPDATE(_java_append_data_time, _append);
-
-            int64_t _create = 0;
-            RETURN_IF_ERROR(
-                    _jni_scanner_obj
-                            .call_long_method(env, 
_jni_scanner_get_create_vector_table_time)
-                            .call(&_create));
-
-            COUNTER_UPDATE(_java_create_vector_table_time, _create);
-
-            COUNTER_UPDATE(_java_scan_time, _java_scan_watcher - _append - 
_create);
+    if (_closed) {
+        return Status::OK();
+    }
+    if (!_scanner_opened) {
+        _closed = true;
+        return Status::OK();
+    }
 
-            _max_time_split_weight_counter->conditional_update(
-                    _jni_scanner_open_watcher + _fill_block_watcher + 
_java_scan_watcher,
-                    _self_split_weight);
+    JNIEnv* env = nullptr;
+    RETURN_IF_ERROR(Jni::Env::Get(&env));
 
-            // _fill_block may be failed and returned, we should release table 
in close.
-            // org.apache.doris.common.jni.JniScanner#releaseTable is 
idempotent
-            RETURN_IF_ERROR(
-                    _jni_scanner_obj.call_void_method(env, 
_jni_scanner_release_table).call());
-            RETURN_IF_ERROR(_jni_scanner_obj.call_void_method(env, 
_jni_scanner_close).call());
-        }
+    // _fill_block may fail before releasing the current Java table. 
JniScanner::releaseTable()
+    // is idempotent, so close always retries it. Java close must still run 
when that release
+    // fails, otherwise connector resources such as Paimon's static 
table-cache lease can leak.
+    auto close_status = _jni_scanner_obj.call_void_method(env, 
_jni_scanner_release_table).call();
+    auto java_close_status = _jni_scanner_obj.call_void_method(env, 
_jni_scanner_close).call();
+    if (close_status.ok() && !java_close_status.ok()) {
+        close_status = std::move(java_close_status);
     }
-    return Status::OK();
+    if (close_status.ok()) {
+        _scanner_opened = false;
+        _closed = true;
+    }
+    return close_status;
 }
 
 Status JniConnector::_init_jni_scanner(JNIEnv* env, int batch_size) {
@@ -833,6 +821,35 @@ void JniConnector::_collect_profile_before_close() {
             LOG(WARNING) << "failed to get jni env when collect profile: " << 
st;
             return;
         }
+        COUNTER_UPDATE(_open_scanner_time, _jni_scanner_open_watcher);
+        COUNTER_UPDATE(_fill_block_time, _fill_block_watcher);
+
+        int64_t append_data_time = 0;
+        auto append_time_status =
+                _jni_scanner_obj.call_long_method(env, 
_jni_scanner_get_append_data_time)
+                        .call(&append_data_time);
+        int64_t create_vector_table_time = 0;
+        auto create_table_time_status =
+                _jni_scanner_obj.call_long_method(env, 
_jni_scanner_get_create_vector_table_time)
+                        .call(&create_vector_table_time);
+        if (!append_time_status.ok()) {
+            LOG(WARNING) << "failed to collect JNI append-data time before 
close: "
+                         << append_time_status;
+        }
+        if (!create_table_time_status.ok()) {
+            LOG(WARNING) << "failed to collect JNI vector-table time before 
close: "
+                         << create_table_time_status;
+        }
+        if (append_time_status.ok() && create_table_time_status.ok()) {
+            COUNTER_UPDATE(_java_append_data_time, append_data_time);
+            COUNTER_UPDATE(_java_create_vector_table_time, 
create_vector_table_time);
+            COUNTER_UPDATE(_java_scan_time,
+                           _java_scan_watcher - append_data_time - 
create_vector_table_time);
+            _max_time_split_weight_counter->conditional_update(
+                    _jni_scanner_open_watcher + _fill_block_watcher + 
_java_scan_watcher,
+                    _self_split_weight);
+        }
+
         // update scanner metrics
         std::map<std::string, std::string> statistics_result;
         st = get_statistics(env, &statistics_result);
diff --git a/be/src/exec/connector/jni_connector.h 
b/be/src/exec/connector/jni_connector.h
index 40549963cfd..d069eb58193 100644
--- a/be/src/exec/connector/jni_connector.h
+++ b/be/src/exec/connector/jni_connector.h
@@ -254,6 +254,12 @@ public:
      */
     Status close();
 
+#ifdef BE_TEST
+    const std::map<std::string, std::string>& TEST_scanner_params() const {
+        return _scanner_params;
+    }
+#endif
+
     /**
      * Set column name to block index map from FileScanner to avoid repeated 
map creation.
      */
diff --git a/be/src/format/table/paimon_jni_reader.cpp 
b/be/src/format/table/paimon_jni_reader.cpp
index 10f6323f30b..04dcc9daeec 100644
--- a/be/src/format/table/paimon_jni_reader.cpp
+++ b/be/src/format/table/paimon_jni_reader.cpp
@@ -27,6 +27,8 @@
 #include "runtime/exec_env.h"
 #include "runtime/runtime_state.h"
 #include "util/string_util.h"
+#include "util/uid_util.h"
+
 namespace doris {
 class RuntimeProfile;
 class RuntimeState;
@@ -78,6 +80,11 @@ PaimonJniReader::PaimonJniReader(const 
std::vector<SlotDescriptor*>& file_slot_d
     if (range_params->__isset.serialized_table) {
         params["serialized_table"] = range_params->serialized_table;
     }
+    params["serialized_table_cache_key"] =
+            range_params->__isset.serialized_table_cache_key &&
+                            !range_params->serialized_table_cache_key.empty()
+                    ? range_params->serialized_table_cache_key
+                    : generate_uuid_string();
     if (range.table_format_params.__isset.table_level_row_count) {
         _remaining_table_level_row_count = 
range.table_format_params.table_level_row_count;
     } else {
diff --git a/be/src/format/table/paimon_jni_reader.h 
b/be/src/format/table/paimon_jni_reader.h
index feab10b2d39..11895eb6e15 100644
--- a/be/src/format/table/paimon_jni_reader.h
+++ b/be/src/format/table/paimon_jni_reader.h
@@ -18,6 +18,7 @@
 #pragma once
 
 #include <cstddef>
+#include <map>
 #include <string>
 #include <unordered_map>
 #include <unordered_set>
@@ -60,6 +61,12 @@ public:
 
     Status init_reader();
 
+#ifdef BE_TEST
+    const std::map<std::string, std::string>& TEST_scanner_params() const {
+        return _jni_connector->TEST_scanner_params();
+    }
+#endif
+
 private:
     int64_t _remaining_table_level_row_count;
 };
diff --git a/be/src/format_v2/jni/paimon_jni_reader.cpp 
b/be/src/format_v2/jni/paimon_jni_reader.cpp
index 730d431d4cb..01f33c5cdf0 100644
--- a/be/src/format_v2/jni/paimon_jni_reader.cpp
+++ b/be/src/format_v2/jni/paimon_jni_reader.cpp
@@ -22,6 +22,7 @@
 #include "runtime/exec_env.h"
 #include "runtime/runtime_state.h"
 #include "util/string_util.h"
+#include "util/uid_util.h"
 
 namespace doris::format::paimon {
 namespace {
@@ -94,6 +95,12 @@ Status 
PaimonJniReader::build_scanner_params(std::map<std::string, std::string>*
     (*params)["paimon_split"] = paimon_params.paimon_split;
     (*params)["paimon_predicate"] = *paimon_predicate;
     (*params)["serialized_table"] = _scan_params->serialized_table;
+    // if old Version FE not have set it, generate uuid in BE, so no need to 
compatible
+    (*params)["serialized_table_cache_key"] =
+            _scan_params->__isset.serialized_table_cache_key &&
+                            !_scan_params->serialized_table_cache_key.empty()
+                    ? _scan_params->serialized_table_cache_key
+                    : generate_uuid_string();
 
     if (_scan_params->__isset.paimon_options && 
!_scan_params->paimon_options.empty()) {
         for (const auto& kv : _scan_params->paimon_options) {
diff --git a/be/test/format/table/paimon_jni_reader_test.cpp 
b/be/test/format/table/paimon_jni_reader_test.cpp
new file mode 100644
index 00000000000..f4530f10fed
--- /dev/null
+++ b/be/test/format/table/paimon_jni_reader_test.cpp
@@ -0,0 +1,67 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#include "format/table/paimon_jni_reader.h"
+
+#include <gtest/gtest.h>
+
+#include <string>
+#include <utility>
+#include <vector>
+
+#include "gen_cpp/PlanNodes_types.h"
+#include "runtime/runtime_state.h"
+
+namespace doris {
+namespace {
+
+TFileRangeDesc make_legacy_paimon_jni_range() {
+    TFileRangeDesc range;
+    TTableFormatFileDesc table_format_params;
+    table_format_params.__set_table_format_type("paimon");
+    TPaimonFileDesc paimon_params;
+    paimon_params.__set_paimon_split("serialized-split");
+    table_format_params.__set_paimon_params(std::move(paimon_params));
+    range.__set_table_format_params(std::move(table_format_params));
+    return range;
+}
+
+TEST(LegacyPaimonJniReaderTest, 
GeneratesMissingOrEmptySerializedTableCacheKey) {
+    const auto range = make_legacy_paimon_jni_range();
+    TFileScanRangeParams scan_params;
+    scan_params.__set_serialized_table("serialized-table");
+    scan_params.__set_paimon_predicate("serialized-predicate");
+    RuntimeState state {TQueryOptions(), TQueryGlobals()};
+    const std::vector<SlotDescriptor*> file_slot_descs;
+
+    PaimonJniReader missing_key_reader(file_slot_descs, &state, nullptr, 
range, &scan_params);
+    const auto& missing_params = missing_key_reader.TEST_scanner_params();
+    EXPECT_EQ(missing_params.at("serialized_table"), "serialized-table");
+    const auto& missing_key = missing_params.at("serialized_table_cache_key");
+    EXPECT_FALSE(missing_key.empty());
+
+    scan_params.__set_serialized_table_cache_key("");
+    PaimonJniReader empty_key_reader(file_slot_descs, &state, nullptr, range, 
&scan_params);
+    const auto& empty_params = empty_key_reader.TEST_scanner_params();
+    EXPECT_EQ(empty_params.at("serialized_table"), "serialized-table");
+    const auto& empty_key = empty_params.at("serialized_table_cache_key");
+    EXPECT_FALSE(empty_key.empty());
+    EXPECT_NE(missing_key, empty_key);
+}
+
+} // namespace
+} // namespace doris
diff --git a/be/test/format_v2/jni/paimon_jni_reader_test.cpp 
b/be/test/format_v2/jni/paimon_jni_reader_test.cpp
index 921b6e70ab9..b97e6e6126b 100644
--- a/be/test/format_v2/jni/paimon_jni_reader_test.cpp
+++ b/be/test/format_v2/jni/paimon_jni_reader_test.cpp
@@ -84,6 +84,43 @@ TEST(PaimonJniReaderTest, 
UsesScanLevelPredicateBeforeLegacySplitPredicate) {
     EXPECT_EQ(params["paimon_predicate"], "scan-predicate");
 }
 
+TEST(PaimonJniReaderTest, ForwardsSerializedTableCacheKey) {
+    auto range = make_paimon_jni_range();
+    
range.table_format_params.paimon_params.__set_paimon_predicate("serialized-predicate");
+
+    auto scan_params = make_scan_params();
+    scan_params.__set_serialized_table_cache_key("table-cache-key");
+
+    PaimonJniReader reader;
+    ASSERT_TRUE(init_reader(&reader, &scan_params).ok());
+
+    std::map<std::string, std::string> params;
+    ASSERT_TRUE(build_params(&reader, range, &params).ok());
+    EXPECT_EQ(params["serialized_table_cache_key"], "table-cache-key");
+}
+
+TEST(PaimonJniReaderTest, GeneratesMissingOrEmptySerializedTableCacheKey) {
+    auto range = make_paimon_jni_range();
+    
range.table_format_params.paimon_params.__set_paimon_predicate("serialized-predicate");
+    auto scan_params = make_scan_params();
+
+    PaimonJniReader reader;
+    ASSERT_TRUE(init_reader(&reader, &scan_params).ok());
+
+    std::map<std::string, std::string> params;
+    ASSERT_TRUE(build_params(&reader, range, &params).ok());
+    EXPECT_EQ(params["serialized_table"], "serialized-table");
+    const std::string missing_key = params["serialized_table_cache_key"];
+    EXPECT_FALSE(missing_key.empty());
+
+    scan_params.__set_serialized_table_cache_key("");
+    ASSERT_TRUE(build_params(&reader, range, &params).ok());
+    EXPECT_EQ(params["serialized_table"], "serialized-table");
+    const std::string empty_key = params["serialized_table_cache_key"];
+    EXPECT_FALSE(empty_key.empty());
+    EXPECT_NE(missing_key, empty_key);
+}
+
 TEST(PaimonJniReaderTest, 
FallsBackToLegacySplitPredicateWhenScanPredicateIsMissing) {
     auto range = make_paimon_jni_range();
     
range.table_format_params.paimon_params.__set_paimon_predicate("legacy-predicate");
diff --git a/be/test/format_v2/table/paimon_reader_test.cpp 
b/be/test/format_v2/table/paimon_reader_test.cpp
index 32b82ab12ac..8a215961c0e 100644
--- a/be/test/format_v2/table/paimon_reader_test.cpp
+++ b/be/test/format_v2/table/paimon_reader_test.cpp
@@ -354,6 +354,7 @@ TFileRangeDesc 
make_legacy_paimon_native_range(TFileFormatType::type physical_fo
 TFileScanRangeParams make_paimon_jni_scan_params() {
     TFileScanRangeParams scan_params;
     scan_params.__set_serialized_table("serialized-paimon-table");
+    
scan_params.__set_serialized_table_cache_key("serialized-paimon-table-cache-key");
     scan_params.__set_paimon_predicate("serialized-paimon-predicate");
     return scan_params;
 }
@@ -943,6 +944,7 @@ TEST(PaimonJniReaderTest, 
BuildScannerParamsKeepsExplicitIOManagerTempDir) {
     EXPECT_EQ(params["paimon.jni.enable_jni_io_manager"], "true");
     EXPECT_EQ(params["paimon.jni.io_manager.tmp_dir"], 
"/tmp/explicit-paimon-spill");
     EXPECT_EQ(params["paimon.jni.io_manager.impl_class"], 
"org.example.CustomIOManager");
+    EXPECT_EQ(params["serialized_table_cache_key"], 
"serialized-paimon-table-cache-key");
 }
 
 TEST(PaimonJniReaderTest, 
BuildScannerParamsInjectsStorageRootTmpDirForEnabledIOManager) {
diff --git 
a/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonJniScanner.java
 
b/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonJniScanner.java
index 4924986b4ad..49c77314d4c 100644
--- 
a/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonJniScanner.java
+++ 
b/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonJniScanner.java
@@ -76,11 +76,14 @@ public class PaimonJniScanner extends JniScanner {
     private static final String PAIMON_OPTION_PREFIX = "paimon.";
     private static final String ASYNC_READER_THREAD_NAME_PREFIX = 
"paimon-reader-async-thread";
     private static final String FILE_READER_ASYNC_THRESHOLD = 
"file-reader-async-threshold";
+    private static final String SERIALIZED_TABLE = "serialized_table";
     private static final int MAX_MANIFEST_PARALLELISM = 256;
     static final String DORIS_MANIFEST_PARALLELISM_CAP =
             "doris.scan.manifest.parallelism-cap";
     static final String DORIS_SERIALIZED_SYSTEM_SOURCE = 
"doris.serialized-system-source";
     static final String DORIS_SYSTEM_TABLE_TYPE = "doris.system-table-type";
+    private static final String SERIALIZED_SYSTEM_SOURCE =
+            PAIMON_OPTION_PREFIX + DORIS_SERIALIZED_SYSTEM_SOURCE;
     static final String ENABLE_JNI_IO_MANAGER = 
"paimon.jni.enable_jni_io_manager";
     static final String JNI_IO_MANAGER_TMP_DIR = 
"paimon.jni.io_manager.tmp_dir";
     static final String JNI_IO_MANAGER_IMPL_CLASS = 
"paimon.jni.io_manager.impl_class";
@@ -95,7 +98,9 @@ public class PaimonJniScanner extends JniScanner {
     private final Map<String, String> hadoopOptionParams;
     private final String paimonSplit;
     private final String paimonPredicate;
+    private final String tableCacheKey;
     private Table table;
+    private PaimonTableCache.TableCacheEntry tableCacheEntry;
     private RecordReader<InternalRow> reader;
     private IOManager ioManager;
     private String ioManagerTempDirs;
@@ -134,6 +139,9 @@ public class PaimonJniScanner extends JniScanner {
         }
         paimonSplit = params.get("paimon_split");
         paimonPredicate = params.get("paimon_predicate");
+        tableCacheKey = params.get("serialized_table_cache_key");
+        Preconditions.checkState(tableCacheKey != null && 
!tableCacheKey.isEmpty(),
+                "Missing required Paimon scanner parameter: 
serialized_table_cache_key");
         String timeZone = params.getOrDefault("time_zone", 
TimeZone.getDefault().getID());
         columnValue.setTimeZone(timeZone);
         initTableInfo(columnTypes, requiredFields, batchSize);
@@ -156,8 +164,7 @@ public class PaimonJniScanner extends JniScanner {
             Thread.currentThread().setContextClassLoader(classLoader);
             preExecutionAuthenticator.execute(() -> {
                 PaimonJdbcDriverUtils.registerDriverIfNeeded(params, 
classLoader);
-                initTable();
-                initReader();
+                initTableAndReader();
                 return null;
             });
             resetDatetimeV2Precision();
@@ -361,6 +368,7 @@ public class PaimonJniScanner extends JniScanner {
                 }
             }
         } finally {
+            releaseCachedTable();
             markScannerClosedForMetrics();
         }
         if (exception != null) {
@@ -632,11 +640,13 @@ public class PaimonJniScanner extends JniScanner {
     }
 
     private void initTable() {
-        Preconditions.checkState(params.containsKey("serialized_table"));
-        table = PaimonUtils.deserialize(params.get("serialized_table"));
-        String encodedSystemSource = params.get(PAIMON_OPTION_PREFIX + 
DORIS_SERIALIZED_SYSTEM_SOURCE);
+        Preconditions.checkState(params.containsKey(SERIALIZED_TABLE));
+        table = PaimonUtils.deserialize(params.get(SERIALIZED_TABLE));
+        params.remove(SERIALIZED_TABLE);
+        String encodedSystemSource = params.get(SERIALIZED_SYSTEM_SOURCE);
         FileStoreTable systemSource = encodedSystemSource == null
                 ? null : PaimonUtils.deserialize(encodedSystemSource);
+        params.remove(SERIALIZED_SYSTEM_SOURCE);
         table = applyBackendManifestParallelism(table,
                 params.get(PAIMON_OPTION_PREFIX + 
DORIS_MANIFEST_PARALLELISM_CAP),
                 Runtime.getRuntime().availableProcessors(), systemSource,
@@ -890,6 +900,40 @@ public class PaimonJniScanner extends JniScanner {
         }
     }
 
+    private boolean initTableFromCache() {
+        PaimonTableCache.TableCacheEntry cachedEntry = 
PaimonTableCache.acquire(tableCacheKey);
+        if (cachedEntry == null) {
+            return false;
+        }
+        tableCacheEntry = cachedEntry;
+        table = cachedEntry.table();
+        paimonAllFieldNames = cachedEntry.fieldNames();
+        params.remove(SERIALIZED_TABLE);
+        params.remove(SERIALIZED_SYSTEM_SOURCE);
+        return true;
+    }
+
+    private void initTableAndReader() throws IOException {
+        if (initTableFromCache()) {
+            initReader();
+            return;
+        }
+        initTable();
+        initReader();
+        PaimonTableCache.TableCacheEntry candidate =
+                new PaimonTableCache.TableCacheEntry(table, 
paimonAllFieldNames);
+        if (PaimonTableCache.publish(tableCacheKey, candidate)) {
+            tableCacheEntry = candidate;
+        }
+    }
+
+    private void releaseCachedTable() {
+        if (tableCacheEntry != null) {
+            PaimonTableCache.release(tableCacheKey, tableCacheEntry);
+            tableCacheEntry = null;
+        }
+    }
+
     private static String[] splitParam(String value, String delimiter) {
         if (value == null || value.isEmpty()) {
             return new String[0];
diff --git 
a/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonTableCache.java
 
b/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonTableCache.java
new file mode 100644
index 00000000000..e401f7afe02
--- /dev/null
+++ 
b/fe/be-java-extensions/paimon-connector/src/main/java/org/apache/doris/paimon/PaimonTableCache.java
@@ -0,0 +1,82 @@
+// 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 com.google.common.base.Preconditions;
+import org.apache.paimon.table.Table;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.ConcurrentHashMap;
+
+final class PaimonTableCache {
+    private static final ConcurrentHashMap<String, TableCacheEntry> 
TABLE_CACHE = new ConcurrentHashMap<>();
+
+    private PaimonTableCache() {
+    }
+
+    static TableCacheEntry acquire(String cacheKey) {
+        return TABLE_CACHE.computeIfPresent(cacheKey, (key, entry) -> {
+            entry.users++;
+            return entry;
+        });
+    }
+
+    static boolean publish(String cacheKey, TableCacheEntry entry) {
+        return TABLE_CACHE.putIfAbsent(cacheKey, entry) == null;
+    }
+
+    static void release(String cacheKey, TableCacheEntry expectedEntry) {
+        TABLE_CACHE.compute(cacheKey, (key, currentEntry) -> {
+            Preconditions.checkState(currentEntry == expectedEntry,
+                    "Paimon table cache entry changed unexpectedly for key 
%s", cacheKey);
+            Preconditions.checkState(currentEntry.users > 0,
+                    "Paimon table cache reference count is invalid for key 
%s", cacheKey);
+            currentEntry.users--;
+            return currentEntry.users == 0 ? null : currentEntry;
+        });
+    }
+
+    static int size() {
+        return TABLE_CACHE.size();
+    }
+
+    static void clearForTest() {
+        TABLE_CACHE.clear();
+    }
+
+    static final class TableCacheEntry {
+        private final Table table;
+        private final List<String> fieldNames;
+        private int users = 1;
+
+        TableCacheEntry(Table table, List<String> fieldNames) {
+            this.table = table;
+            this.fieldNames = Collections.unmodifiableList(new 
ArrayList<>(fieldNames));
+        }
+
+        Table table() {
+            return table;
+        }
+
+        List<String> fieldNames() {
+            return fieldNames;
+        }
+    }
+}
diff --git 
a/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonJniScannerTest.java
 
b/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonJniScannerTest.java
index 9918618c07a..8d040bfc40d 100644
--- 
a/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonJniScannerTest.java
+++ 
b/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonJniScannerTest.java
@@ -45,6 +45,7 @@ import org.apache.paimon.table.Table;
 import org.apache.paimon.table.system.SystemTableLoader;
 import org.apache.paimon.types.RowType;
 import org.apache.paimon.utils.InstantiationUtil;
+import org.junit.After;
 import org.junit.Assert;
 import org.junit.Rule;
 import org.junit.Test;
@@ -72,14 +73,63 @@ import java.util.concurrent.atomic.AtomicInteger;
 import java.util.concurrent.atomic.AtomicLong;
 
 public class PaimonJniScannerTest {
+    private static final String SERIALIZED_TABLE = "serialized_table";
+    private static final String SERIALIZED_SYSTEM_SOURCE =
+            "paimon.doris.serialized-system-source";
+    private static final String SERIALIZED_TABLE_CACHE_KEY = 
"serialized_table_cache_key";
+
     @Rule
     public TemporaryFolder temporaryFolder = new TemporaryFolder();
 
+    @After
+    public void clearTableCache() {
+        PaimonTableCache.clearForTest();
+    }
+
     @Test
     public void testConstructorAcceptsEmptyProjection() {
         new PaimonJniScanner(128, createBaseParams());
     }
 
+    @Test
+    public void testConstructorRejectsMissingOrEmptyTableCacheKey() {
+        Map<String, String> params = createBaseParams();
+        params.remove(SERIALIZED_TABLE_CACHE_KEY);
+        assertInvalidTableCacheKey(params);
+
+        params.put(SERIALIZED_TABLE_CACHE_KEY, "");
+        assertInvalidTableCacheKey(params);
+    }
+
+    @Test
+    public void testWarmTableCacheHitReleasesSerializedTablePayloads() throws 
Exception {
+        String cacheKey = "warm-table-cache-hit";
+        Map<String, String> params = createBaseParams();
+        params.put(SERIALIZED_TABLE_CACHE_KEY, cacheKey);
+        params.put(SERIALIZED_TABLE, "serialized-table-payload");
+        params.put(SERIALIZED_SYSTEM_SOURCE, 
"serialized-system-source-payload");
+        Table cachedTable = tableWithOptions(Collections.emptyMap());
+        PaimonTableCache.TableCacheEntry cacheEntry =
+                new PaimonTableCache.TableCacheEntry(cachedTable, 
Collections.emptyList());
+        Assert.assertTrue(PaimonTableCache.publish(cacheKey, cacheEntry));
+
+        PaimonJniScanner scanner = new PaimonJniScanner(128, params);
+        Method initTableFromCache = 
PaimonJniScanner.class.getDeclaredMethod("initTableFromCache");
+        initTableFromCache.setAccessible(true);
+
+        Assert.assertTrue((Boolean) initTableFromCache.invoke(scanner));
+        Assert.assertFalse(params.containsKey(SERIALIZED_TABLE));
+        Assert.assertFalse(params.containsKey(SERIALIZED_SYSTEM_SOURCE));
+        Field tableField = PaimonJniScanner.class.getDeclaredField("table");
+        tableField.setAccessible(true);
+        Assert.assertSame(cachedTable, tableField.get(scanner));
+
+        scanner.close();
+        Assert.assertEquals(1, PaimonTableCache.size());
+        PaimonTableCache.release(cacheKey, cacheEntry);
+        Assert.assertEquals(0, PaimonTableCache.size());
+    }
+
     @Test
     public void testOldFeSerializedAsyncThresholdIsRejectedInEveryChild() 
throws Exception {
         Table visible = (Table) 
Proxy.newProxyInstance(Table.class.getClassLoader(),
@@ -647,8 +697,11 @@ public class PaimonJniScannerTest {
     }
 
     @Test
-    public void testFailedCloseRetainsResourcesForRetry() throws Exception {
-        PaimonJniScanner scanner = new PaimonJniScanner(128, 
createBaseParams());
+    public void testFailedCloseReleasesCacheAndRetainsResourcesForRetry() 
throws Exception {
+        String cacheKey = "retryable-close";
+        Map<String, String> params = createBaseParams();
+        params.put(SERIALIZED_TABLE_CACHE_KEY, cacheKey);
+        PaimonJniScanner scanner = new PaimonJniScanner(128, params);
         AtomicInteger iteratorCloseCalls = new AtomicInteger();
         RecordReader.RecordIterator<InternalRow> recordIterator =
                 new RecordReader.RecordIterator<InternalRow>() {
@@ -689,6 +742,13 @@ public class PaimonJniScannerTest {
         Field ioManagerField = 
PaimonJniScanner.class.getDeclaredField("ioManager");
         ioManagerField.setAccessible(true);
         ioManagerField.set(scanner, ioManager);
+        PaimonTableCache.TableCacheEntry cacheEntry =
+                new 
PaimonTableCache.TableCacheEntry(tableWithOptions(Collections.emptyMap()),
+                        Collections.emptyList());
+        Assert.assertTrue(PaimonTableCache.publish(cacheKey, cacheEntry));
+        Field cacheEntryField = 
PaimonJniScanner.class.getDeclaredField("tableCacheEntry");
+        cacheEntryField.setAccessible(true);
+        cacheEntryField.set(scanner, cacheEntry);
 
         try {
             scanner.close();
@@ -699,6 +759,7 @@ public class PaimonJniScannerTest {
         Assert.assertSame(recordIterator, recordIteratorField.get(scanner));
         Assert.assertSame(reader, readerField.get(scanner));
         Assert.assertSame(ioManager, ioManagerField.get(scanner));
+        Assert.assertEquals(0, PaimonTableCache.size());
 
         scanner.close();
         Assert.assertNull(recordIteratorField.get(scanner));
@@ -707,6 +768,7 @@ public class PaimonJniScannerTest {
         Assert.assertEquals(2, iteratorCloseCalls.get());
         Assert.assertEquals(2, readerCloseCalls.get());
         Assert.assertEquals(2, ioManager.closeCalls.get());
+        Assert.assertEquals(0, PaimonTableCache.size());
     }
 
     private Map<String, String> createBaseParams() {
@@ -715,9 +777,19 @@ public class PaimonJniScannerTest {
         params.put("columns_types", "");
         params.put("paimon_split", "");
         params.put("paimon_predicate", "");
+        params.put(SERIALIZED_TABLE_CACHE_KEY, "test-table-cache-key");
         return params;
     }
 
+    private void assertInvalidTableCacheKey(Map<String, String> params) {
+        try {
+            new PaimonJniScanner(128, params);
+            Assert.fail("expected constructor to reject an invalid table cache 
key");
+        } catch (IllegalStateException e) {
+            
Assert.assertTrue(e.getMessage().contains(SERIALIZED_TABLE_CACHE_KEY));
+        }
+    }
+
     private String encodeFields(String... fields) {
         return Arrays.stream(fields)
                 .map(field -> "$" + 
Base64.getEncoder().encodeToString(field.getBytes(StandardCharsets.UTF_8)))
diff --git 
a/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonTableCacheTest.java
 
b/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonTableCacheTest.java
new file mode 100644
index 00000000000..0ec066dc45c
--- /dev/null
+++ 
b/fe/be-java-extensions/paimon-connector/src/test/java/org/apache/doris/paimon/PaimonTableCacheTest.java
@@ -0,0 +1,121 @@
+// 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.paimon.table.Table;
+import org.junit.After;
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.lang.reflect.Proxy;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+import java.util.stream.IntStream;
+
+public class PaimonTableCacheTest {
+    @After
+    public void tearDown() {
+        PaimonTableCache.clearForTest();
+    }
+
+    @Test
+    public void testAcquireIncrementsAndReleaseRemovesAtZero() {
+        String cacheKey = "reference-count";
+        PaimonTableCache.TableCacheEntry publishedEntry = newEntry();
+        Assert.assertTrue(PaimonTableCache.publish(cacheKey, publishedEntry));
+
+        PaimonTableCache.TableCacheEntry acquiredEntry = 
PaimonTableCache.acquire(cacheKey);
+        Assert.assertSame(publishedEntry, acquiredEntry);
+
+        PaimonTableCache.release(cacheKey, acquiredEntry);
+        Assert.assertEquals(1, PaimonTableCache.size());
+
+        PaimonTableCache.release(cacheKey, publishedEntry);
+        Assert.assertEquals(0, PaimonTableCache.size());
+        Assert.assertNull(PaimonTableCache.acquire(cacheKey));
+    }
+
+    @Test
+    public void 
testConcurrentAcquireAndReleaseKeepsEntryWhilePublisherUsesIt() throws 
Exception {
+        String cacheKey = "concurrent-reference-count";
+        PaimonTableCache.TableCacheEntry publishedEntry = newEntry();
+        Assert.assertTrue(PaimonTableCache.publish(cacheKey, publishedEntry));
+
+        int threadCount = 16;
+        int iterations = 1000;
+        CountDownLatch start = new CountDownLatch(1);
+        ExecutorService executor = Executors.newFixedThreadPool(threadCount);
+        List<Future<?>> futures = IntStream.range(0, threadCount)
+                .mapToObj(ignored -> executor.submit(() -> {
+                    start.await();
+                    for (int i = 0; i < iterations; i++) {
+                        PaimonTableCache.TableCacheEntry entry = 
PaimonTableCache.acquire(cacheKey);
+                        Assert.assertSame(publishedEntry, entry);
+                        PaimonTableCache.release(cacheKey, entry);
+                    }
+                    return null;
+                }))
+                .collect(Collectors.toList());
+
+        try {
+            start.countDown();
+            for (Future<?> future : futures) {
+                future.get(30, TimeUnit.SECONDS);
+            }
+        } finally {
+            executor.shutdownNow();
+        }
+
+        Assert.assertEquals(1, PaimonTableCache.size());
+        PaimonTableCache.release(cacheKey, publishedEntry);
+        Assert.assertEquals(0, PaimonTableCache.size());
+    }
+
+    @Test
+    public void testOnlyFirstEntryIsPublished() {
+        String cacheKey = "publish-race";
+        PaimonTableCache.TableCacheEntry first = newEntry();
+        PaimonTableCache.TableCacheEntry second = newEntry();
+
+        Assert.assertTrue(PaimonTableCache.publish(cacheKey, first));
+        Assert.assertFalse(PaimonTableCache.publish(cacheKey, second));
+        PaimonTableCache.TableCacheEntry acquiredEntry = 
PaimonTableCache.acquire(cacheKey);
+        Assert.assertSame(first, acquiredEntry);
+
+        PaimonTableCache.release(cacheKey, acquiredEntry);
+        PaimonTableCache.release(cacheKey, first);
+        Assert.assertEquals(0, PaimonTableCache.size());
+    }
+
+    private PaimonTableCache.TableCacheEntry newEntry() {
+        Table table = (Table) Proxy.newProxyInstance(
+                Table.class.getClassLoader(), new Class[] {Table.class}, 
(proxy, method, args) -> {
+                    if ("toString".equals(method.getName())) {
+                        return "TestPaimonTable";
+                    }
+                    throw new UnsupportedOperationException(method.getName());
+                });
+        return new PaimonTableCache.TableCacheEntry(table, 
Collections.singletonList("field"));
+    }
+}
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/datasource/FileQueryScanNode.java 
b/fe/fe-core/src/main/java/org/apache/doris/datasource/FileQueryScanNode.java
index 873ab200350..d83326a1838 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/datasource/FileQueryScanNode.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/datasource/FileQueryScanNode.java
@@ -343,6 +343,11 @@ public abstract class FileQueryScanNode extends 
FileScanNode {
         return Optional.empty();
     }
 
+    // Identify scanner instances that may share the deserialized JNI table.
+    protected Optional<String> getSerializedTableCacheKey() {
+        return Optional.empty();
+    }
+
     @Override
     public void createScanRangeLocations() throws UserException {
         long start = System.currentTimeMillis();
@@ -473,6 +478,7 @@ public abstract class FileQueryScanNode extends 
FileScanNode {
         }
 
         getSerializedTable().ifPresent(params::setSerializedTable);
+        
getSerializedTableCacheKey().ifPresent(params::setSerializedTableCacheKey);
 
         if (executor != null) {
             executor.getSummaryProfile().setCreateScanRangeFinishTime();
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/source/PaimonScanNode.java
 
b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/source/PaimonScanNode.java
index 3e9136a0565..7707aed9a62 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/source/PaimonScanNode.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/source/PaimonScanNode.java
@@ -89,6 +89,7 @@ import java.util.Locale;
 import java.util.Map;
 import java.util.Optional;
 import java.util.OptionalInt;
+import java.util.UUID;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.stream.Collectors;
 
@@ -169,6 +170,7 @@ public class PaimonScanNode extends FileQueryScanNode {
     private int paimonSplitNum = 0;
     private List<SplitStat> splitStats = new ArrayList<>();
     private String serializedTable;
+    private final String serializedTableCacheKey = 
UUID.randomUUID().toString();
     // Store PropertiesMap, including vended credentials or static credentials
     // get them in doInitialize() to ensure internal consistency of ScanNode
     private Map<StorageProperties.Type, StorageProperties> 
storagePropertiesMap;
@@ -291,6 +293,11 @@ public class PaimonScanNode extends FileQueryScanNode {
         return Optional.of(serializedTable);
     }
 
+    @Override
+    protected Optional<String> getSerializedTableCacheKey() {
+        return Optional.of(serializedTableCacheKey);
+    }
+
     @Override
     public void createScanRangeLocations() throws UserException {
         super.createScanRangeLocations();
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/source/PaimonScanNodeTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/source/PaimonScanNodeTest.java
index 8480a48b71e..9cb649b02ac 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/source/PaimonScanNodeTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/datasource/paimon/source/PaimonScanNodeTest.java
@@ -104,6 +104,17 @@ public class PaimonScanNodeTest {
     @Mock
     private PaimonFileExternalCatalog paimonFileExternalCatalog;
 
+    @Test
+    public void testSerializedTableCacheKeyIsStablePerScanNode() {
+        PaimonScanNode first = newTestNode(new PlanNodeId(0), new TupleId(0), 
sv);
+        PaimonScanNode second = newTestNode(new PlanNodeId(1), new TupleId(1), 
sv);
+
+        String firstKey = first.getSerializedTableCacheKey().orElse("");
+        Assert.assertFalse(firstKey.isEmpty());
+        Assert.assertEquals(firstKey, 
first.getSerializedTableCacheKey().orElse(""));
+        Assert.assertNotEquals(firstKey, 
second.getSerializedTableCacheKey().orElse(""));
+    }
+
     @Test
     public void 
testCountColumnKeepsAllSplitsWhileCountStarUsesMergedRowCount() throws 
UserException {
         PaimonScanNode node = Mockito.spy(newTestNode(new PlanNodeId(1), new 
TupleId(3), sv));
diff --git a/gensrc/thrift/PlanNodes.thrift b/gensrc/thrift/PlanNodes.thrift
index a76e4543b38..a78188fe543 100644
--- a/gensrc/thrift/PlanNodes.thrift
+++ b/gensrc/thrift/PlanNodes.thrift
@@ -534,6 +534,8 @@ struct TFileScanRangeParams {
     // behavior during a BE-first rolling upgrade; version 1 enables file-wide 
ID projection and
     // logical initial-default materialization.
     34: optional i32 iceberg_scan_semantics_version
+    // FE-generated identity for sharing a deserialized table across JNI 
scanners in one scan node.
+    35: optional string serialized_table_cache_key
 }
 
 struct TFileRangeDesc {


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

Reply via email to