This is an automated email from the ASF dual-hosted git repository.
zhangstar333 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new 74227a80e46 [feature](paimon) Support Paimon table option passthrough
(#65955)
74227a80e46 is described below
commit 74227a80e46676d863ca6a398a23350eb4e905c2
Author: zhangstar333 <[email protected]>
AuthorDate: Mon Jul 27 15:40:37 2026 +0800
[feature](paimon) Support Paimon table option passthrough (#65955)
### What problem does this PR solve?
doc: https://github.com/apache/doris-website/pull/4012
Paimon provides many configurable table options, now introduces the
`paimon.table-option.*` namespace.
Doris removes the prefix, validates the option through Paimon's
`SupportedTableOptions`, and applies it to the Paimon `Table` before
serialization.
Explicit options, such as `paimon.table-option.read.batch-size`, are
preserved by the JNI scanner.
Invalid or unsupported options fail fast during Catalog initialization.
eg:
```
CREATE CATALOG paimon_catalog PROPERTIES (
"type" = "paimon",
"paimon.catalog.type" = "filesystem",
"warehouse" = "hdfs://127.0.0.1:8020/user/paimon/warehouse",
"paimon.table-option.read.batch-size" = "4096",
"paimon.jni.enable_jni_io_manager" = "true",
"paimon.jni.io_manager.tmp_dir" = "/data/doris/paimon_jni_tmp"
);
```
2. some jni params all named with paimon.jni.xxxx
---
be/src/format/table/paimon_jni_reader.cpp | 4 +-
be/src/format_v2/jni/paimon_jni_reader.cpp | 4 +-
be/test/format_v2/table/paimon_reader_test.cpp | 18 ++--
.../org/apache/doris/paimon/PaimonJniScanner.java | 19 +---
.../datasource/paimon/PaimonExternalCatalog.java | 16 ++-
.../datasource/paimon/source/PaimonScanNode.java | 10 +-
.../metastore/AbstractPaimonProperties.java | 109 ++++++++++++++++++++-
.../paimon/source/PaimonScanNodeTest.java | 12 +--
.../metastore/AbstractPaimonPropertiesTest.java | 97 ++++++++++++++++++
9 files changed, 241 insertions(+), 48 deletions(-)
diff --git a/be/src/format/table/paimon_jni_reader.cpp
b/be/src/format/table/paimon_jni_reader.cpp
index 3c734c377d9..e03e28186a2 100644
--- a/be/src/format/table/paimon_jni_reader.cpp
+++ b/be/src/format/table/paimon_jni_reader.cpp
@@ -42,8 +42,8 @@ constexpr std::string_view PAIMON_JNI_SCANNER_IO_TMP_DIR =
"paimon_jni_scanner_i
const std::string PaimonJniReader::PAIMON_OPTION_PREFIX = "paimon.";
const std::string PaimonJniReader::HADOOP_OPTION_PREFIX = "hadoop.";
-const std::string PaimonJniReader::DORIS_ENABLE_JNI_IO_MANAGER =
"doris.enable_jni_io_manager";
-const std::string PaimonJniReader::DORIS_JNI_IO_MANAGER_TMP_DIR =
"doris.jni_io_manager.tmp_dir";
+const std::string PaimonJniReader::DORIS_ENABLE_JNI_IO_MANAGER =
"jni.enable_jni_io_manager";
+const std::string PaimonJniReader::DORIS_JNI_IO_MANAGER_TMP_DIR =
"jni.io_manager.tmp_dir";
PaimonJniReader::PaimonJniReader(const std::vector<SlotDescriptor*>&
file_slot_descs,
RuntimeState* state, RuntimeProfile* profile,
diff --git a/be/src/format_v2/jni/paimon_jni_reader.cpp
b/be/src/format_v2/jni/paimon_jni_reader.cpp
index 86d16ce6f7d..730d431d4cb 100644
--- a/be/src/format_v2/jni/paimon_jni_reader.cpp
+++ b/be/src/format_v2/jni/paimon_jni_reader.cpp
@@ -28,8 +28,8 @@ namespace {
constexpr std::string_view PAIMON_OPTION_PREFIX = "paimon.";
constexpr std::string_view HADOOP_OPTION_PREFIX = "hadoop.";
-constexpr std::string_view DORIS_ENABLE_JNI_IO_MANAGER =
"doris.enable_jni_io_manager";
-constexpr std::string_view DORIS_JNI_IO_MANAGER_TMP_DIR =
"doris.jni_io_manager.tmp_dir";
+constexpr std::string_view DORIS_ENABLE_JNI_IO_MANAGER =
"jni.enable_jni_io_manager";
+constexpr std::string_view DORIS_JNI_IO_MANAGER_TMP_DIR =
"jni.io_manager.tmp_dir";
constexpr std::string_view PAIMON_JNI_SCANNER_IO_TMP_DIR =
"paimon_jni_scanner_io_tmp";
const std::string* get_paimon_predicate(const TFileScanRangeParams*
scan_params,
diff --git a/be/test/format_v2/table/paimon_reader_test.cpp
b/be/test/format_v2/table/paimon_reader_test.cpp
index 06301815b49..6ebec512f42 100644
--- a/be/test/format_v2/table/paimon_reader_test.cpp
+++ b/be/test/format_v2/table/paimon_reader_test.cpp
@@ -888,17 +888,17 @@ TEST(PaimonHybridReaderTest,
FirstNativeAndJniChildInitAreCountedOnce) {
TEST(PaimonJniReaderTest, BuildScannerParamsKeepsExplicitIOManagerTempDir) {
auto scan_params = make_paimon_jni_scan_params();
scan_params.__set_paimon_options({
- {"doris.enable_jni_io_manager", "true"},
- {"doris.jni_io_manager.tmp_dir", "/tmp/explicit-paimon-spill"},
- {"doris.jni_io_manager.impl_class", "org.example.CustomIOManager"},
+ {"jni.enable_jni_io_manager", "true"},
+ {"jni.io_manager.tmp_dir", "/tmp/explicit-paimon-spill"},
+ {"jni.io_manager.impl_class", "org.example.CustomIOManager"},
});
RuntimeState state {TQueryOptions(), TQueryGlobals()};
state.set_exec_env(ExecEnv::GetInstance());
auto params = build_paimon_jni_scanner_params(&scan_params, &state);
- EXPECT_EQ(params["paimon.doris.enable_jni_io_manager"], "true");
- EXPECT_EQ(params["paimon.doris.jni_io_manager.tmp_dir"],
"/tmp/explicit-paimon-spill");
- EXPECT_EQ(params["paimon.doris.jni_io_manager.impl_class"],
"org.example.CustomIOManager");
+ 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");
}
TEST(PaimonJniReaderTest,
BuildScannerParamsInjectsStorageRootTmpDirForEnabledIOManager) {
@@ -908,14 +908,14 @@ TEST(PaimonJniReaderTest,
BuildScannerParamsInjectsStorageRootTmpDirForEnabledIO
});
auto scan_params = make_paimon_jni_scan_params();
scan_params.__set_paimon_options({
- {"doris.enable_jni_io_manager", "true"},
+ {"jni.enable_jni_io_manager", "true"},
});
RuntimeState state {TQueryOptions(), TQueryGlobals()};
state.set_exec_env(ExecEnv::GetInstance());
auto params = build_paimon_jni_scanner_params(&scan_params, &state);
- EXPECT_EQ(params["paimon.doris.enable_jni_io_manager"], "true");
- EXPECT_EQ(params["paimon.doris.jni_io_manager.tmp_dir"],
+ EXPECT_EQ(params["paimon.jni.enable_jni_io_manager"], "true");
+ EXPECT_EQ(params["paimon.jni.io_manager.tmp_dir"],
"/data1/doris/paimon_jni_scanner_io_tmp:/data2/doris/"
"paimon_jni_scanner_io_tmp");
}
diff --git
a/fe/be-java-extensions/paimon-scanner/src/main/java/org/apache/doris/paimon/PaimonJniScanner.java
b/fe/be-java-extensions/paimon-scanner/src/main/java/org/apache/doris/paimon/PaimonJniScanner.java
index 46caf45aaee..3279aa3df74 100644
---
a/fe/be-java-extensions/paimon-scanner/src/main/java/org/apache/doris/paimon/PaimonJniScanner.java
+++
b/fe/be-java-extensions/paimon-scanner/src/main/java/org/apache/doris/paimon/PaimonJniScanner.java
@@ -24,7 +24,6 @@ import
org.apache.doris.common.security.authentication.PreExecutionAuthenticator
import
org.apache.doris.common.security.authentication.PreExecutionAuthenticatorCache;
import com.google.common.base.Preconditions;
-import org.apache.paimon.CoreOptions;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.disk.IOManager;
import org.apache.paimon.disk.IOManagerImpl;
@@ -65,12 +64,10 @@ 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";
- static final String ENABLE_JNI_IO_MANAGER =
"paimon.doris.enable_jni_io_manager";
- static final String JNI_IO_MANAGER_TMP_DIR =
"paimon.doris.jni_io_manager.tmp_dir";
- static final String JNI_IO_MANAGER_IMPL_CLASS =
"paimon.doris.jni_io_manager.impl_class";
+ 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";
private static final AtomicInteger ACTIVE_SCANNERS = new AtomicInteger();
- static final String DORIS_ENABLE_FILE_READER_ASYNC =
"paimon.jni.enable_file_reader_async";
- static final String MAX_ASYNC_READ_THRESHOLD = Long.MAX_VALUE + "b"; //
max threshold means disable
private final Map<String, String> params;
private final Map<String, String> hadoopOptionParams;
@@ -571,7 +568,6 @@ public class PaimonJniScanner extends JniScanner {
private void initTable() {
Preconditions.checkState(params.containsKey("serialized_table"));
table = PaimonUtils.deserialize(params.get("serialized_table"));
- table = table.copy(buildTableOptions(table.options()));
paimonAllFieldNames = PaimonUtils.getFieldNames(this.table.rowType());
if (LOG.isDebugEnabled()) {
LOG.debug("paimonAllFieldNames:{}", paimonAllFieldNames);
@@ -584,13 +580,4 @@ public class PaimonJniScanner extends JniScanner {
}
return value.split(delimiter);
}
-
- private Map<String, String> buildTableOptions(Map<String, String>
tableOptions) {
- Map<String, String> options = new HashMap<>(tableOptions);
- options.put(CoreOptions.READ_BATCH_SIZE.key(),
String.valueOf(batchSize));
- if
(Boolean.parseBoolean(params.getOrDefault(DORIS_ENABLE_FILE_READER_ASYNC,
"true")) == false) {
- options.put(CoreOptions.FILE_READER_ASYNC_THRESHOLD.key(),
MAX_ASYNC_READ_THRESHOLD);
- }
- return options;
- }
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonExternalCatalog.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonExternalCatalog.java
index 75e86769cd9..829ff78b9e1 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonExternalCatalog.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/paimon/PaimonExternalCatalog.java
@@ -33,6 +33,7 @@ import org.apache.logging.log4j.Logger;
import org.apache.paimon.catalog.Catalog;
import org.apache.paimon.catalog.Identifier;
import org.apache.paimon.partition.Partition;
+import org.apache.paimon.table.Table;
import java.util.ArrayList;
import java.util.List;
@@ -113,12 +114,11 @@ public class PaimonExternalCatalog extends
ExternalCatalog {
}
}
- public org.apache.paimon.table.Table getPaimonTable(NameMapping
nameMapping) {
+ public Table getPaimonTable(NameMapping nameMapping) {
return getPaimonTable(nameMapping, null, null);
}
- public org.apache.paimon.table.Table getPaimonTable(NameMapping
nameMapping, String branch,
- String queryType) {
+ public Table getPaimonTable(NameMapping nameMapping, String branch, String
queryType) {
makeSureInitialized();
try {
Identifier identifier;
@@ -134,7 +134,12 @@ public class PaimonExternalCatalog extends ExternalCatalog
{
} else {
identifier = new Identifier(nameMapping.getRemoteDbName(),
nameMapping.getRemoteTblName());
}
- return executionAuthenticator.execute(() ->
catalog.getTable(identifier));
+ return executionAuthenticator.execute(() -> {
+ Table table = catalog.getTable(identifier);
+ Map<String, String> tableOptions =
+
paimonProperties.getTableOptionsForCopy(table.options());
+ return tableOptions.isEmpty() ? table :
table.copy(tableOptions);
+ });
} catch (Exception e) {
throw new RuntimeException("Failed to get Paimon table:" +
getName() + "."
+ nameMapping.getRemoteDbName() + "." +
nameMapping.getRemoteTblName() + "$" + queryType
@@ -173,7 +178,8 @@ public class PaimonExternalCatalog extends ExternalCatalog {
public void notifyPropertiesUpdated(Map<String, String> updatedProps) {
super.notifyPropertiesUpdated(updatedProps);
if (updatedProps.keySet().stream()
- .anyMatch(key -> CacheSpec.isMetaCacheKeyForEngine(key,
PaimonExternalMetaCache.ENGINE))) {
+ .anyMatch(key -> CacheSpec.isMetaCacheKeyForEngine(key,
PaimonExternalMetaCache.ENGINE)
+ ||
AbstractPaimonProperties.isTableOptionProperty(key))) {
Env.getCurrentEnv().getExtMetaCacheMgr().removeCatalogByEngine(getId(),
PaimonExternalMetaCache.ENGINE);
}
}
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 3bfada3e9ac..56ade9d2ac9 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
@@ -92,15 +92,13 @@ public class PaimonScanNode extends FileQueryScanNode {
private static final String DORIS_END_TIMESTAMP = "endTimestamp";
private static final String DORIS_INCREMENTAL_BETWEEN_SCAN_MODE =
"incrementalBetweenScanMode";
private static final String PAIMON_PROPERTY_PREFIX = "paimon.";
- private static final String DORIS_ENABLE_FILE_READER_ASYNC =
"jni.enable_file_reader_async";
- private static final String DORIS_ENABLE_JNI_IO_MANAGER =
"doris.enable_jni_io_manager";
- private static final String DORIS_JNI_IO_MANAGER_TMP_DIR =
"doris.jni_io_manager.tmp_dir";
- private static final String DORIS_JNI_IO_MANAGER_IMPL_CLASS =
"doris.jni_io_manager.impl_class";
+ private static final String DORIS_ENABLE_JNI_IO_MANAGER =
"jni.enable_jni_io_manager";
+ private static final String DORIS_JNI_IO_MANAGER_TMP_DIR =
"jni.io_manager.tmp_dir";
+ private static final String DORIS_JNI_IO_MANAGER_IMPL_CLASS =
"jni.io_manager.impl_class";
private static final List<String> BACKEND_PAIMON_OPTIONS = Arrays.asList(
DORIS_ENABLE_JNI_IO_MANAGER,
DORIS_JNI_IO_MANAGER_TMP_DIR,
- DORIS_JNI_IO_MANAGER_IMPL_CLASS,
- DORIS_ENABLE_FILE_READER_ASYNC);
+ DORIS_JNI_IO_MANAGER_IMPL_CLASS);
private static final String PAIMON_BINLOG_SYSTEM_TABLE_TYPE = "binlog";
private static final String PAIMON_AUDIT_LOG_SYSTEM_TABLE_TYPE =
"audit_log";
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/property/metastore/AbstractPaimonProperties.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/property/metastore/AbstractPaimonProperties.java
index 9c9631a516f..4ca7836b86c 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/property/metastore/AbstractPaimonProperties.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/property/metastore/AbstractPaimonProperties.java
@@ -25,12 +25,18 @@ import com.google.common.collect.ImmutableList;
import lombok.Getter;
import org.apache.commons.lang3.StringUtils;
import org.apache.hadoop.conf.Configuration;
+import org.apache.paimon.CoreOptions;
import org.apache.paimon.catalog.Catalog;
import org.apache.paimon.options.CatalogOptions;
+import org.apache.paimon.options.ConfigOption;
+import org.apache.paimon.options.FallbackKey;
import org.apache.paimon.options.Options;
+import java.util.Collections;
import java.util.HashMap;
+import java.util.LinkedHashMap;
import java.util.List;
+import java.util.Locale;
import java.util.Map;
import java.util.concurrent.atomic.AtomicReference;
@@ -53,6 +59,12 @@ public abstract class AbstractPaimonProperties extends
MetastoreProperties {
public abstract String getPaimonCatalogType();
private static final String USER_PROPERTY_PREFIX = "paimon.";
+ private static final String DORIS_JNI_PROPERTY_PREFIX = "paimon.jni.";
+ /** The suffix after this prefix is passed to Paimon as a dynamic table
option. */
+ public static final String TABLE_OPTION_PREFIX = "paimon.table-option.";
+ private static final SupportedTableOptions SUPPORTED_TABLE_OPTIONS =
SupportedTableOptions.build();
+
+ private Map<String, String> tableOptionsMap = Collections.emptyMap();
protected AbstractPaimonProperties(Map<String, String> props) {
super(Type.PAIMON, props);
@@ -60,6 +72,12 @@ public abstract class AbstractPaimonProperties extends
MetastoreProperties {
public abstract Catalog initializeCatalog(String catalogName,
List<StorageProperties> storagePropertiesList);
+ @Override
+ public void initNormalizeAndCheckProps() {
+ super.initNormalizeAndCheckProps();
+ tableOptionsMap = extractTableOptions();
+ }
+
protected void appendCatalogOptions() {
if (StringUtils.isNotBlank(warehouse)) {
catalogOptions.set(CatalogOptions.WAREHOUSE.key(), warehouse);
@@ -68,10 +86,12 @@ public abstract class AbstractPaimonProperties extends
MetastoreProperties {
// FIXME(cmy): Rethink these custom properties
origProps.forEach((k, v) -> {
- if (k.toLowerCase().startsWith(USER_PROPERTY_PREFIX)) {
+ if (k.toLowerCase(Locale.ROOT).startsWith(USER_PROPERTY_PREFIX)) {
String newKey = k.substring(USER_PROPERTY_PREFIX.length());
if (StringUtils.isNotBlank(newKey)) {
- boolean excluded =
userStoragePrefixes.stream().anyMatch(k::startsWith);
+ boolean excluded = isTableOptionProperty(k)
+ ||
k.toLowerCase(Locale.ROOT).startsWith(DORIS_JNI_PROPERTY_PREFIX)
+ ||
userStoragePrefixes.stream().anyMatch(k::startsWith);
if (!excluded) {
catalogOptions.set(newKey, v);
}
@@ -121,6 +141,91 @@ public abstract class AbstractPaimonProperties extends
MetastoreProperties {
}
}
+ public Map<String, String> getTableOptionsMap() {
+ return tableOptionsMap;
+ }
+
+ /**
+ * Returns Catalog table options which are not explicitly configured by
the Paimon table.
+ *
+ * <p>The comparison is based on Paimon {@link ConfigOption}s so canonical
and fallback keys
+ * follow the same precedence rule.
+ */
+ public Map<String, String> getTableOptionsForCopy(Map<String, String>
currentTableOptions) {
+ if (tableOptionsMap.isEmpty() || currentTableOptions.isEmpty()) {
+ return tableOptionsMap;
+ }
+
+ Options existingOptions = new Options(currentTableOptions);
+ Map<String, String> optionsForCopy = new LinkedHashMap<>();
+ tableOptionsMap.forEach((key, value) -> {
+ ConfigOption<?> option = SUPPORTED_TABLE_OPTIONS.find(key);
+ if (!existingOptions.contains(option)) {
+ optionsForCopy.put(key, value);
+ }
+ });
+ return Collections.unmodifiableMap(optionsForCopy);
+ }
+
+ public static boolean isTableOptionProperty(String key) {
+ return key.toLowerCase(Locale.ROOT).startsWith(TABLE_OPTION_PREFIX);
+ }
+
+ private Map<String, String> extractTableOptions() {
+ Map<String, String> tableOptions = new LinkedHashMap<>();
+ origProps.forEach((key, value) -> {
+ if (isTableOptionProperty(key)) {
+ String tableOptionKey =
key.substring(TABLE_OPTION_PREFIX.length());
+ if (StringUtils.isBlank(tableOptionKey)) {
+ throw new IllegalArgumentException(
+ "Paimon table option name must not be empty after
prefix " + TABLE_OPTION_PREFIX);
+ }
+ validateTableOption(tableOptionKey, value);
+ tableOptions.put(tableOptionKey, value);
+ }
+ });
+ return Collections.unmodifiableMap(tableOptions);
+ }
+
+ private void validateTableOption(String key, String value) {
+ ConfigOption<?> option = SUPPORTED_TABLE_OPTIONS.find(key);
+ if (option == null) {
+ throw new IllegalArgumentException("Unsupported Paimon table
option '" + key
+ + "' for the bundled Paimon version");
+ }
+
+ try {
+ new Options(Collections.singletonMap(key, value)).get(option);
+ } catch (IllegalArgumentException e) {
+ throw new IllegalArgumentException("Invalid value for Paimon table
option '" + key + "': "
+ + e.getMessage(), e);
+ }
+ }
+
+ private static final class SupportedTableOptions {
+ /** Canonical and fallback option names which support direct lookup. */
+ private final Map<String, ConfigOption<?>> exactOptions;
+
+ private SupportedTableOptions(Map<String, ConfigOption<?>>
exactOptions) {
+ this.exactOptions = exactOptions;
+ }
+
+ private static SupportedTableOptions build() {
+ Map<String, ConfigOption<?>> exactOptions = new HashMap<>();
+ for (ConfigOption<?> option : CoreOptions.getOptions()) {
+ exactOptions.put(option.key(), option);
+ for (FallbackKey fallbackKey : option.fallbackKeys()) {
+ exactOptions.put(fallbackKey.getKey(), option);
+ }
+ }
+ return new
SupportedTableOptions(Collections.unmodifiableMap(exactOptions));
+ }
+
+ private ConfigOption<?> find(String key) {
+ return exactOptions.get(key);
+ }
+ }
+
/**
* @See org.apache.paimon.s3.S3FileIO
* Possible S3 config key prefixes:
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 30d79c47626..dc2cb92937b 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
@@ -549,9 +549,9 @@ public class PaimonScanNodeTest {
@Test
public void testGetBackendPaimonOptionsForJniIOManager() {
Map<String, String> props = new HashMap<>();
- props.put("paimon.doris.enable_jni_io_manager", "true");
- props.put("paimon.doris.jni_io_manager.tmp_dir", "/tmp/doris-paimon");
- props.put("paimon.doris.jni_io_manager.impl_class",
"org.example.CustomIOManager");
+ props.put("paimon.jni.enable_jni_io_manager", "true");
+ props.put("paimon.jni.io_manager.tmp_dir", "/tmp/doris-paimon");
+ props.put("paimon.jni.io_manager.impl_class",
"org.example.CustomIOManager");
CatalogProperty catalogProperty = Mockito.mock(CatalogProperty.class);
Mockito.when(catalogProperty.getProperties()).thenReturn(props);
@@ -567,10 +567,10 @@ public class PaimonScanNodeTest {
node.setSource(source);
Map<String, String> backendOptions = node.getBackendPaimonOptions();
- Assert.assertEquals("true",
backendOptions.get("doris.enable_jni_io_manager"));
- Assert.assertEquals("/tmp/doris-paimon",
backendOptions.get("doris.jni_io_manager.tmp_dir"));
+ Assert.assertEquals("true",
backendOptions.get("jni.enable_jni_io_manager"));
+ Assert.assertEquals("/tmp/doris-paimon",
backendOptions.get("jni.io_manager.tmp_dir"));
Assert.assertEquals("org.example.CustomIOManager",
- backendOptions.get("doris.jni_io_manager.impl_class"));
+ backendOptions.get("jni.io_manager.impl_class"));
Assert.assertEquals(3, backendOptions.size());
}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/datasource/property/metastore/AbstractPaimonPropertiesTest.java
b/fe/fe-core/src/test/java/org/apache/doris/datasource/property/metastore/AbstractPaimonPropertiesTest.java
index e5a775ba6e3..9cac43830ff 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/datasource/property/metastore/AbstractPaimonPropertiesTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/datasource/property/metastore/AbstractPaimonPropertiesTest.java
@@ -24,6 +24,7 @@ import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
+import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@@ -86,4 +87,100 @@ public class AbstractPaimonPropertiesTest {
Assertions.assertTrue("3".equals(result.get("fs.s3a.replication.factor")));
}
+ @Test
+ void testExtractAndValidateTableOptions() {
+ Map<String, String> input = new HashMap<>();
+ input.put("warehouse", "s3://tmp/warehouse");
+ input.put("paimon.jni.enable_jni_io_manager", "true");
+ input.put("paimon.table-option.read.batch-size", "4096");
+ input.put("paimon.table-option.file.compression.per.level",
"0:lz4,1:zstd");
+ TestPaimonProperties testProps = new TestPaimonProperties(input);
+
+ testProps.initNormalizeAndCheckProps();
+ testProps.buildCatalogOptions();
+
+ Assertions.assertEquals("4096",
testProps.getTableOptionsMap().get("read.batch-size"));
+ Assertions.assertEquals(
+ "0:lz4,1:zstd",
testProps.getTableOptionsMap().get("file.compression.per.level"));
+
Assertions.assertFalse(testProps.getCatalogOptionsMap().containsKey("table-option.read.batch-size"));
+
Assertions.assertFalse(testProps.getCatalogOptionsMap().containsKey("jni.enable_jni_io_manager"));
+ }
+
+ @Test
+ void testPaimonTableOptionsTakePrecedenceOverCatalogOptions() {
+ Map<String, String> input = new HashMap<>();
+ input.put("warehouse", "s3://tmp/warehouse");
+ input.put("paimon.table-option.read.batch-size", "4096");
+ input.put("paimon.table-option.write.batch-size", "2048");
+ input.put("paimon.table-option.file.compression.per.level",
"0:lz4,1:zstd");
+ TestPaimonProperties testProps = new TestPaimonProperties(input);
+ testProps.initNormalizeAndCheckProps();
+
+ Map<String, String> currentTableOptions = new HashMap<>();
+ currentTableOptions.put("read.batch-size", "1024");
+ currentTableOptions.put("orc.write.batch-size", "512");
+ currentTableOptions.put("file.compression.per.level", "0:snappy");
+
+ Map<String, String> optionsForCopy =
+ testProps.getTableOptionsForCopy(currentTableOptions);
+
+ Assertions.assertFalse(optionsForCopy.containsKey("read.batch-size"));
+ Assertions.assertFalse(optionsForCopy.containsKey("write.batch-size"));
+
Assertions.assertFalse(optionsForCopy.containsKey("file.compression.per.level"));
+ }
+
+ @Test
+ void testCatalogTableOptionsFillMissingPaimonTableOptions() {
+ Map<String, String> input = new HashMap<>();
+ input.put("warehouse", "s3://tmp/warehouse");
+ input.put("paimon.table-option.read.batch-size", "4096");
+ TestPaimonProperties testProps = new TestPaimonProperties(input);
+ testProps.initNormalizeAndCheckProps();
+
+ Map<String, String> optionsForCopy =
+ testProps.getTableOptionsForCopy(Collections.singletonMap(
+ "path", "s3://tmp/warehouse/test.db/test"));
+
+ Assertions.assertEquals("4096", optionsForCopy.get("read.batch-size"));
+ }
+
+ @Test
+ void testRejectUnknownTableOption() {
+ Map<String, String> input = new HashMap<>();
+ input.put("warehouse", "s3://tmp/warehouse");
+ input.put("paimon.table-option.option-does-not-exist", "value");
+ TestPaimonProperties testProps = new TestPaimonProperties(input);
+
+ IllegalArgumentException exception = Assertions.assertThrows(
+ IllegalArgumentException.class,
testProps::initNormalizeAndCheckProps);
+
+
Assertions.assertTrue(exception.getMessage().contains("option-does-not-exist"));
+ }
+
+ @Test
+ void testRejectPrefixMapTableOption() {
+ Map<String, String> input = new HashMap<>();
+ input.put("warehouse", "s3://tmp/warehouse");
+ input.put("paimon.table-option.file.compression.per.level.0", "lz4");
+ TestPaimonProperties testProps = new TestPaimonProperties(input);
+
+ IllegalArgumentException exception = Assertions.assertThrows(
+ IllegalArgumentException.class,
testProps::initNormalizeAndCheckProps);
+
+
Assertions.assertTrue(exception.getMessage().contains("file.compression.per.level.0"));
+ }
+
+ @Test
+ void testRejectInvalidTableOptionValue() {
+ Map<String, String> input = new HashMap<>();
+ input.put("warehouse", "s3://tmp/warehouse");
+ input.put("paimon.table-option.read.batch-size", "not-an-integer");
+ TestPaimonProperties testProps = new TestPaimonProperties(input);
+
+ IllegalArgumentException exception = Assertions.assertThrows(
+ IllegalArgumentException.class,
testProps::initNormalizeAndCheckProps);
+
+
Assertions.assertTrue(exception.getMessage().contains("read.batch-size"));
+ }
+
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]