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

FANNG1 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 90c3919a45  [#10368] improvement(flink): persist Paimon bucket 
distribution when creating table via Flink connector (#10721)
90c3919a45 is described below

commit 90c3919a45b6569544b66558282dc3942ffc96e2
Author: EdG <[email protected]>
AuthorDate: Fri May 8 16:21:28 2026 +0800

     [#10368] improvement(flink): persist Paimon bucket distribution when 
creating table via Flink connector (#10721)
    
    ### What changes were proposed in this pull request?
    
    When creating a Paimon table via the Flink connector with bucket options
    (`bucket-key`, `bucket`), the distribution metadata was silently dropped
    because Flink's `createTable` always
      passed `Distributions.NONE`.
    
    This PR adds end-to-end support for Paimon bucket distribution through
    the Flink connector, using a **pass-through approach** — the Flink
    connector simply translates between Flink table
    options and Gravitino's Distribution model without enforcing
    Paimon-specific validation rules. All Paimon-related logic (e.g.,
    defaulting bucket-key to primary keys, validating bucket
      numbers) is deferred to the Paimon catalog service.
    
      **Flink connector (write path — `createTable`):**
    - Added a `toGravitinoDistribution` hook in `BaseCatalog` (defaults to
    `Distributions.NONE`) to allow catalog-specific distribution parsing
    from Flink table options.
    - Overrode the hook in `GravitinoPaimonCatalog` to parse `bucket-key`
    and `bucket` into a Gravitino `Distribution` object (HASH strategy).
    - Overrode `toGravitinoTableProperties` in `PaimonPropertiesConverter`
    to strip `bucket-key` and `bucket` before passing properties to
    Gravitino, since they are now represented as
      Distribution metadata.
    
      **Flink connector (read path — `toFlinkTable`):**
    - Added a `fromGravitinoDistribution` hook in `BaseCatalog` (defaults to
    empty map) to convert Distribution back into Flink table options.
    - Overrode the hook in `GravitinoPaimonCatalog` to reconstruct
    `bucket-key` and `bucket` from the Gravitino Distribution.
    - Overrode `toFlinkTableProperties` in `PaimonPropertiesConverter` to
    strip bucket properties from raw table properties, ensuring Distribution
    is the single source of truth.
    
      **Catalog-side (`catalog-lakehouse-paimon`):**
    - Refactored `GravitinoPaimonTable.getDistribution` to handle `-1` as
    AUTO and support bucket-only distribution (no explicit bucket-key).
    - Relaxed `PaimonCatalogOperations.validateDistribution` to allow HASH
    distribution without explicit bucket-key expressions, deferring the
    primary-key-subset check to Paimon itself.
    
      ### Why are the changes needed?
    
    Flink's `createTable` always passed `Distributions.NONE`, so
    `bucket-key` and `bucket` configs specified in Flink SQL `WITH` options
    were silently lost. Additionally, the server-side
    validation was overly strict — it rejected distributions without
    explicit bucket-key even though Paimon automatically falls back to
    primary keys.
    
      Fix: #10368
    
      ### Does this PR introduce _any_ user-facing change?
    
      Yes:
    - Users can now specify `bucket-key` and `bucket` in Flink SQL `WITH`
    options when creating Paimon tables, and the distribution will be
    correctly persisted in Gravitino metadata and
      round-tripped back through queries like `SHOW CREATE TABLE`.
    - Specifying only `bucket` (without `bucket-key`) is now supported —
    Paimon will use primary keys as the default bucket keys.
    
      ### How was this patch tested?
    
    - Added unit test in `TestBaseCatalog` verifying the default hooks
    return `Distributions.NONE` / empty map.
      - Added unit tests in `TestPaimonPropertiesConverter` covering:
    - `getDistribution`: null properties, blank bucket-key, missing
    bucket-key, missing bucket, both present, `-1` as AUTO, invalid bucket
    number, negative bucket number pass-through.
    - `distributionToProperties`: HASH with bucket/bucket-key, NONE
    distribution, null distribution, AUTO distribution, AUTO with no
    expressions.
    - `toGravitinoTableProperties`: verifying bucket properties are
    stripped.
    - Round-trip idempotency (bucket-key only ↔ bucket-key with explicit
    `-1`).
    - Added integration test `testBucketDistributionRoundTrip` in
    `FlinkPaimonCatalogIT` asserting the distribution is correctly persisted
    via Gravitino and round-trips back through Flink's
       `SHOW CREATE TABLE`.
---
 .../lakehouse/paimon/GravitinoPaimonTable.java     |  43 +++--
 .../lakehouse/paimon/PaimonCatalogOperations.java  |  51 ++----
 .../lakehouse/paimon/TestGravitinoPaimonTable.java |  36 -----
 docs/flink-connector/flink-catalog-paimon.md       |  22 +++
 flink-connector/flink/build.gradle.kts             |   1 +
 .../flink/connector/catalog/BaseCatalog.java       |  21 ++-
 .../connector/paimon/GravitinoPaimonCatalog.java   |  99 ++++++++++++
 .../paimon/PaimonPropertiesConverter.java          |  24 +++
 .../flink/connector/catalog/TestBaseCatalog.java   |  13 ++
 .../test/paimon/FlinkPaimonCatalogIT.java          | 127 +++++++++++++++
 .../paimon/TestPaimonPropertiesConverter.java      | 180 +++++++++++++++++++++
 11 files changed, 522 insertions(+), 95 deletions(-)

diff --git 
a/catalogs/catalog-lakehouse-paimon/src/main/java/org/apache/gravitino/catalog/lakehouse/paimon/GravitinoPaimonTable.java
 
b/catalogs/catalog-lakehouse-paimon/src/main/java/org/apache/gravitino/catalog/lakehouse/paimon/GravitinoPaimonTable.java
index 4e1714ee7b..eec7d5ce7e 100644
--- 
a/catalogs/catalog-lakehouse-paimon/src/main/java/org/apache/gravitino/catalog/lakehouse/paimon/GravitinoPaimonTable.java
+++ 
b/catalogs/catalog-lakehouse-paimon/src/main/java/org/apache/gravitino/catalog/lakehouse/paimon/GravitinoPaimonTable.java
@@ -232,30 +232,41 @@ public class GravitinoPaimonTable extends BaseTable {
     if (properties == null) {
       return Distributions.NONE;
     }
-    String bucketKeys = properties.get(BUCKET_KEY);
-    if (StringUtils.isBlank(bucketKeys)) {
+
+    String bucketKeyStr = properties.get(BUCKET_KEY);
+    String bucketNumStr = properties.get(BUCKET_NUM);
+
+    boolean hasBucketKey = StringUtils.isNotBlank(bucketKeyStr);
+    boolean hasBucket = StringUtils.isNotBlank(bucketNumStr);
+
+    if (!hasBucketKey && !hasBucket) {
       return Distributions.NONE;
     }
-    List<String> bucketKeyList =
-        Arrays.stream(bucketKeys.split(","))
-            .map(String::trim)
-            .filter(StringUtils::isNotBlank)
-            .collect(Collectors.toList());
-    if (bucketKeyList.isEmpty()) {
-      return Distributions.NONE;
+
+    Expression[] expressions = new Expression[0];
+    if (hasBucketKey) {
+      expressions =
+          Arrays.stream(bucketKeyStr.split(","))
+              .map(String::trim)
+              .filter(StringUtils::isNotBlank)
+              .map(NamedReference::field)
+              .toArray(Expression[]::new);
     }
-    Expression[] expressions =
-        
bucketKeyList.stream().map(NamedReference::field).toArray(Expression[]::new);
-    String bucketValue = properties.get(BUCKET_NUM);
-    if (StringUtils.isBlank(bucketValue)) {
+
+    if (!hasBucket) {
       return Distributions.auto(Strategy.HASH, expressions);
     }
-    String trimmedBucketValue = bucketValue.trim();
+
     try {
-      return Distributions.hash(Integer.parseInt(trimmedBucketValue), 
expressions);
+      int parsedBucket = Integer.parseInt(bucketNumStr.trim());
+      if (parsedBucket == -1) {
+        return Distributions.auto(Strategy.HASH, expressions);
+      }
+      return Distributions.hash(parsedBucket, expressions);
     } catch (NumberFormatException e) {
       throw new IllegalArgumentException(
-          String.format("Paimon bucket number must be a valid integer, but was 
'%s'.", bucketValue),
+          String.format(
+              "Paimon bucket number must be a valid integer, but was '%s'.", 
bucketNumStr),
           e);
     }
   }
diff --git 
a/catalogs/catalog-lakehouse-paimon/src/main/java/org/apache/gravitino/catalog/lakehouse/paimon/PaimonCatalogOperations.java
 
b/catalogs/catalog-lakehouse-paimon/src/main/java/org/apache/gravitino/catalog/lakehouse/paimon/PaimonCatalogOperations.java
index c1e7f08a30..aa2a298a62 100644
--- 
a/catalogs/catalog-lakehouse-paimon/src/main/java/org/apache/gravitino/catalog/lakehouse/paimon/PaimonCatalogOperations.java
+++ 
b/catalogs/catalog-lakehouse-paimon/src/main/java/org/apache/gravitino/catalog/lakehouse/paimon/PaimonCatalogOperations.java
@@ -29,7 +29,6 @@ import com.google.common.base.Preconditions;
 import com.google.common.collect.Maps;
 import java.time.Instant;
 import java.util.Arrays;
-import java.util.Collections;
 import java.util.List;
 import java.util.Map;
 import java.util.Optional;
@@ -354,7 +353,7 @@ public class PaimonCatalogOperations implements 
CatalogOperations, SupportsSchem
         sortOrders == null || sortOrders.length == 0,
         "Sort orders are not supported for Paimon in Gravitino.");
     checkPaimonIndexes(indexes);
-    validateDistribution(distribution, columns, indexes);
+    validateDistribution(distribution, columns);
     String currentUser = currentUser();
     GravitinoPaimonTable createdTable =
         GravitinoPaimonTable.builder()
@@ -502,7 +501,7 @@ public class PaimonCatalogOperations implements 
CatalogOperations, SupportsSchem
                     "Paimon only supports primary key Index."));
   }
 
-  private void validateDistribution(Distribution distribution, Column[] 
columns, Index[] indexes) {
+  private void validateDistribution(Distribution distribution, Column[] 
columns) {
     if (distribution == null || distribution.strategy() == 
Distributions.NONE.strategy()) {
       return;
     }
@@ -511,30 +510,21 @@ public class PaimonCatalogOperations implements 
CatalogOperations, SupportsSchem
         distribution.strategy() == Strategy.HASH,
         "Paimon only supports HASH distribution strategy.");
 
-    Preconditions.checkArgument(
-        distribution.expressions() != null && 
distribution.expressions().length > 0,
-        "Paimon bucket keys must be specified for HASH distribution.");
-
     int bucketNumber = distribution.number();
     Preconditions.checkArgument(
         bucketNumber == Distributions.AUTO || bucketNumber > 0,
         "Paimon bucket number must be positive or AUTO.");
 
-    List<String> bucketKeys = extractBucketKeys(distribution);
-    List<String> columnNames =
-        Arrays.stream(columns).map(Column::name).collect(Collectors.toList());
-    bucketKeys.forEach(
-        bucketKey ->
-            Preconditions.checkArgument(
-                columnNames.stream().anyMatch(name -> name.equals(bucketKey)),
-                "Distribution column %s does not exist in table columns.",
-                bucketKey));
-
-    List<String> primaryKeys = extractPrimaryKeys(indexes);
-    if (!primaryKeys.isEmpty()) {
-      Preconditions.checkArgument(
-          primaryKeys.containsAll(bucketKeys),
-          "Paimon bucket keys must be a subset of primary key columns for 
primary key tables.");
+    if (distribution.expressions() != null && 
distribution.expressions().length > 0) {
+      List<String> bucketKeys = extractBucketKeys(distribution);
+      List<String> columnNames =
+          
Arrays.stream(columns).map(Column::name).collect(Collectors.toList());
+      bucketKeys.forEach(
+          bucketKey ->
+              Preconditions.checkArgument(
+                  columnNames.stream().anyMatch(name -> 
name.equals(bucketKey)),
+                  "Distribution column %s does not exist in table columns.",
+                  bucketKey));
     }
   }
 
@@ -554,23 +544,6 @@ public class PaimonCatalogOperations implements 
CatalogOperations, SupportsSchem
         .collect(Collectors.toList());
   }
 
-  private static List<String> extractPrimaryKeys(Index[] indexes) {
-    if (indexes == null || indexes.length == 0) {
-      return Collections.emptyList();
-    }
-    // Paimon supports at most one index; this is enforced in {@code 
checkPaimonIndexes()}.
-    Index primaryKeyIndex = indexes[0];
-    return Arrays.stream(primaryKeyIndex.fieldNames())
-        .map(
-            fieldName -> {
-              Preconditions.checkArgument(
-                  fieldName != null && fieldName.length == 1,
-                  "Paimon primary keys must be single columns.");
-              return fieldName[0];
-            })
-        .collect(Collectors.toList());
-  }
-
   /**
    * Performs rename table change with the provided identifier.
    *
diff --git 
a/catalogs/catalog-lakehouse-paimon/src/test/java/org/apache/gravitino/catalog/lakehouse/paimon/TestGravitinoPaimonTable.java
 
b/catalogs/catalog-lakehouse-paimon/src/test/java/org/apache/gravitino/catalog/lakehouse/paimon/TestGravitinoPaimonTable.java
index 6133b13c83..23c2666054 100644
--- 
a/catalogs/catalog-lakehouse-paimon/src/test/java/org/apache/gravitino/catalog/lakehouse/paimon/TestGravitinoPaimonTable.java
+++ 
b/catalogs/catalog-lakehouse-paimon/src/test/java/org/apache/gravitino/catalog/lakehouse/paimon/TestGravitinoPaimonTable.java
@@ -419,42 +419,6 @@ public class TestGravitinoPaimonTable {
         exception.getMessage().contains("Paimon bucket number must be a valid 
integer"));
   }
 
-  @Test
-  void testCreatePaimonPrimaryKeyTableWithInvalidBucketKey() {
-    String paimonTableName = "test_paimon_primary_key_table_invalid_bucket";
-    NameIdentifier tableIdentifier = NameIdentifier.of(paimonSchema.name(), 
paimonTableName);
-
-    Column[] columns =
-        new Column[] {
-          fromPaimonColumn(new DataField(0, "col_1", 
DataTypes.INT().notNull(), PAIMON_COMMENT)),
-          fromPaimonColumn(new DataField(1, "col_2", 
DataTypes.STRING().notNull(), PAIMON_COMMENT))
-        };
-
-    Index[] indexes =
-        Collections.singletonList(
-                primary(
-                    PAIMON_PRIMARY_KEY_INDEX_NAME,
-                    new String[][] {new String[] {"col_2"}},
-                    Map.of()))
-            .toArray(new Index[0]);
-
-    IllegalArgumentException exception =
-        Assertions.assertThrows(
-            IllegalArgumentException.class,
-            () ->
-                paimonCatalogOperations.createTable(
-                    tableIdentifier,
-                    columns,
-                    PAIMON_COMMENT,
-                    Maps.newHashMap(),
-                    new Transform[0],
-                    Distributions.hash(2, NamedReference.field("col_1")),
-                    new SortOrder[0],
-                    indexes));
-    Assertions.assertTrue(
-        exception.getMessage().contains("bucket keys must be a subset of 
primary key columns"));
-  }
-
   @Test
   void testDropPaimonTable() {
     NameIdentifier tableIdentifier = NameIdentifier.of(paimonSchema.name(), 
genRandomName());
diff --git a/docs/flink-connector/flink-catalog-paimon.md 
b/docs/flink-connector/flink-catalog-paimon.md
index e994233b3b..ebb3780a7d 100644
--- a/docs/flink-connector/flink-catalog-paimon.md
+++ b/docs/flink-connector/flink-catalog-paimon.md
@@ -12,6 +12,13 @@ This document provides a comprehensive guide on configuring 
and using Apache Gra
 ### Supported Paimon Table Types
 
 * AppendOnly Table
+* Primary Key Table (with bucket distribution)
+
+### Supported Distribution
+
+* HASH distribution via `bucket-key` and `bucket` table properties.
+* Only HASH strategy is supported. Range or other strategies are not 
applicable.
+* When `bucket-key` is specified without `bucket`, the bucket number defaults 
to auto.
 
 ### Supported Operation Types
 
@@ -96,6 +103,21 @@ SELECT * FROM paimon_table_a;
 -- 1 row in set
 ```
 
+#### Distribution Example
+
+```sql
+-- Create a primary key table with HASH distribution on the 'id' column with 4 
buckets
+-- The distribution metadata is persisted in Gravitino and can be verified via 
the Gravitino API or client.
+CREATE TABLE paimon_bucketed_table (
+    id BIGINT,
+    name STRING,
+    PRIMARY KEY (id) NOT ENFORCED
+) WITH (
+    'bucket-key' = 'id',
+    'bucket' = '4'
+);
+```
+
 ## Catalog properties
 
 Gravitino Flink connector will transform below property names which are 
defined in catalog properties to Flink Paimon connector configuration.
diff --git a/flink-connector/flink/build.gradle.kts 
b/flink-connector/flink/build.gradle.kts
index 0e3b595366..53dca3e9db 100644
--- a/flink-connector/flink/build.gradle.kts
+++ b/flink-connector/flink/build.gradle.kts
@@ -45,6 +45,7 @@ dependencies {
     exclude("org.apache.logging.log4j")
   }
   implementation(libs.guava)
+  implementation(libs.commons.lang3)
 
   compileOnly(project(":clients:client-java-runtime", configuration = 
"shadow"))
 
diff --git 
a/flink-connector/flink/src/main/java/org/apache/gravitino/flink/connector/catalog/BaseCatalog.java
 
b/flink-connector/flink/src/main/java/org/apache/gravitino/flink/connector/catalog/BaseCatalog.java
index afa0242711..02d297bce8 100644
--- 
a/flink-connector/flink/src/main/java/org/apache/gravitino/flink/connector/catalog/BaseCatalog.java
+++ 
b/flink-connector/flink/src/main/java/org/apache/gravitino/flink/connector/catalog/BaseCatalog.java
@@ -26,6 +26,7 @@ import com.google.common.collect.MapDifference;
 import com.google.common.collect.Maps;
 import java.util.Arrays;
 import java.util.Collections;
+import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.Objects;
@@ -78,6 +79,7 @@ import org.apache.gravitino.flink.connector.utils.TypeUtils;
 import org.apache.gravitino.rel.Column;
 import org.apache.gravitino.rel.Table;
 import org.apache.gravitino.rel.TableChange;
+import org.apache.gravitino.rel.expressions.distributions.Distribution;
 import org.apache.gravitino.rel.expressions.distributions.Distributions;
 import org.apache.gravitino.rel.expressions.sorts.SortOrder;
 import org.apache.gravitino.rel.expressions.transforms.Transform;
@@ -289,8 +291,9 @@ public abstract class BaseCatalog extends AbstractCatalog {
             .map(this::toGravitinoColumn)
             .toArray(Column[]::new);
     String comment = table.getComment();
+    Map<String, String> flinkOptions = table.getOptions();
     Map<String, String> properties =
-        
schemaAndTablePropertiesConverter.toGravitinoTableProperties(table.getOptions());
+        
schemaAndTablePropertiesConverter.toGravitinoTableProperties(flinkOptions);
     Transform[] partitions =
         partitionConverter.toGravitinoPartitions(((CatalogTable) 
table).getPartitionKeys());
     Index[] indices = getGrivatinoIndices(resolvedTable);
@@ -304,7 +307,7 @@ public abstract class BaseCatalog extends AbstractCatalog {
               comment,
               properties,
               partitions,
-              Distributions.NONE,
+              toGravitinoDistribution(flinkOptions),
               new SortOrder[0],
               indices);
     } catch (NoSuchSchemaException e) {
@@ -567,8 +570,10 @@ public abstract class BaseCatalog extends AbstractCatalog {
     Optional<List<String>> flinkPrimaryKey = getFlinkPrimaryKey(table);
     flinkPrimaryKey.ifPresent(builder::primaryKey);
     Map<String, String> flinkTableProperties =
-        schemaAndTablePropertiesConverter.toFlinkTableProperties(
-            catalogOptions, table.properties(), tablePath);
+        new HashMap<>(
+            schemaAndTablePropertiesConverter.toFlinkTableProperties(
+                catalogOptions, table.properties(), tablePath));
+    
flinkTableProperties.putAll(fromGravitinoDistribution(table.distribution()));
     List<String> partitionKeys = 
partitionConverter.toFlinkPartitionKeys(table.partitioning());
     return CatalogTable.of(builder.build(), table.comment(), partitionKeys, 
flinkTableProperties);
   }
@@ -734,4 +739,12 @@ public abstract class BaseCatalog extends AbstractCatalog {
   protected String catalogName() {
     return getName();
   }
+
+  protected Distribution toGravitinoDistribution(Map<String, String> 
properties) {
+    return Distributions.NONE;
+  }
+
+  protected Map<String, String> fromGravitinoDistribution(Distribution 
distribution) {
+    return Collections.emptyMap();
+  }
 }
diff --git 
a/flink-connector/flink/src/main/java/org/apache/gravitino/flink/connector/paimon/GravitinoPaimonCatalog.java
 
b/flink-connector/flink/src/main/java/org/apache/gravitino/flink/connector/paimon/GravitinoPaimonCatalog.java
index fa4fd86322..abe16a6418 100644
--- 
a/flink-connector/flink/src/main/java/org/apache/gravitino/flink/connector/paimon/GravitinoPaimonCatalog.java
+++ 
b/flink-connector/flink/src/main/java/org/apache/gravitino/flink/connector/paimon/GravitinoPaimonCatalog.java
@@ -19,7 +19,14 @@
 
 package org.apache.gravitino.flink.connector.paimon;
 
+import com.google.common.annotations.VisibleForTesting;
+import com.google.common.base.Preconditions;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.Map;
 import java.util.Optional;
+import java.util.stream.Collectors;
+import org.apache.commons.lang3.StringUtils;
 import org.apache.flink.table.catalog.AbstractCatalog;
 import org.apache.flink.table.catalog.ObjectPath;
 import org.apache.flink.table.catalog.exceptions.CatalogException;
@@ -27,9 +34,15 @@ import 
org.apache.flink.table.catalog.exceptions.TableNotExistException;
 import org.apache.flink.table.factories.CatalogFactory;
 import org.apache.flink.table.factories.Factory;
 import org.apache.gravitino.NameIdentifier;
+import org.apache.gravitino.catalog.lakehouse.paimon.PaimonConstants;
 import org.apache.gravitino.flink.connector.PartitionConverter;
 import org.apache.gravitino.flink.connector.SchemaAndTablePropertiesConverter;
 import org.apache.gravitino.flink.connector.catalog.BaseCatalog;
+import org.apache.gravitino.rel.expressions.Expression;
+import org.apache.gravitino.rel.expressions.NamedReference;
+import org.apache.gravitino.rel.expressions.distributions.Distribution;
+import org.apache.gravitino.rel.expressions.distributions.Distributions;
+import org.apache.gravitino.rel.expressions.distributions.Strategy;
 import org.apache.paimon.flink.FlinkCatalogFactory;
 import org.apache.paimon.flink.FlinkTableFactory;
 
@@ -77,4 +90,90 @@ public class GravitinoPaimonCatalog extends BaseCatalog {
   public Optional<Factory> getFactory() {
     return Optional.of(new FlinkTableFactory());
   }
+
+  @Override
+  protected Distribution toGravitinoDistribution(Map<String, String> 
properties) {
+    return getDistribution(properties);
+  }
+
+  @Override
+  protected Map<String, String> fromGravitinoDistribution(Distribution 
distribution) {
+    return distributionToProperties(distribution);
+  }
+
+  @VisibleForTesting
+  static Map<String, String> distributionToProperties(Distribution 
distribution) {
+    if (distribution == null || distribution.strategy() == Strategy.NONE) {
+      return new HashMap<>();
+    }
+    Map<String, String> properties = new HashMap<>();
+    int number = distribution.number();
+    Expression[] expressions = distribution.expressions();
+    boolean hasExpressions = expressions != null && expressions.length > 0;
+
+    if (number == Distributions.AUTO && !hasExpressions) {
+      return properties;
+    }
+
+    String bucketKey =
+        Arrays.stream(expressions)
+            .map(
+                e -> {
+                  Preconditions.checkArgument(
+                      e instanceof NamedReference,
+                      "Paimon bucket-key expressions must be NamedReference, 
but got: %s",
+                      e.getClass().getSimpleName());
+                  return ((NamedReference) e).fieldName()[0];
+                })
+            .collect(Collectors.joining(","));
+    if (StringUtils.isNotBlank(bucketKey)) {
+      properties.put(PaimonConstants.BUCKET_KEY, bucketKey);
+    }
+    properties.put(PaimonConstants.BUCKET_NUM, String.valueOf(number));
+    return properties;
+  }
+
+  @VisibleForTesting
+  static Distribution getDistribution(Map<String, String> properties) {
+    if (properties == null) {
+      return Distributions.NONE;
+    }
+
+    String bucketKeyStr = properties.get(PaimonConstants.BUCKET_KEY);
+    String bucketNumStr = properties.get(PaimonConstants.BUCKET_NUM);
+
+    boolean hasBucketKey = StringUtils.isNotBlank(bucketKeyStr);
+    boolean hasBucket = StringUtils.isNotBlank(bucketNumStr);
+
+    if (!hasBucketKey && !hasBucket) {
+      return Distributions.NONE;
+    }
+
+    Expression[] expressions = new Expression[0];
+    if (hasBucketKey) {
+      expressions =
+          Arrays.stream(bucketKeyStr.split(","))
+              .map(String::trim)
+              .filter(StringUtils::isNotBlank)
+              .map(NamedReference::field)
+              .toArray(Expression[]::new);
+    }
+
+    if (!hasBucket) {
+      return Distributions.auto(Strategy.HASH, expressions);
+    }
+
+    try {
+      int parsedBucket = Integer.parseInt(bucketNumStr.trim());
+      if (parsedBucket == -1) {
+        return Distributions.auto(Strategy.HASH, expressions);
+      }
+      return Distributions.hash(parsedBucket, expressions);
+    } catch (NumberFormatException e) {
+      throw new IllegalArgumentException(
+          String.format(
+              "Paimon bucket number must be a valid integer, but was '%s'.", 
bucketNumStr),
+          e);
+    }
+  }
 }
diff --git 
a/flink-connector/flink/src/main/java/org/apache/gravitino/flink/connector/paimon/PaimonPropertiesConverter.java
 
b/flink-connector/flink/src/main/java/org/apache/gravitino/flink/connector/paimon/PaimonPropertiesConverter.java
index 17d7001b6b..f10cd9f2c6 100644
--- 
a/flink-connector/flink/src/main/java/org/apache/gravitino/flink/connector/paimon/PaimonPropertiesConverter.java
+++ 
b/flink-connector/flink/src/main/java/org/apache/gravitino/flink/connector/paimon/PaimonPropertiesConverter.java
@@ -19,6 +19,9 @@
 
 package org.apache.gravitino.flink.connector.paimon;
 
+import java.util.HashMap;
+import java.util.Map;
+import org.apache.flink.table.catalog.ObjectPath;
 import org.apache.gravitino.catalog.lakehouse.paimon.PaimonConstants;
 import org.apache.gravitino.catalog.lakehouse.paimon.PaimonPropertiesUtils;
 import org.apache.gravitino.flink.connector.CatalogPropertiesConverter;
@@ -47,6 +50,27 @@ public class PaimonPropertiesConverter
     return PaimonPropertiesUtils.GRAVITINO_CONFIG_TO_PAIMON.get(configKey);
   }
 
+  @Override
+  public Map<String, String> toGravitinoTableProperties(Map<String, String> 
flinkProperties) {
+    Map<String, String> properties = new HashMap<>(flinkProperties);
+    properties.remove(PaimonConstants.BUCKET_KEY);
+    properties.remove(PaimonConstants.BUCKET_NUM);
+    return properties;
+  }
+
+  @Override
+  public Map<String, String> toFlinkTableProperties(
+      Map<String, String> flinkCatalogProperties,
+      Map<String, String> gravitinoTableProperties,
+      ObjectPath tablePath) {
+    Map<String, String> properties = new HashMap<>(gravitinoTableProperties);
+    // Strip bucket properties from raw table properties; Distribution is the 
single source
+    // of truth and will be merged separately via fromGravitinoDistribution().
+    properties.remove(PaimonConstants.BUCKET_KEY);
+    properties.remove(PaimonConstants.BUCKET_NUM);
+    return properties;
+  }
+
   @Override
   public String getFlinkCatalogType() {
     return GravitinoPaimonCatalogFactoryOptions.IDENTIFIER;
diff --git 
a/flink-connector/flink/src/test/java/org/apache/gravitino/flink/connector/catalog/TestBaseCatalog.java
 
b/flink-connector/flink/src/test/java/org/apache/gravitino/flink/connector/catalog/TestBaseCatalog.java
index 66f63e60ff..630d9ba81d 100644
--- 
a/flink-connector/flink/src/test/java/org/apache/gravitino/flink/connector/catalog/TestBaseCatalog.java
+++ 
b/flink-connector/flink/src/test/java/org/apache/gravitino/flink/connector/catalog/TestBaseCatalog.java
@@ -35,6 +35,7 @@ import org.apache.flink.table.catalog.TableChange;
 import org.apache.gravitino.SchemaChange;
 import org.apache.gravitino.flink.connector.PartitionConverter;
 import org.apache.gravitino.flink.connector.SchemaAndTablePropertiesConverter;
+import org.apache.gravitino.rel.expressions.distributions.Distributions;
 import org.apache.gravitino.rel.types.Types;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
@@ -122,6 +123,18 @@ public class TestBaseCatalog {
     Assertions.assertArrayEquals(expected.toArray(), tableChanges);
   }
 
+  @Test
+  public void testToGravitinoDistributionDefaultsToNone() {
+    TestableBaseCatalog catalog = new 
TestableBaseCatalog(Mockito.mock(AbstractCatalog.class));
+
+    Assertions.assertEquals(
+        Distributions.NONE,
+        catalog.toGravitinoDistribution(ImmutableMap.of("bucket-key", "id", 
"bucket", "4")));
+    Assertions.assertEquals(
+        Distributions.NONE, 
catalog.toGravitinoDistribution(Collections.emptyMap()));
+    Assertions.assertEquals(Distributions.NONE, 
catalog.toGravitinoDistribution(null));
+  }
+
   @Test
   public void testListViewsReturnsEmptyWithoutDelegation() throws Exception {
     AbstractCatalog delegate = Mockito.mock(AbstractCatalog.class);
diff --git 
a/flink-connector/flink/src/test/java/org/apache/gravitino/flink/connector/integration/test/paimon/FlinkPaimonCatalogIT.java
 
b/flink-connector/flink/src/test/java/org/apache/gravitino/flink/connector/integration/test/paimon/FlinkPaimonCatalogIT.java
index e2dd51b5ad..adb8bc836e 100644
--- 
a/flink-connector/flink/src/test/java/org/apache/gravitino/flink/connector/integration/test/paimon/FlinkPaimonCatalogIT.java
+++ 
b/flink-connector/flink/src/test/java/org/apache/gravitino/flink/connector/integration/test/paimon/FlinkPaimonCatalogIT.java
@@ -19,9 +19,17 @@
 package org.apache.gravitino.flink.connector.integration.test.paimon;
 
 import com.google.common.base.Preconditions;
+import com.google.common.collect.Lists;
+import java.util.List;
 import java.util.Map;
+import org.apache.flink.table.api.TableResult;
+import org.apache.flink.types.Row;
 import org.apache.gravitino.Catalog;
+import org.apache.gravitino.NameIdentifier;
 import org.apache.gravitino.flink.connector.integration.test.FlinkCommonIT;
+import org.apache.gravitino.rel.Table;
+import org.apache.gravitino.rel.expressions.distributions.Distributions;
+import org.apache.gravitino.rel.expressions.distributions.Strategy;
 import org.junit.jupiter.api.AfterAll;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.BeforeAll;
@@ -80,6 +88,125 @@ public abstract class FlinkPaimonCatalogIT extends 
FlinkCommonIT {
 
   protected abstract String getWarehouse();
 
+  @Test
+  public void testBucketDistributionRoundTrip() {
+    String databaseName = "test_bucket_distribution_db";
+    String tableName = "test_bucket_table";
+
+    doWithSchema(
+        currentCatalog(),
+        databaseName,
+        catalog -> {
+          sql(
+              "CREATE TABLE %s (id BIGINT, name STRING) "
+                  + "WITH ('bucket' = '4', 'bucket-key' = 'id')",
+              tableName);
+
+          Table table =
+              
catalog.asTableCatalog().loadTable(NameIdentifier.of(databaseName, tableName));
+          Assertions.assertEquals(Strategy.HASH, 
table.distribution().strategy());
+          Assertions.assertEquals(4, table.distribution().number());
+          Assertions.assertEquals(1, 
table.distribution().expressions().length);
+          Assertions.assertEquals("id", 
table.distribution().expressions()[0].toString());
+
+          TableResult showResult = sql("SHOW CREATE TABLE %s", tableName);
+          List<Row> rows = Lists.newArrayList(showResult.collect());
+          Assertions.assertEquals(1, rows.size());
+          String createTableDDL = rows.get(0).getField(0).toString();
+          Assertions.assertTrue(
+              createTableDDL.contains("'bucket' = '4'"),
+              "SHOW CREATE TABLE should contain bucket number, but was: " + 
createTableDDL);
+          Assertions.assertTrue(
+              createTableDDL.contains("'bucket-key' = 'id'"),
+              "SHOW CREATE TABLE should contain bucket-key, but was: " + 
createTableDDL);
+        },
+        true,
+        supportDropCascade());
+  }
+
+  @Test
+  public void testDynamicBucketDistributionRoundTrip() {
+    String databaseName = "test_dynamic_bucket_db";
+    String tableName = "test_dynamic_bucket_table";
+
+    doWithSchema(
+        currentCatalog(),
+        databaseName,
+        catalog -> {
+          sql(
+              "CREATE TABLE %s (id BIGINT, name STRING, PRIMARY KEY (id) NOT 
ENFORCED) "
+                  + "WITH ('bucket' = '-1')",
+              tableName);
+
+          Table table =
+              
catalog.asTableCatalog().loadTable(NameIdentifier.of(databaseName, tableName));
+          Assertions.assertEquals(Strategy.HASH, 
table.distribution().strategy());
+          Assertions.assertEquals(Distributions.AUTO, 
table.distribution().number());
+          Assertions.assertEquals(0, 
table.distribution().expressions().length);
+        },
+        true,
+        supportDropCascade());
+  }
+
+  @Test
+  public void testBucketKeyWithDynamicBucketNumRejected() {
+    String databaseName = "test_bucket_key_dynamic_rejected_db";
+    String tableName = "test_rejected_table";
+
+    // Dynamic bucket mode ('-1') does not accept bucket-ley statement
+    // ref:
+    // 
https://github.com/apache/paimon/blob/dd2273f70d2f5298a3a35a557c6b462f961e3647/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java#L568-L572
+    doWithSchema(
+        currentCatalog(),
+        databaseName,
+        catalog -> {
+          Exception exception =
+              Assertions.assertThrows(
+                  Exception.class,
+                  () ->
+                      sql(
+                          "CREATE TABLE %s (id BIGINT, name STRING) "
+                              + "WITH ('bucket' = '-1', 'bucket-key' = 'id')",
+                          tableName));
+          Assertions.assertTrue(
+              exception.getMessage().contains("bucket-key")
+                  || exception.getCause().getMessage().contains("bucket-key"),
+              "Error should mention bucket-key, but was: " + 
exception.getMessage());
+        },
+        true,
+        supportDropCascade());
+  }
+
+  @Test
+  public void testOnlyBucketKeyWithoutNoBucketNumRejected() {
+    String databaseName = "test_only_bucket_key_rejected_db";
+    String tableName = "test_rejected_table";
+
+    doWithSchema(
+        currentCatalog(),
+        databaseName,
+        catalog -> {
+          // Only bucket-key without bucket maps to auto(HASH, [id]) which sets
+          // bucket=-1 and bucket-key=id in Paimon options. Paimon rejects 
this.
+          // ref:
+          // 
https://github.com/apache/paimon/blob/dd2273f70d2f5298a3a35a557c6b462f961e3647/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java#L568-L572
+          Exception exception =
+              Assertions.assertThrows(
+                  Exception.class,
+                  () ->
+                      sql(
+                          "CREATE TABLE %s (id BIGINT, name STRING) "
+                              + "WITH ('bucket-key' = 'id')",
+                          tableName));
+          Assertions.assertTrue(
+              exception.getMessage().contains("bucket-key")
+                  || exception.getCause().getMessage().contains("bucket-key"),
+              "Error should mention bucket-key, but was: " + 
exception.getMessage());
+        },
+        true,
+        supportDropCascade());
+  }
+
   @Test
   public void testCreateGravitinoPaimonCatalogUsingSQL() {
     tableEnv.useCatalog(DEFAULT_CATALOG);
diff --git 
a/flink-connector/flink/src/test/java/org/apache/gravitino/flink/connector/paimon/TestPaimonPropertiesConverter.java
 
b/flink-connector/flink/src/test/java/org/apache/gravitino/flink/connector/paimon/TestPaimonPropertiesConverter.java
index 7ead6999a4..c59135b64b 100644
--- 
a/flink-connector/flink/src/test/java/org/apache/gravitino/flink/connector/paimon/TestPaimonPropertiesConverter.java
+++ 
b/flink-connector/flink/src/test/java/org/apache/gravitino/flink/connector/paimon/TestPaimonPropertiesConverter.java
@@ -19,10 +19,15 @@
 package org.apache.gravitino.flink.connector.paimon;
 
 import com.google.common.collect.ImmutableMap;
+import java.util.HashMap;
 import java.util.Map;
 import org.apache.flink.configuration.Configuration;
 import org.apache.gravitino.catalog.lakehouse.paimon.PaimonConstants;
 import org.apache.gravitino.flink.connector.CatalogPropertiesConverter;
+import org.apache.gravitino.rel.expressions.NamedReference;
+import org.apache.gravitino.rel.expressions.distributions.Distribution;
+import org.apache.gravitino.rel.expressions.distributions.Distributions;
+import org.apache.gravitino.rel.expressions.distributions.Strategy;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
 
@@ -98,4 +103,179 @@ public class TestPaimonPropertiesConverter {
     Assertions.assertEquals(testUri, properties.get(PaimonConstants.URI));
     Assertions.assertEquals(testBackend, 
properties.get(PaimonConstants.CATALOG_BACKEND));
   }
+
+  @Test
+  public void testToGravitinoTablePropertiesStripesBucketProperties() {
+    Map<String, String> flinkProperties = new HashMap<>();
+    flinkProperties.put(PaimonConstants.BUCKET_KEY, "id");
+    flinkProperties.put(PaimonConstants.BUCKET_NUM, "4");
+    flinkProperties.put("some-other-key", "some-value");
+
+    Map<String, String> result = 
CONVERTER.toGravitinoTableProperties(flinkProperties);
+
+    Assertions.assertFalse(result.containsKey(PaimonConstants.BUCKET_KEY));
+    Assertions.assertFalse(result.containsKey(PaimonConstants.BUCKET_NUM));
+    Assertions.assertEquals("some-value", result.get("some-other-key"));
+  }
+
+  @Test
+  public void testToGravitinoTablePropertiesWithoutBucketProperties() {
+    Map<String, String> flinkProperties = ImmutableMap.of("some-key", 
"some-value");
+
+    Map<String, String> result = 
CONVERTER.toGravitinoTableProperties(flinkProperties);
+
+    Assertions.assertFalse(result.containsKey(PaimonConstants.BUCKET_KEY));
+    Assertions.assertFalse(result.containsKey(PaimonConstants.BUCKET_NUM));
+    Assertions.assertEquals("some-value", result.get("some-key"));
+  }
+
+  @Test
+  public void testGetDistributionWithBlankBucketKey() {
+    Map<String, String> options = ImmutableMap.of(PaimonConstants.BUCKET_KEY, 
"  ");
+    Distribution distribution = 
GravitinoPaimonCatalog.getDistribution(options);
+    Assertions.assertEquals(Distributions.NONE, distribution);
+  }
+
+  @Test
+  public void testGetDistributionWithNullProperties() {
+    Distribution distribution = GravitinoPaimonCatalog.getDistribution(null);
+    Assertions.assertEquals(Distributions.NONE, distribution);
+  }
+
+  @Test
+  public void testGetDistributionWithNoBucketKeyOrBucket() {
+    Map<String, String> options = ImmutableMap.of();
+    Distribution distribution = 
GravitinoPaimonCatalog.getDistribution(options);
+    Assertions.assertEquals(Distributions.NONE, distribution);
+  }
+
+  @Test
+  public void testGetDistributionWithBothBucketKeyAndBucket() {
+    Map<String, String> options =
+        ImmutableMap.of(PaimonConstants.BUCKET_KEY, "col_1,col_2", 
PaimonConstants.BUCKET_NUM, "4");
+    Distribution distribution = 
GravitinoPaimonCatalog.getDistribution(options);
+    Assertions.assertEquals(Strategy.HASH, distribution.strategy());
+    Assertions.assertEquals(4, distribution.number());
+    Assertions.assertEquals(2, distribution.expressions().length);
+    Assertions.assertEquals(
+        "col_1", ((NamedReference) 
distribution.expressions()[0]).fieldName()[0]);
+    Assertions.assertEquals(
+        "col_2", ((NamedReference) 
distribution.expressions()[1]).fieldName()[0]);
+  }
+
+  @Test
+  public void testGetDistributionWithOnlyBucketKey() {
+    Map<String, String> options = ImmutableMap.of(PaimonConstants.BUCKET_KEY, 
"col_a");
+    Distribution distribution = 
GravitinoPaimonCatalog.getDistribution(options);
+    Assertions.assertEquals(Strategy.HASH, distribution.strategy());
+    Assertions.assertEquals(Distributions.AUTO, distribution.number());
+    Assertions.assertEquals(1, distribution.expressions().length);
+    Assertions.assertEquals(
+        "col_a", ((NamedReference) 
distribution.expressions()[0]).fieldName()[0]);
+  }
+
+  @Test
+  public void testGetDistributionWithOnlyBucket() {
+    Map<String, String> options = ImmutableMap.of(PaimonConstants.BUCKET_NUM, 
"8");
+    Distribution distribution = 
GravitinoPaimonCatalog.getDistribution(options);
+    Assertions.assertEquals(Strategy.HASH, distribution.strategy());
+    Assertions.assertEquals(8, distribution.number());
+    Assertions.assertEquals(0, distribution.expressions().length);
+  }
+
+  @Test
+  public void testGetDistributionWithOnlyBucketMinusOne() {
+    Map<String, String> options = ImmutableMap.of(PaimonConstants.BUCKET_NUM, 
"-1");
+    Distribution distribution = 
GravitinoPaimonCatalog.getDistribution(options);
+    Assertions.assertEquals(Strategy.HASH, distribution.strategy());
+    Assertions.assertEquals(Distributions.AUTO, distribution.number());
+    Assertions.assertEquals(0, distribution.expressions().length);
+  }
+
+  @Test
+  public void testGetDistributionWithBucketKeyAndExplicitMinusOne() {
+    Map<String, String> options =
+        ImmutableMap.of(PaimonConstants.BUCKET_KEY, "col_a", 
PaimonConstants.BUCKET_NUM, "-1");
+    Distribution distribution = 
GravitinoPaimonCatalog.getDistribution(options);
+    Assertions.assertEquals(Strategy.HASH, distribution.strategy());
+    Assertions.assertEquals(Distributions.AUTO, distribution.number());
+    Assertions.assertEquals(1, distribution.expressions().length);
+    Assertions.assertEquals(
+        "col_a", ((NamedReference) 
distribution.expressions()[0]).fieldName()[0]);
+  }
+
+  @Test
+  public void testGetDistributionWithInvalidBucketNumber() {
+    Map<String, String> options =
+        ImmutableMap.of(
+            PaimonConstants.BUCKET_KEY, "col_1", PaimonConstants.BUCKET_NUM, 
"not_a_number");
+    IllegalArgumentException exception =
+        Assertions.assertThrows(
+            IllegalArgumentException.class, () -> 
GravitinoPaimonCatalog.getDistribution(options));
+    Assertions.assertTrue(
+        exception.getMessage().contains("Paimon bucket number must be a valid 
integer"));
+  }
+
+  @Test
+  public void testGetDistributionWithNegativeBucketNumberPassesThrough() {
+    Map<String, String> options =
+        ImmutableMap.of(PaimonConstants.BUCKET_KEY, "col_1", 
PaimonConstants.BUCKET_NUM, "-4");
+    Distribution distribution = 
GravitinoPaimonCatalog.getDistribution(options);
+    Assertions.assertEquals(Strategy.HASH, distribution.strategy());
+    Assertions.assertEquals(-4, distribution.number());
+    Assertions.assertEquals(1, distribution.expressions().length);
+  }
+
+  @Test
+  public void testDistributionToPropertiesWithBucketAndBucketKey() {
+    Distribution distribution = Distributions.hash(4, 
NamedReference.field("id"));
+    Map<String, String> properties = 
GravitinoPaimonCatalog.distributionToProperties(distribution);
+    Assertions.assertEquals("id", properties.get(PaimonConstants.BUCKET_KEY));
+    Assertions.assertEquals("4", properties.get(PaimonConstants.BUCKET_NUM));
+    Assertions.assertEquals(2, properties.size());
+  }
+
+  @Test
+  public void testDistributionToPropertiesWithNoDistribution() {
+    Map<String, String> properties =
+        GravitinoPaimonCatalog.distributionToProperties(Distributions.NONE);
+    Assertions.assertTrue(properties.isEmpty());
+  }
+
+  @Test
+  public void testDistributionToPropertiesWithNullDistribution() {
+    Map<String, String> properties = 
GravitinoPaimonCatalog.distributionToProperties(null);
+    Assertions.assertTrue(properties.isEmpty());
+  }
+
+  @Test
+  public void testDistributionToPropertiesWithAutoDistribution() {
+    Distribution distribution = Distributions.auto(Strategy.HASH, 
NamedReference.field("col_a"));
+    Map<String, String> properties = 
GravitinoPaimonCatalog.distributionToProperties(distribution);
+    Assertions.assertEquals("col_a", 
properties.get(PaimonConstants.BUCKET_KEY));
+    Assertions.assertEquals("-1", properties.get(PaimonConstants.BUCKET_NUM));
+    Assertions.assertEquals(2, properties.size());
+  }
+
+  @Test
+  public void testDistributionToPropertiesWithAutoNoExpressions() {
+    // AUTO with no expressions is Paimon's default — nothing to output.
+    Distribution distribution = Distributions.auto(Strategy.HASH);
+    Map<String, String> properties = 
GravitinoPaimonCatalog.distributionToProperties(distribution);
+    Assertions.assertTrue(properties.isEmpty());
+  }
+
+  @Test
+  public void testDistributionRoundTripIdempotent() {
+    Map<String, String> optionsBlank = 
ImmutableMap.of(PaimonConstants.BUCKET_KEY, "col_a");
+    Map<String, String> optionsExplicit =
+        ImmutableMap.of(PaimonConstants.BUCKET_KEY, "col_a", 
PaimonConstants.BUCKET_NUM, "-1");
+    Map<String, String> fromBlank =
+        GravitinoPaimonCatalog.distributionToProperties(
+            GravitinoPaimonCatalog.getDistribution(optionsBlank));
+    Map<String, String> fromExplicit =
+        GravitinoPaimonCatalog.distributionToProperties(
+            GravitinoPaimonCatalog.getDistribution(optionsExplicit));
+    Assertions.assertEquals(fromBlank, fromExplicit);
+  }
 }


Reply via email to