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

clintropolis pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/druid.git


The following commit(s) were added to refs/heads/master by this push:
     new ea690003d09 feat: add support for specifying DimensionSchema to 
catalog clustered segment spec (#19770)
ea690003d09 is described below

commit ea690003d09f5f7ef4376bd2b65ab07da2c2f822
Author: Clint Wylie <[email protected]>
AuthorDate: Tue Jul 28 15:44:52 2026 -0700

    feat: add support for specifying DimensionSchema to catalog clustered 
segment spec (#19770)
---
 .../catalog/CatalogIngestAndQueryTest.java         |   2 +-
 .../org/apache/druid/msq/exec/MSQInsertTest.java   |   9 +-
 .../ClusteredValueGroupsBaseTableMetadata.java     | 136 +++++++++++-
 .../ClusteredValueGroupsBaseTableMetadataTest.java | 234 ++++++++++++++++++++-
 .../catalog/model/table/DatasourceTableTest.java   |  13 +-
 .../calcite/CalciteCatalogIngestionDmlTest.java    |   7 +
 6 files changed, 384 insertions(+), 17 deletions(-)

diff --git 
a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/catalog/CatalogIngestAndQueryTest.java
 
b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/catalog/CatalogIngestAndQueryTest.java
index 1dcca2c4a8f..d1078439d73 100644
--- 
a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/catalog/CatalogIngestAndQueryTest.java
+++ 
b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/catalog/CatalogIngestAndQueryTest.java
@@ -460,7 +460,7 @@ public abstract class CatalogIngestAndQueryTest extends 
CatalogTestBase
         .sealed(true)
         .property(
             DatasourceDefn.BASE_TABLE_PROPERTY,
-            new 
ClusteredValueGroupsBaseTableMetadata(ImmutableList.of("varchar_col2"), null)
+            new 
ClusteredValueGroupsBaseTableMetadata(ImmutableList.of("varchar_col2"), null, 
null)
         )
         .property(
             DatasourceDefn.PROJECTIONS_KEYS_PROPERTY,
diff --git 
a/multi-stage-query/src/test/java/org/apache/druid/msq/exec/MSQInsertTest.java 
b/multi-stage-query/src/test/java/org/apache/druid/msq/exec/MSQInsertTest.java
index 8e00d2780b9..f425fdfbd0e 100644
--- 
a/multi-stage-query/src/test/java/org/apache/druid/msq/exec/MSQInsertTest.java
+++ 
b/multi-stage-query/src/test/java/org/apache/druid/msq/exec/MSQInsertTest.java
@@ -69,6 +69,7 @@ import 
org.apache.druid.query.groupby.orderby.DefaultLimitSpec;
 import org.apache.druid.query.groupby.orderby.OrderByColumnSpec;
 import org.apache.druid.query.spec.MultipleIntervalSegmentSpec;
 import org.apache.druid.segment.AggregateProjectionMetadata;
+import org.apache.druid.segment.AutoTypeColumnSchema;
 import org.apache.druid.segment.column.ColumnType;
 import org.apache.druid.segment.column.RowSignature;
 import org.apache.druid.segment.column.ValueType;
@@ -248,7 +249,13 @@ public class MSQInsertTest extends MSQTestBase
                     .column("delta", Columns.SQL_BIGINT)
                     .sealed(true)
                     .baseTable(
-                        new 
ClusteredValueGroupsBaseTableMetadata(ImmutableList.of("channel"), null)
+                        new ClusteredValueGroupsBaseTableMetadata(
+                            ImmutableList.of("channel"),
+                            null,
+                            // customize the segment-creation schema of a 
declared column: 'user' is stored as an
+                            // auto column cast to its declared VARCHAR type 
(the logical schema is unchanged)
+                            ImmutableList.of(new AutoTypeColumnSchema("user", 
ColumnType.STRING, null))
+                        )
                     )
                     .property(
                         DatasourceDefn.PROJECTIONS_KEYS_PROPERTY,
diff --git 
a/server/src/main/java/org/apache/druid/catalog/model/ClusteredValueGroupsBaseTableMetadata.java
 
b/server/src/main/java/org/apache/druid/catalog/model/ClusteredValueGroupsBaseTableMetadata.java
index e84cc853c63..f0a3321285d 100644
--- 
a/server/src/main/java/org/apache/druid/catalog/model/ClusteredValueGroupsBaseTableMetadata.java
+++ 
b/server/src/main/java/org/apache/druid/catalog/model/ClusteredValueGroupsBaseTableMetadata.java
@@ -26,16 +26,20 @@ import com.fasterxml.jackson.annotation.JsonTypeName;
 import 
org.apache.druid.data.input.impl.ClusteredValueGroupsBaseTableProjectionSpec;
 import org.apache.druid.data.input.impl.DimensionSchema;
 import org.apache.druid.error.InvalidInput;
+import org.apache.druid.segment.AutoTypeColumnSchema;
 import org.apache.druid.segment.NestedDataColumnSchema;
 import org.apache.druid.segment.VirtualColumns;
+import org.apache.druid.segment.column.ColumnHolder;
 import org.apache.druid.segment.column.ColumnType;
 import org.apache.druid.utils.CollectionUtils;
 
 import javax.annotation.Nullable;
 import java.util.ArrayList;
 import java.util.Collections;
+import java.util.HashMap;
 import java.util.HashSet;
 import java.util.List;
+import java.util.Map;
 import java.util.Objects;
 import java.util.Set;
 
@@ -46,6 +50,10 @@ import java.util.Set;
  * (column names, types, and order) is taken from the catalog column list by 
{@link #createSpec(List)}. The declared
  * column order is the physical segment order, so, mirroring the physical 
spec, the clustering columns must be
  * declared as the leading prefix of the column list.
+ * <p>
+ * The optional {@link #columnSchemas} do NOT define columns, the catalog 
column list remains the logical schema,
+ * they are per-column customizations of the exact {@link DimensionSchema} 
used to create segments for a declared
+ * column, replacing the default schema derived from the declared column type.
  */
 @JsonTypeName(ClusteredValueGroupsBaseTableMetadata.TYPE_NAME)
 public class ClusteredValueGroupsBaseTableMetadata implements 
DatasourceBaseTableMetadata
@@ -54,15 +62,18 @@ public class ClusteredValueGroupsBaseTableMetadata 
implements DatasourceBaseTabl
 
   private final List<String> clusteringColumns;
   private final VirtualColumns virtualColumns;
+  private final List<DimensionSchema> columnSchemas;
 
   @JsonCreator
   public ClusteredValueGroupsBaseTableMetadata(
       @JsonProperty("clusteringColumns") List<String> clusteringColumns,
-      @JsonProperty("virtualColumns") @Nullable VirtualColumns virtualColumns
+      @JsonProperty("virtualColumns") @Nullable VirtualColumns virtualColumns,
+      @JsonProperty("columnSchemas") @Nullable List<DimensionSchema> 
columnSchemas
   )
   {
     this.clusteringColumns = clusteringColumns == null ? 
Collections.emptyList() : clusteringColumns;
     this.virtualColumns = virtualColumns == null ? VirtualColumns.EMPTY : 
virtualColumns;
+    this.columnSchemas = columnSchemas == null ? Collections.emptyList() : 
columnSchemas;
   }
 
   @Override
@@ -86,12 +97,30 @@ public class ClusteredValueGroupsBaseTableMetadata 
implements DatasourceBaseTabl
     return virtualColumns;
   }
 
+  /**
+   * Per-column customizations of the {@link DimensionSchema} used during 
segment creation, keyed by
+   * {@link DimensionSchema#getName()}; empty when every declared column uses 
the schema derived from its declared
+   * type. These do not define columns: every entry must customize a declared, 
non-clustering column.
+   */
+  @JsonProperty("columnSchemas")
+  @JsonInclude(JsonInclude.Include.NON_EMPTY)
+  public List<DimensionSchema> getColumnSchemas()
+  {
+    return columnSchemas;
+  }
+
   /**
    * Creates the physical spec from the declared catalog columns, used 
verbatim: the declared column order is the
    * physical segment order, so the clustering columns must be declared as the 
leading prefix of the column list (in
    * {@link #clusteringColumns} order); anything else is a validation error. 
Every clustering column must be a declared
    * column, a clustering column computed by a virtual column at ingest time 
is still a stored, queryable column, so it
    * too must appear in the column list. All ordering and layout rules are 
enforced by the spec itself.
+   * <p>
+   * A column with an entry in {@link #columnSchemas} uses that {@link 
DimensionSchema} verbatim in place of the
+   * default derived from its declared type. Customizations may only target 
declared, non-clustering columns (a
+   * clustering column's physical representation is fixed by the clustered 
format, and {@code __time} is always a
+   * long), and the schema's type must match the declared logical type so the 
physical schema cannot silently
+   * contradict the SQL schema.
    */
   @Override
   public ClusteredValueGroupsBaseTableProjectionSpec 
createSpec(List<ColumnSpec> columns)
@@ -102,11 +131,12 @@ public class ClusteredValueGroupsBaseTableMetadata 
implements DatasourceBaseTabl
           TYPE_NAME
       );
     }
+    final Map<String, DimensionSchema> customSchemas = indexColumnSchemas();
     final Set<String> declaredNames = new HashSet<>();
     final List<DimensionSchema> specColumns = new ArrayList<>(columns.size());
     for (ColumnSpec column : columns) {
       declaredNames.add(column.name());
-      specColumns.add(toDimensionSchema(column));
+      specColumns.add(toDimensionSchema(column, 
customSchemas.get(column.name())));
     }
     for (String clusteringColumn : clusteringColumns) {
       if (!declaredNames.contains(clusteringColumn)) {
@@ -118,6 +148,16 @@ public class ClusteredValueGroupsBaseTableMetadata 
implements DatasourceBaseTabl
         );
       }
     }
+    for (String customized : customSchemas.keySet()) {
+      if (!declaredNames.contains(customized)) {
+        throw InvalidInput.exception(
+            "columnSchemas entry [%s] does not customize a declared column; 
column schemas do not define columns,"
+            + " declare [%s] in the table's column list",
+            customized,
+            customized
+        );
+      }
+    }
     return ClusteredValueGroupsBaseTableProjectionSpec.builder()
                                                       
.virtualColumns(virtualColumns)
                                                       .columns(specColumns)
@@ -125,12 +165,30 @@ public class ClusteredValueGroupsBaseTableMetadata 
implements DatasourceBaseTabl
                                                       .build();
   }
 
-  private static DimensionSchema toDimensionSchema(ColumnSpec column)
+  private Map<String, DimensionSchema> indexColumnSchemas()
+  {
+    final Map<String, DimensionSchema> customSchemas = new HashMap<>();
+    for (DimensionSchema schema : columnSchemas) {
+      if (schema == null) {
+        throw InvalidInput.exception("columnSchemas must not contain null 
entries");
+      }
+      if (customSchemas.put(schema.getName(), schema) != null) {
+        throw InvalidInput.exception("columnSchemas contains duplicate entries 
for column [%s]", schema.getName());
+      }
+    }
+    return customSchemas;
+  }
+
+  private DimensionSchema toDimensionSchema(ColumnSpec column, @Nullable 
DimensionSchema customSchema)
   {
     ColumnType druidType = Columns.druidType(column);
     if (druidType == null) {
       druidType = ColumnType.STRING;
     }
+    if (customSchema != null) {
+      validateColumnSchemaCustomization(column, customSchema, druidType);
+      return customSchema;
+    }
     if (druidType.isPrimitive() || druidType.isPrimitiveArray()) {
       // The declared type is retained in the ingestion schema (primitive 
arrays are cast, rather than left to an
       // untyped auto column whose type is inferred from the ingested values; 
note that the auto schema stores
@@ -150,6 +208,72 @@ public class ClusteredValueGroupsBaseTableMetadata 
implements DatasourceBaseTabl
     );
   }
 
+  private void validateColumnSchemaCustomization(
+      ColumnSpec column,
+      DimensionSchema customSchema,
+      ColumnType declaredType
+  )
+  {
+    if (ColumnHolder.TIME_COLUMN_NAME.equals(column.name())) {
+      throw InvalidInput.exception(
+          "columnSchemas cannot customize [%s]: the time column is always 
stored as a long",
+          ColumnHolder.TIME_COLUMN_NAME
+      );
+    }
+    if (clusteringColumns.contains(column.name())) {
+      throw InvalidInput.exception(
+          "columnSchemas cannot customize clustering column [%s]: the physical 
representation of clustering columns"
+          + " is fixed by the clustered segment format",
+          column.name()
+      );
+    }
+    // The schema's type must match the declared logical type, so the physical 
schema cannot silently contradict the
+    // SQL schema that INSERT/REPLACE queries are validated and coerced 
against.
+    ColumnType expectedType = declaredType;
+    if (customSchema instanceof AutoTypeColumnSchema) {
+      // An uncast auto column stores values as they are ingested (inferring 
the physical type) rather than coercing
+      // them to the declared type; only a column declared COMPLEX<json> may 
store arbitrary shapes.
+      if (((AutoTypeColumnSchema) customSchema).getCastToType() == null
+          && !ColumnType.NESTED_DATA.equals(declaredType)) {
+        throw InvalidInput.exception(
+            "columnSchemas entry [%s] is an auto column schema without a 
castToType; an uncast auto column stores"
+            + " values as they are ingested rather than coercing them to the 
column's declared type [%s], set"
+            + " castToType to match the declared type",
+            column.name(),
+            declaredType
+        );
+      }
+      // The auto schema stores FLOAT as DOUBLE.
+      expectedType = autoColumnType(declaredType);
+    }
+    // A NestedDataColumnSchema always reports COMPLEX<json>, so this also 
restricts json schemas to columns declared
+    // as COMPLEX<json>.
+    if (!expectedType.equals(customSchema.getColumnType())) {
+      throw InvalidInput.exception(
+          "columnSchemas entry [%s] of type [%s] does not match the column's 
declared type [%s]; column schemas"
+          + " customize the physical representation of a declared column, not 
its type",
+          column.name(),
+          customSchema.getColumnType(),
+          declaredType
+      );
+    }
+  }
+
+  /**
+   * The type the auto schema stores for a declared type: {@link 
AutoTypeColumnSchema} coerces FLOAT to DOUBLE (the
+   * default derivation for declared FLOAT ARRAY columns relies on the same 
coercion).
+   */
+  private static ColumnType autoColumnType(ColumnType declaredType)
+  {
+    if (ColumnType.FLOAT.equals(declaredType)) {
+      return ColumnType.DOUBLE;
+    }
+    if (ColumnType.FLOAT_ARRAY.equals(declaredType)) {
+      return ColumnType.DOUBLE_ARRAY;
+    }
+    return declaredType;
+  }
+
   @Override
   public boolean equals(Object o)
   {
@@ -158,13 +282,14 @@ public class ClusteredValueGroupsBaseTableMetadata 
implements DatasourceBaseTabl
     }
     ClusteredValueGroupsBaseTableMetadata that = 
(ClusteredValueGroupsBaseTableMetadata) o;
     return Objects.equals(clusteringColumns, that.clusteringColumns)
-           && Objects.equals(virtualColumns, that.virtualColumns);
+           && Objects.equals(virtualColumns, that.virtualColumns)
+           && Objects.equals(columnSchemas, that.columnSchemas);
   }
 
   @Override
   public int hashCode()
   {
-    return Objects.hash(clusteringColumns, virtualColumns);
+    return Objects.hash(clusteringColumns, virtualColumns, columnSchemas);
   }
 
   @Override
@@ -173,6 +298,7 @@ public class ClusteredValueGroupsBaseTableMetadata 
implements DatasourceBaseTabl
     return "ClusteredValueGroupsBaseTableMetadata{" +
            "clusteringColumns=" + clusteringColumns +
            ", virtualColumns=" + virtualColumns +
+           ", columnSchemas=" + columnSchemas +
            '}';
   }
 }
diff --git 
a/server/src/test/java/org/apache/druid/catalog/model/ClusteredValueGroupsBaseTableMetadataTest.java
 
b/server/src/test/java/org/apache/druid/catalog/model/ClusteredValueGroupsBaseTableMetadataTest.java
index af3a956970a..e7891dd9c2a 100644
--- 
a/server/src/test/java/org/apache/druid/catalog/model/ClusteredValueGroupsBaseTableMetadataTest.java
+++ 
b/server/src/test/java/org/apache/druid/catalog/model/ClusteredValueGroupsBaseTableMetadataTest.java
@@ -24,6 +24,7 @@ import com.fasterxml.jackson.databind.InjectableValues;
 import com.fasterxml.jackson.databind.ObjectMapper;
 import nl.jqno.equalsverifier.EqualsVerifier;
 import 
org.apache.druid.data.input.impl.ClusteredValueGroupsBaseTableProjectionSpec;
+import org.apache.druid.data.input.impl.DimensionSchema;
 import org.apache.druid.data.input.impl.DoubleDimensionSchema;
 import org.apache.druid.data.input.impl.LongDimensionSchema;
 import org.apache.druid.data.input.impl.StringDimensionSchema;
@@ -31,6 +32,7 @@ import org.apache.druid.error.DruidException;
 import org.apache.druid.jackson.DefaultObjectMapper;
 import org.apache.druid.math.expr.ExprMacroTable;
 import org.apache.druid.segment.AutoTypeColumnSchema;
+import org.apache.druid.segment.DefaultColumnFormatConfig;
 import org.apache.druid.segment.NestedDataColumnSchema;
 import org.apache.druid.segment.VirtualColumns;
 import org.apache.druid.segment.column.ColumnType;
@@ -47,7 +49,9 @@ import java.util.Map;
 public class ClusteredValueGroupsBaseTableMetadataTest extends 
InitializedNullHandlingTest
 {
   private final ObjectMapper mapper = new 
DefaultObjectMapper().setInjectableValues(
-      new InjectableValues.Std().addValue(ExprMacroTable.class, 
ExprMacroTable.nil())
+      new InjectableValues.Std()
+          .addValue(ExprMacroTable.class, ExprMacroTable.nil())
+          .addValue(DefaultColumnFormatConfig.class, new 
DefaultColumnFormatConfig(null, null, null, null))
   );
 
   // Declared order is the physical segment order: clustering columns lead, 
__time is an explicit positional column.
@@ -66,7 +70,7 @@ public class ClusteredValueGroupsBaseTableMetadataTest 
extends InitializedNullHa
         Collections.singletonList("tenant_lower"),
         VirtualColumns.create(
             new ExpressionVirtualColumn("tenant_lower", "lower(\"tenant\")", 
ColumnType.STRING, ExprMacroTable.nil())
-        )
+        ), null
     );
     final String json = mapper.writeValueAsString(metadata);
     final DatasourceBaseTableMetadata fromJson = mapper.readValue(json, 
DatasourceBaseTableMetadata.class);
@@ -78,10 +82,29 @@ public class ClusteredValueGroupsBaseTableMetadataTest 
extends InitializedNullHa
   {
     final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
         Arrays.asList("tenant", "region"),
+        null,
         null
     );
     final String json = mapper.writeValueAsString(metadata);
     Assert.assertFalse(json.contains("virtualColumns"));
+    Assert.assertFalse(json.contains("columnSchemas"));
+    final DatasourceBaseTableMetadata fromJson = mapper.readValue(json, 
DatasourceBaseTableMetadata.class);
+    Assert.assertEquals(metadata, fromJson);
+  }
+
+  @Test
+  public void testSerdeWithColumnSchemas() throws Exception
+  {
+    final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
+        Collections.singletonList("tenant"),
+        null,
+        Arrays.asList(
+            new StringDimensionSchema("region", 
DimensionSchema.MultiValueHandling.ARRAY, false),
+            new AutoTypeColumnSchema("value", ColumnType.DOUBLE, null),
+            new NestedDataColumnSchema("attrs", 
NestedDataColumnSchema.DEFAULT_FORMAT_VERSION)
+        )
+    );
+    final String json = mapper.writeValueAsString(metadata);
     final DatasourceBaseTableMetadata fromJson = mapper.readValue(json, 
DatasourceBaseTableMetadata.class);
     Assert.assertEquals(metadata, fromJson);
   }
@@ -93,6 +116,7 @@ public class ClusteredValueGroupsBaseTableMetadataTest 
extends InitializedNullHa
     // the type discriminator must survive that path (it is an 
EXISTING_PROPERTY for this reason).
     final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
         Collections.singletonList("tenant"),
+        null,
         null
     );
     final String json = 
mapper.writeValueAsString(Collections.singletonMap("baseTable", metadata));
@@ -109,6 +133,7 @@ public class ClusteredValueGroupsBaseTableMetadataTest 
extends InitializedNullHa
   {
     final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
         Collections.singletonList("tenant"),
+        null,
         null
     );
     // The declared column order is the physical segment order, used verbatim; 
types map through Columns.druidType
@@ -139,7 +164,8 @@ public class ClusteredValueGroupsBaseTableMetadataTest 
extends InitializedNullHa
     );
     final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
         Collections.singletonList("tenant_lower"),
-        virtualColumns
+        virtualColumns,
+        null
     );
     final List<ColumnSpec> columns = Arrays.asList(
         new ColumnSpec("tenant_lower", Columns.SQL_VARCHAR, null),
@@ -162,11 +188,203 @@ public class ClusteredValueGroupsBaseTableMetadataTest 
extends InitializedNullHa
     );
   }
 
+  @Test
+  public void testCreateSpecWithColumnSchemas()
+  {
+    // Column schemas customize the exact DimensionSchema used at segment 
creation without defining columns: the
+    // declared column list remains the logical schema. 'region' keeps its 
declared STRING type but customizes
+    // multi-value handling and disables the bitmap index; 'value' is stored 
as an auto column cast to its declared
+    // DOUBLE type.
+    final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
+        Collections.singletonList("tenant"),
+        null,
+        Arrays.asList(
+            new StringDimensionSchema("region", 
DimensionSchema.MultiValueHandling.ARRAY, false),
+            new AutoTypeColumnSchema("value", ColumnType.DOUBLE, null)
+        )
+    );
+    Assert.assertEquals(
+        ClusteredValueGroupsBaseTableProjectionSpec.builder()
+                                                   .columns(
+                                                       new 
StringDimensionSchema("tenant"),
+                                                       new 
LongDimensionSchema(Columns.TIME_COLUMN),
+                                                       new 
StringDimensionSchema("region", DimensionSchema.MultiValueHandling.ARRAY, 
false),
+                                                       new 
LongDimensionSchema("delta"),
+                                                       new 
AutoTypeColumnSchema("value", ColumnType.DOUBLE, null)
+                                                   )
+                                                   .clusteringColumns("tenant")
+                                                   .build(),
+        metadata.createSpec(COLUMNS)
+    );
+  }
+
+  @Test
+  public void testCreateSpecAutoColumnSchemaCoercions()
+  {
+    // A declared FLOAT column may be customized with an auto schema cast to 
FLOAT: the auto schema itself stores
+    // FLOAT as DOUBLE (the same coercion the default derivation relies on for 
FLOAT arrays). A column declared
+    // COMPLEX<json> may use an uncast auto schema, since arbitrary shapes are 
that type's contract.
+    final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
+        Collections.singletonList("tenant"),
+        null,
+        Arrays.asList(
+            new AutoTypeColumnSchema("ratio", ColumnType.FLOAT, null),
+            AutoTypeColumnSchema.of("attrs")
+        )
+    );
+    final List<ColumnSpec> columns = Arrays.asList(
+        new ColumnSpec("tenant", Columns.SQL_VARCHAR, null),
+        new ColumnSpec(Columns.TIME_COLUMN, null, null),
+        new ColumnSpec("ratio", Columns.SQL_FLOAT, null),
+        new ColumnSpec("attrs", ColumnType.NESTED_DATA.asTypeString(), null)
+    );
+    Assert.assertEquals(
+        ClusteredValueGroupsBaseTableProjectionSpec.builder()
+                                                   .columns(
+                                                       new 
StringDimensionSchema("tenant"),
+                                                       new 
LongDimensionSchema(Columns.TIME_COLUMN),
+                                                       new 
AutoTypeColumnSchema("ratio", ColumnType.FLOAT, null),
+                                                       
AutoTypeColumnSchema.of("attrs")
+                                                   )
+                                                   .clusteringColumns("tenant")
+                                                   .build(),
+        metadata.createSpec(columns)
+    );
+  }
+
+  @Test
+  public void testCreateSpecUncastAutoColumnSchemaFails()
+  {
+    // An uncast auto column stores values as ingested rather than coercing to 
the declared type, so it is only
+    // legal for columns declared COMPLEX<json>.
+    final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
+        Collections.singletonList("tenant"),
+        null,
+        Collections.singletonList(AutoTypeColumnSchema.of("value"))
+    );
+    final DruidException e = Assert.assertThrows(DruidException.class, () -> 
metadata.createSpec(COLUMNS));
+    Assert.assertTrue(
+        e.getMessage().contains("columnSchemas entry [value] is an auto column 
schema without a castToType")
+    );
+  }
+
+  @Test
+  public void testCreateSpecAutoColumnSchemaCastMismatchFails()
+  {
+    final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
+        Collections.singletonList("tenant"),
+        null,
+        Collections.singletonList(new AutoTypeColumnSchema("region", 
ColumnType.LONG, null))
+    );
+    final DruidException e = Assert.assertThrows(DruidException.class, () -> 
metadata.createSpec(COLUMNS));
+    Assert.assertTrue(
+        e.getMessage().contains("columnSchemas entry [region] of type [LONG] 
does not match the column's declared type [STRING]")
+    );
+  }
+
+  @Test
+  public void testCreateSpecJsonColumnSchemaForNonJsonColumnFails()
+  {
+    // A json schema always stores COMPLEX<json>, so it may only customize 
columns declared as COMPLEX<json>.
+    final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
+        Collections.singletonList("tenant"),
+        null,
+        Collections.singletonList(new NestedDataColumnSchema("region", 
NestedDataColumnSchema.DEFAULT_FORMAT_VERSION))
+    );
+    final DruidException e = Assert.assertThrows(DruidException.class, () -> 
metadata.createSpec(COLUMNS));
+    Assert.assertTrue(
+        e.getMessage().contains(
+            "columnSchemas entry [region] of type [COMPLEX<json>] does not 
match the column's declared type [STRING]"
+        )
+    );
+  }
+
+  @Test
+  public void testCreateSpecColumnSchemaForUndeclaredColumnFails()
+  {
+    final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
+        Collections.singletonList("tenant"),
+        null,
+        Collections.singletonList(new StringDimensionSchema("no_such_column"))
+    );
+    final DruidException e = Assert.assertThrows(DruidException.class, () -> 
metadata.createSpec(COLUMNS));
+    Assert.assertTrue(
+        e.getMessage().contains("columnSchemas entry [no_such_column] does not 
customize a declared column")
+    );
+  }
+
+  @Test
+  public void testCreateSpecColumnSchemaForClusteringColumnFails()
+  {
+    final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
+        Collections.singletonList("tenant"),
+        null,
+        Collections.singletonList(new StringDimensionSchema("tenant", 
DimensionSchema.MultiValueHandling.ARRAY, false))
+    );
+    final DruidException e = Assert.assertThrows(DruidException.class, () -> 
metadata.createSpec(COLUMNS));
+    Assert.assertTrue(e.getMessage().contains("columnSchemas cannot customize 
clustering column [tenant]"));
+  }
+
+  @Test
+  public void testCreateSpecColumnSchemaForTimeColumnFails()
+  {
+    final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
+        Collections.singletonList("tenant"),
+        null,
+        Collections.singletonList(new LongDimensionSchema(Columns.TIME_COLUMN))
+    );
+    final DruidException e = Assert.assertThrows(DruidException.class, () -> 
metadata.createSpec(COLUMNS));
+    Assert.assertTrue(e.getMessage().contains("columnSchemas cannot customize 
[__time]"));
+  }
+
+  @Test
+  public void testCreateSpecColumnSchemaTypeMismatchFails()
+  {
+    // 'region' is declared as STRING; a LONG dimension schema would 
contradict the logical schema SQL writes are
+    // validated against.
+    final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
+        Collections.singletonList("tenant"),
+        null,
+        Collections.singletonList(new LongDimensionSchema("region"))
+    );
+    final DruidException e = Assert.assertThrows(DruidException.class, () -> 
metadata.createSpec(COLUMNS));
+    Assert.assertTrue(
+        e.getMessage().contains("columnSchemas entry [region] of type [LONG] 
does not match the column's declared type [STRING]")
+    );
+  }
+
+  @Test
+  public void testCreateSpecNullColumnSchemaEntryFails() throws Exception
+  {
+    final DatasourceBaseTableMetadata metadata = mapper.readValue(
+        
"{\"type\":\"clusteredValueGroups\",\"clusteringColumns\":[\"tenant\"],\"columnSchemas\":[null]}",
+        DatasourceBaseTableMetadata.class
+    );
+    final DruidException e = Assert.assertThrows(DruidException.class, () -> 
metadata.createSpec(COLUMNS));
+    Assert.assertTrue(e.getMessage().contains("columnSchemas must not contain 
null entries"));
+  }
+
+  @Test
+  public void testCreateSpecDuplicateColumnSchemasFails()
+  {
+    final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
+        Collections.singletonList("tenant"),
+        null,
+        Arrays.asList(
+            new StringDimensionSchema("region"),
+            new StringDimensionSchema("region", 
DimensionSchema.MultiValueHandling.ARRAY, false)
+        )
+    );
+    final DruidException e = Assert.assertThrows(DruidException.class, () -> 
metadata.createSpec(COLUMNS));
+    Assert.assertTrue(e.getMessage().contains("columnSchemas contains 
duplicate entries for column [region]"));
+  }
+
   @Test
   public void testCreateSpecMultipleClusteringColumns()
   {
     final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
         Arrays.asList("region", "tenant"),
+        null,
         null
     );
     final List<ColumnSpec> columns = Arrays.asList(
@@ -194,6 +412,7 @@ public class ClusteredValueGroupsBaseTableMetadataTest 
extends InitializedNullHa
   {
     final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
         Collections.singletonList("tenant"),
+        null,
         null
     );
     final List<ColumnSpec> columns = Arrays.asList(
@@ -228,6 +447,7 @@ public class ClusteredValueGroupsBaseTableMetadataTest 
extends InitializedNullHa
   {
     final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
         Collections.singletonList("tenant"),
+        null,
         null
     );
     final List<ColumnSpec> columns = Arrays.asList(
@@ -246,6 +466,7 @@ public class ClusteredValueGroupsBaseTableMetadataTest 
extends InitializedNullHa
     // physical segment order, so this is an error rather than a silent 
reorder.
     final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
         Collections.singletonList("region"),
+        null,
         null
     );
     final DruidException e = Assert.assertThrows(DruidException.class, () -> 
metadata.createSpec(COLUMNS));
@@ -257,6 +478,7 @@ public class ClusteredValueGroupsBaseTableMetadataTest 
extends InitializedNullHa
   {
     final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
         Arrays.asList("tenant", "region"),
+        null,
         null
     );
     final List<ColumnSpec> columns = Arrays.asList(
@@ -273,6 +495,7 @@ public class ClusteredValueGroupsBaseTableMetadataTest 
extends InitializedNullHa
   {
     final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
         Collections.singletonList("no_such_column"),
+        null,
         null
     );
     final DruidException e = Assert.assertThrows(DruidException.class, () -> 
metadata.createSpec(COLUMNS));
@@ -284,6 +507,7 @@ public class ClusteredValueGroupsBaseTableMetadataTest 
extends InitializedNullHa
   {
     final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
         Collections.singletonList("tenant"),
+        null,
         null
     );
     final List<ColumnSpec> columns = Arrays.asList(
@@ -299,6 +523,7 @@ public class ClusteredValueGroupsBaseTableMetadataTest 
extends InitializedNullHa
   {
     final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
         Collections.singletonList("tags"),
+        null,
         null
     );
     final List<ColumnSpec> columns = Arrays.asList(
@@ -312,7 +537,7 @@ public class ClusteredValueGroupsBaseTableMetadataTest 
extends InitializedNullHa
   @Test
   public void testCreateSpecEmptyClusteringColumnsFails()
   {
-    final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(null, null);
+    final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(null, null, null);
     final DruidException e = Assert.assertThrows(DruidException.class, () -> 
metadata.createSpec(COLUMNS));
     Assert.assertTrue(e.getMessage().contains("clusteringColumns must be 
non-empty"));
   }
@@ -322,6 +547,7 @@ public class ClusteredValueGroupsBaseTableMetadataTest 
extends InitializedNullHa
   {
     final DatasourceBaseTableMetadata metadata = new 
ClusteredValueGroupsBaseTableMetadata(
         Collections.singletonList("tenant"),
+        null,
         null
     );
     final DruidException e = Assert.assertThrows(
diff --git 
a/server/src/test/java/org/apache/druid/catalog/model/table/DatasourceTableTest.java
 
b/server/src/test/java/org/apache/druid/catalog/model/table/DatasourceTableTest.java
index b4e06cb1a01..287e9477c89 100644
--- 
a/server/src/test/java/org/apache/druid/catalog/model/table/DatasourceTableTest.java
+++ 
b/server/src/test/java/org/apache/druid/catalog/model/table/DatasourceTableTest.java
@@ -162,7 +162,7 @@ public class DatasourceTableTest
           ImmutableMap.of(
               DatasourceDefn.SEALED_PROPERTY, true,
               DatasourceDefn.BASE_TABLE_PROPERTY,
-              new 
ClusteredValueGroupsBaseTableMetadata(Collections.singletonList("tenant"), null)
+              new 
ClusteredValueGroupsBaseTableMetadata(Collections.singletonList("tenant"), 
null, null)
           ),
           columns
       );
@@ -176,7 +176,7 @@ public class DatasourceTableTest
           DatasourceDefn.TABLE_TYPE,
           ImmutableMap.of(
               DatasourceDefn.BASE_TABLE_PROPERTY,
-              new 
ClusteredValueGroupsBaseTableMetadata(Collections.singletonList("tenant"), null)
+              new 
ClusteredValueGroupsBaseTableMetadata(Collections.singletonList("tenant"), 
null, null)
           ),
           columns
       );
@@ -193,7 +193,7 @@ public class DatasourceTableTest
           ImmutableMap.of(
               DatasourceDefn.SEALED_PROPERTY, true,
               DatasourceDefn.BASE_TABLE_PROPERTY,
-              new 
ClusteredValueGroupsBaseTableMetadata(Collections.singletonList("no_such_column"),
 null)
+              new 
ClusteredValueGroupsBaseTableMetadata(Collections.singletonList("no_such_column"),
 null, null)
           ),
           columns
       );
@@ -209,7 +209,7 @@ public class DatasourceTableTest
           ImmutableMap.of(
               DatasourceDefn.SEALED_PROPERTY, true,
               DatasourceDefn.BASE_TABLE_PROPERTY,
-              new 
ClusteredValueGroupsBaseTableMetadata(Collections.singletonList("region"), null)
+              new 
ClusteredValueGroupsBaseTableMetadata(Collections.singletonList("region"), 
null, null)
           ),
           columns
       );
@@ -224,7 +224,7 @@ public class DatasourceTableTest
           ImmutableMap.of(
               DatasourceDefn.SEALED_PROPERTY, true,
               DatasourceDefn.BASE_TABLE_PROPERTY,
-              new 
ClusteredValueGroupsBaseTableMetadata(Collections.singletonList("tenant"), null)
+              new 
ClusteredValueGroupsBaseTableMetadata(Collections.singletonList("tenant"), 
null, null)
           ),
           Collections.singletonList(new ColumnSpec("tenant", 
Columns.SQL_VARCHAR, null))
       );
@@ -244,7 +244,8 @@ public class DatasourceTableTest
         Collections.singletonList("tenant_lower"),
         VirtualColumns.create(
             new ExpressionVirtualColumn("tenant_lower", "lower(\"tenant\")", 
ColumnType.STRING, ExprMacroTable.nil())
-        )
+        ),
+        null
     );
     TableSpec spec = new TableSpec(
         DatasourceDefn.TABLE_TYPE,
diff --git 
a/sql/src/test/java/org/apache/druid/sql/calcite/CalciteCatalogIngestionDmlTest.java
 
b/sql/src/test/java/org/apache/druid/sql/calcite/CalciteCatalogIngestionDmlTest.java
index fcd498450c3..02958704557 100644
--- 
a/sql/src/test/java/org/apache/druid/sql/calcite/CalciteCatalogIngestionDmlTest.java
+++ 
b/sql/src/test/java/org/apache/druid/sql/calcite/CalciteCatalogIngestionDmlTest.java
@@ -33,7 +33,9 @@ import org.apache.druid.catalog.model.facade.DatasourceFacade;
 import org.apache.druid.catalog.model.table.ClusterKeySpec;
 import org.apache.druid.catalog.model.table.DatasourceDefn;
 import org.apache.druid.data.input.impl.CsvInputFormat;
+import org.apache.druid.data.input.impl.DimensionSchema;
 import org.apache.druid.data.input.impl.InlineInputSource;
+import org.apache.druid.data.input.impl.StringDimensionSchema;
 import org.apache.druid.error.DruidException;
 import org.apache.druid.jackson.DefaultObjectMapper;
 import org.apache.druid.java.util.common.StringUtils;
@@ -402,6 +404,11 @@ public abstract class CalciteCatalogIngestionDmlTest 
extends CalciteIngestionDml
                                         ColumnType.STRING,
                                         ExprMacroTable.nil()
                                     )
+                                ),
+                                // customizes the segment-creation schema of 
'dim1'; invisible to the planner, since
+                                // the declared column list remains the 
logical schema
+                                ImmutableList.of(
+                                    new StringDimensionSchema("dim1", 
DimensionSchema.MultiValueHandling.ARRAY, false)
                                 )
                             )
                         ),


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

Reply via email to