This is an automated email from the ASF dual-hosted git repository.
yuqi1129 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/main by this push:
new d38268a57d [#11829] fix(doris): support replication allocation
property (#11877)
d38268a57d is described below
commit d38268a57d3a23dba8dd337c68dd78f764fc9741
Author: hutiefang76 <[email protected]>
AuthorDate: Thu Aug 20 20:57:13 2026 +0800
[#11829] fix(doris): support replication allocation property (#11877)
### What changed
This PR lets the JDBC Doris catalog accept the Doris 2.1+
`replication_allocation` table property.
It also keeps the existing single-BE fallback behavior for old-style
`replication_num`, but skips that automatic `replication_num=1`
injection when the user has already provided `replication_allocation`.
Doris treats those two properties as mutually exclusive, so adding both
can make table creation fail.
### Why
`replication_allocation` is valid Doris table syntax, but Gravitino did
not register it in the Doris table property metadata. In single-backend
environments, Gravitino could also add `replication_num` automatically,
which conflicts with the user-supplied allocation policy.
### Tests
```bash
JAVA_HOME=$(/usr/libexec/java_home -v 17) ./gradlew
:catalogs:catalog-jdbc-doris:test --tests
org.apache.gravitino.catalog.doris.TestDorisCatalog --tests
org.apache.gravitino.catalog.doris.operation.TestDorisTableOperationsSqlGeneration
-PskipITs
git diff --check
```
Closes #11829
---------
Co-authored-by: hutiefang <[email protected]>
Co-authored-by: Qi Yu <[email protected]>
---
.../doris/DorisTablePropertiesMetadata.java | 7 +++
.../doris/operation/DorisTableOperations.java | 11 +++-
.../gravitino/catalog/doris/TestDorisCatalog.java | 18 +++++-
.../doris/integration/test/CatalogDorisIT.java | 25 ++++++++
.../TestDorisTableOperationsSqlGeneration.java | 66 ++++++++++++++++++++++
docs/jdbc-doris-catalog.md | 1 +
6 files changed, 124 insertions(+), 4 deletions(-)
diff --git
a/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/DorisTablePropertiesMetadata.java
b/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/DorisTablePropertiesMetadata.java
index 62beaa17e4..6435b3a9e0 100644
---
a/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/DorisTablePropertiesMetadata.java
+++
b/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/DorisTablePropertiesMetadata.java
@@ -29,6 +29,7 @@ public class DorisTablePropertiesMetadata extends
JdbcTablePropertiesMetadata {
// ---- writable properties ----
public static final String REPLICATION_FACTOR = "replication_num";
+ public static final String REPLICATION_ALLOCATION = "replication_allocation";
public static final int DEFAULT_REPLICATION_FACTOR = 1;
public static final int DEFAULT_REPLICATION_FACTOR_IN_SERVER_SIDE = 3;
public static final String COMPRESSION = "compression";
@@ -49,6 +50,12 @@ public class DorisTablePropertiesMetadata extends
JdbcTablePropertiesMetadata {
false /* immutable */,
DEFAULT_REPLICATION_FACTOR, /* default value */
false /* hidden */),
+ PropertyEntry.stringOptionalPropertyEntry(
+ REPLICATION_ALLOCATION,
+ "The replication allocation policy for the table.",
+ false /* immutable */,
+ null /* default value */,
+ false /* hidden */),
PropertyEntry.stringOptionalPropertyEntry(
COMPRESSION,
"The compression type for the table (ZSTD, LZ4, LZ4F, ZLIB)."
diff --git
a/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTableOperations.java
b/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTableOperations.java
index 926c193267..657dc04c12 100644
---
a/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTableOperations.java
+++
b/catalogs/catalog-jdbc-doris/src/main/java/org/apache/gravitino/catalog/doris/operation/DorisTableOperations.java
@@ -20,6 +20,7 @@ package org.apache.gravitino.catalog.doris.operation;
import static
org.apache.gravitino.catalog.doris.DorisCatalog.DORIS_TABLE_PROPERTIES_META;
import static
org.apache.gravitino.catalog.doris.DorisTablePropertiesMetadata.DEFAULT_REPLICATION_FACTOR_IN_SERVER_SIDE;
+import static
org.apache.gravitino.catalog.doris.DorisTablePropertiesMetadata.REPLICATION_ALLOCATION;
import static
org.apache.gravitino.catalog.doris.DorisTablePropertiesMetadata.REPLICATION_FACTOR;
import static
org.apache.gravitino.catalog.doris.utils.DorisUtils.generatePartitionSqlFragment;
import static
org.apache.gravitino.catalog.jdbc.utils.JdbcConnectorUtils.escapeSqlLiteral;
@@ -179,9 +180,17 @@ public class DorisTableOperations extends
JdbcTableOperations {
resultMap = new HashMap<>(properties);
}
+ Preconditions.checkArgument(
+ !resultMap.containsKey(REPLICATION_FACTOR)
+ || !resultMap.containsKey(REPLICATION_ALLOCATION),
+ "Properties '%s' and '%s' cannot be set at the same time",
+ REPLICATION_FACTOR,
+ REPLICATION_ALLOCATION);
+
// If the backend server is less than
DEFAULT_REPLICATION_FACTOR_IN_SERVER_SIDE (3), we need to
// set the property 'replication_num' to 1 explicitly.
- if (!resultMap.containsKey(REPLICATION_FACTOR)) {
+ if (!resultMap.containsKey(REPLICATION_FACTOR)
+ && !resultMap.containsKey(REPLICATION_ALLOCATION)) {
// Try to check the number of backend servers using `show backends`,
this SQL is supported by
// all versions of Doris
String query = "show backends";
diff --git
a/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/TestDorisCatalog.java
b/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/TestDorisCatalog.java
index c014d85080..edc5129cb6 100644
---
a/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/TestDorisCatalog.java
+++
b/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/TestDorisCatalog.java
@@ -24,6 +24,7 @@ import static
org.apache.gravitino.catalog.doris.DorisTablePropertiesMetadata.BL
import static
org.apache.gravitino.catalog.doris.DorisTablePropertiesMetadata.COMPRESSION;
import static
org.apache.gravitino.catalog.doris.DorisTablePropertiesMetadata.ENABLE_UNIQUE_KEY_MERGE_ON_WRITE;
import static
org.apache.gravitino.catalog.doris.DorisTablePropertiesMetadata.LIGHT_SCHEMA_CHANGE;
+import static
org.apache.gravitino.catalog.doris.DorisTablePropertiesMetadata.REPLICATION_ALLOCATION;
import static
org.apache.gravitino.catalog.doris.DorisTablePropertiesMetadata.REPLICATION_FACTOR;
import static
org.apache.gravitino.catalog.doris.DorisTablePropertiesMetadata.STORAGE_POLICY;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
@@ -45,9 +46,9 @@ public class TestDorisCatalog {
dorisTablePropertiesMetadata.specificPropertyEntries();
// Verify the total number of registered properties.
- // 6 = 1 existing (replication_num) + 5 new (compression,
bloom_filter_columns,
- // storage_policy, light_schema_change, enable_unique_key_merge_on_write).
- Assertions.assertEquals(6, propertyEntryMap.size());
+ // 7 = replication_num, replication_allocation, compression,
bloom_filter_columns,
+ // storage_policy, light_schema_change, and
enable_unique_key_merge_on_write.
+ Assertions.assertEquals(7, propertyEntryMap.size());
// ---- replication_num (integerOptional) ----
Assertions.assertTrue(propertyEntryMap.containsKey(REPLICATION_FACTOR));
@@ -61,6 +62,17 @@ public class TestDorisCatalog {
Assertions.assertEquals(
DorisTablePropertiesMetadata.DEFAULT_REPLICATION_FACTOR,
replication.getDefaultValue());
+ // ---- replication_allocation (stringOptional) ----
+
Assertions.assertTrue(propertyEntryMap.containsKey(REPLICATION_ALLOCATION));
+ PropertyEntry<?> replicationAllocation =
propertyEntryMap.get(REPLICATION_ALLOCATION);
+ Assertions.assertEquals(REPLICATION_ALLOCATION,
replicationAllocation.getName());
+ Assertions.assertFalse(replicationAllocation.isRequired());
+ Assertions.assertFalse(replicationAllocation.isImmutable());
+ Assertions.assertFalse(replicationAllocation.isReserved());
+ Assertions.assertFalse(replicationAllocation.isHidden());
+ Assertions.assertEquals(String.class, replicationAllocation.getJavaType());
+ Assertions.assertNull(replicationAllocation.getDefaultValue());
+
// ---- compression (stringOptional) ----
Assertions.assertTrue(propertyEntryMap.containsKey(COMPRESSION));
PropertyEntry<?> compression = propertyEntryMap.get(COMPRESSION);
diff --git
a/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/integration/test/CatalogDorisIT.java
b/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/integration/test/CatalogDorisIT.java
index 792f52dda1..e5aeebf702 100644
---
a/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/integration/test/CatalogDorisIT.java
+++
b/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/integration/test/CatalogDorisIT.java
@@ -20,6 +20,8 @@ package org.apache.gravitino.catalog.doris.integration.test;
import static
org.apache.gravitino.catalog.doris.DorisTablePropertiesMetadata.BLOOM_FILTER_COLUMNS;
import static
org.apache.gravitino.catalog.doris.DorisTablePropertiesMetadata.COMPRESSION;
+import static
org.apache.gravitino.catalog.doris.DorisTablePropertiesMetadata.REPLICATION_ALLOCATION;
+import static
org.apache.gravitino.catalog.doris.DorisTablePropertiesMetadata.REPLICATION_FACTOR;
import static
org.apache.gravitino.integration.test.util.ITUtils.assertPartition;
import static
org.apache.gravitino.rel.Column.DEFAULT_VALUE_OF_CURRENT_TIMESTAMP;
import static org.junit.jupiter.api.Assertions.assertEquals;
@@ -420,6 +422,29 @@ public class CatalogDorisIT extends BaseIT {
renamedTable);
}
+ @Test
+ void testCreateTableWithReplicationAllocation() {
+ NameIdentifier tableIdentifier =
+ NameIdentifier.of(
+ schemaName,
GravitinoITUtils.genRandomName("doris_replication_allocation"));
+ String replicationAllocation = "tag.location.default: 1";
+ Map<String, String> properties = ImmutableMap.of(REPLICATION_ALLOCATION,
replicationAllocation);
+ TableCatalog tableCatalog = catalog.asTableCatalog();
+
+ tableCatalog.createTable(
+ tableIdentifier,
+ createColumns(),
+ table_comment,
+ properties,
+ Transforms.EMPTY_TRANSFORM,
+ createDistribution(),
+ null);
+
+ Table loadedTable = tableCatalog.loadTable(tableIdentifier);
+ assertEquals(replicationAllocation,
loadedTable.properties().get(REPLICATION_ALLOCATION));
+ assertFalse(loadedTable.properties().containsKey(REPLICATION_FACTOR));
+ }
+
@Test
void testDorisIllegalTableName() {
Map<String, String> properties = createTableProperties();
diff --git
a/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/operation/TestDorisTableOperationsSqlGeneration.java
b/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/operation/TestDorisTableOperationsSqlGeneration.java
index 68b442bbe9..ba7b665c01 100644
---
a/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/operation/TestDorisTableOperationsSqlGeneration.java
+++
b/catalogs/catalog-jdbc-doris/src/test/java/org/apache/gravitino/catalog/doris/operation/TestDorisTableOperationsSqlGeneration.java
@@ -18,11 +18,16 @@
*/
package org.apache.gravitino.catalog.doris.operation;
+import static
org.apache.gravitino.catalog.doris.DorisTablePropertiesMetadata.REPLICATION_ALLOCATION;
+import static
org.apache.gravitino.catalog.doris.DorisTablePropertiesMetadata.REPLICATION_FACTOR;
+
import java.sql.Connection;
import java.sql.ResultSet;
import java.sql.ResultSetMetaData;
import java.sql.Statement;
import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
import javax.sql.DataSource;
import org.apache.gravitino.catalog.doris.converter.DorisTypeConverter;
import org.apache.gravitino.catalog.jdbc.JdbcColumn;
@@ -71,6 +76,10 @@ public class TestDorisTableOperationsSqlGeneration {
}
}
+ public void setDataSource(DataSource dataSource) {
+ super.dataSource = dataSource;
+ }
+
public String createTableSql(
String tableName, JdbcColumn[] columns, Distribution distribution) {
return createTableSql(tableName, columns, distribution, "comment");
@@ -495,4 +504,61 @@ public class TestDorisTableOperationsSqlGeneration {
Assertions.assertTrue(DorisTableOperations.isVersionAtLeast("2.1.1", 2, 1,
0));
Assertions.assertFalse(DorisTableOperations.isVersionAtLeast("2.1.0", 2,
1, 1));
}
+
+ @Test
+ public void
testAppendNecessaryPropertiesAddsReplicationNumWhenBackendsAreNotEnough()
+ throws Exception {
+ TestableDorisTableOperations ops = new TestableDorisTableOperations();
+ ops.setDataSource(mockBackendDataSource(1));
+
+ Map<String, String> properties =
ops.appendNecessaryProperties(Collections.emptyMap());
+
+ Assertions.assertEquals("1", properties.get(REPLICATION_FACTOR));
+ }
+
+ @Test
+ public void testAppendNecessaryPropertiesKeepsReplicationAllocation() throws
Exception {
+ TestableDorisTableOperations ops = new TestableDorisTableOperations();
+ ops.setDataSource(mockBackendDataSource(1));
+
+ Map<String, String> properties = new HashMap<>();
+ properties.put(REPLICATION_ALLOCATION, "tag.location.default: 1");
+
+ Map<String, String> result = ops.appendNecessaryProperties(properties);
+
+ Assertions.assertEquals("tag.location.default: 1",
result.get(REPLICATION_ALLOCATION));
+ Assertions.assertFalse(result.containsKey(REPLICATION_FACTOR));
+ }
+
+ @Test
+ public void
testAppendNecessaryPropertiesRejectsConflictingReplicationProperties() {
+ TestableDorisTableOperations ops = new TestableDorisTableOperations();
+ Map<String, String> properties = new HashMap<>();
+ properties.put(REPLICATION_FACTOR, "1");
+ properties.put(REPLICATION_ALLOCATION, "tag.location.default: 1");
+
+ IllegalArgumentException exception =
+ Assertions.assertThrows(
+ IllegalArgumentException.class, () ->
ops.appendNecessaryProperties(properties));
+
+ Assertions.assertEquals(
+ "Properties 'replication_num' and 'replication_allocation' cannot be
set at the same time",
+ exception.getMessage());
+ }
+
+ private static DataSource mockBackendDataSource(int aliveBackendCount)
throws Exception {
+ DataSource dataSource = Mockito.mock(DataSource.class);
+ Connection connection = Mockito.mock(Connection.class);
+ Statement statement = Mockito.mock(Statement.class);
+ ResultSet resultSet = Mockito.mock(ResultSet.class);
+
+ Mockito.when(dataSource.getConnection()).thenReturn(connection);
+ Mockito.when(connection.createStatement()).thenReturn(statement);
+ Mockito.when(statement.executeQuery("show
backends")).thenReturn(resultSet);
+ int[] remainingAliveBackends = new int[] {aliveBackendCount};
+ Mockito.when(resultSet.next()).thenAnswer(invocation ->
remainingAliveBackends[0]-- > 0);
+ Mockito.when(resultSet.getString("Alive")).thenReturn("true");
+
+ return dataSource;
+ }
}
diff --git a/docs/jdbc-doris-catalog.md b/docs/jdbc-doris-catalog.md
index ab36c916e7..9b9647d87e 100644
--- a/docs/jdbc-doris-catalog.md
+++ b/docs/jdbc-doris-catalog.md
@@ -212,6 +212,7 @@ Only Doris built-in table properties are supported;
user-defined properties are
| Property Name | Description
| Default
Value | Required | Reserved | Immutable |
|------------------------------------|---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|---------------|----------|----------|-----------|
| `replication_num` | The number of replications for the
table. If not specified and the number of backend servers less than 3, then the
default value is 1; If BE ≥ 3, the server-side default (3) will be used. | `1`
or `3` | No | No | No |
+| `replication_allocation` | The replication allocation policy for
the table. It cannot be set together with `replication_num`.
| (none)
| No | No | No |
| `compression` | The compression type for the table.
Supported values: `ZSTD`, `LZ4`, `LZ4F`, `ZLIB`. Deprecated as a table-level
property in Doris 4.0+. Cannot be changed after table creation. |
(none) | No | No | Yes |
| `bloom_filter_columns` | Comma-separated list of columns for
which bloom filter indexes are created.
|
(none) | No | No | No |
| `storage_policy` | The name of the storage policy for
cold-hot separation.
|
(none) | No | No | No |