FrankChen021 commented on code in PR #20379:
URL: https://github.com/apache/druid/pull/20379#discussion_r4046649394


##########
indexing-service/src/main/java/org/apache/druid/indexing/common/task/CompactionTask.java:
##########
@@ -798,16 +816,36 @@ private static DataSchema createDataSchema(
       @Nonnull ClientCompactionTaskGranularitySpec granularitySpec,
       @Nullable List<AggregateProjectionSpec> projections,
       @Nullable BaseTableProjectionSpec baseTable,
+      boolean sealed,
       boolean needMultiValuedColumns
   )
   {
     if (baseTable != null) {
       // BaseTable mode: the spec owns dimensions/metrics and (via its 
query-granularity virtual column) query
       // granularity; segment granularity + intervals are the compaction 
config's. Attach the config's query
       // granularity up front (a no-op if absent/NONE); the generator's 
incremental index applies it at build time.
-      // Existing-segment analysis (dims/metrics/rollup inference) is bypassed 
— the spec is authoritative.
-      final BaseTableProjectionSpec baseTableSpec =
-          
baseTable.withQueryGranularity(granularitySpec.getQueryGranularity());
+      // A sealed spec is the complete declaration of the table, so 
existing-segment analysis has nothing to add and is
+      // skipped entirely.
+      BaseTableProjectionSpec baseTableSpec = 
baseTable.withQueryGranularity(granularitySpec.getQueryGranularity());
+      if (!sealed) {
+        // Otherwise the segments may carry columns the spec does not know 
about. Find them and append them after the
+        // declared columns, leaving the shape the operator asked for 
unchanged. Only dimensions are carried over; a
+        // spec that stores metrics does not exist yet, and asking the 
analyzer for them would fail compaction on
+        // segments whose aggregators cannot be merged even though nothing 
would store them.
+        final ExistingSegmentAnalyzer analyzer = new ExistingSegmentAnalyzer(
+            segments,
+            false,
+            false,
+            true,
+            false,

Review Comment:
   [P1] Preserve or reject undeclared metric columns
   
   **Finding:** For an unsealed base-table compaction, the analyzer is 
configured with metric analysis disabled and only its dimension schema is 
appended to the base-table spec. If the datasource has existing segments from a 
legacy or rollup ingestion that contain metric aggregators, those metric 
columns are therefore omitted from the generated schema and silently disappear 
when compaction rewrites the segments, even though the unsealed path is 
intended to preserve undeclared columns.
   
   **Suggestion:** Either extend the base-table schema path to preserve 
existing metric columns, or detect metric-bearing input segments and fail the 
compaction rather than rewriting them without those values.



##########
server/src/main/java/org/apache/druid/server/coordinator/CatalogDataSourceCompactionConfig.java:
##########
@@ -145,12 +151,49 @@ public Integer getMaxRowsPerSegment()
     return null;
   }
 
+  /**
+   * Partitioning comes from the catalog: {@link 
DatasourceDefn#CLUSTER_KEYS_PROPERTY} becomes the range partition
+   * dimensions and {@link DatasourceDefn#TARGET_SEGMENT_ROWS_PROPERTY} sizes 
the segments. A table that declares no
+   * cluster keys falls back to dynamic partitioning.
+   * <p>
+   * A {@link #getBaseTable()} layout does not replace this: its clustering 
columns group rows <em>within</em> a
+   * segment, while cluster keys range-partition <em>across</em> segments, so 
a table that declares both gets both.
+   * <p>
+   * This is not optional the way the other schema fields are: the MSQ engine 
rejects a compaction task with no
+   * partitionsSpec, and a base table layout can only be compacted by MSQ.
+   */
   @JsonIgnore
   @Nullable
   @Override
   public UserCompactionTaskQueryTuningConfig getTuningConfig()
   {
-    return null;
+    final ResolvedTable table = catalog.resolveTable(tableId);
+    if (table == null) {
+      return null;
+    }
+    final Integer targetSegmentRows = 
table.decodeProperty(DatasourceDefn.TARGET_SEGMENT_ROWS_PROPERTY);
+    final List<ClusterKeySpec> clusterKeys = getClusterKeys(table);
+
+    final PartitionsSpec partitionsSpec;
+    if (CollectionUtils.isNullOrEmpty(clusterKeys)) {
+      partitionsSpec = new DynamicPartitionsSpec(targetSegmentRows, null);
+    } else {
+      // Range partitioning is ascending only, DatasourceDefn.ClusterKeysDefn 
already rejects a descending cluster key,
+      // so a catalog table cannot declare one.
+      partitionsSpec = new DimensionRangePartitionsSpec(

Review Comment:
   [P2] Validate catalog cluster-key types before scheduling
   
   **Finding:** Every catalog cluster key is converted to a range partition 
without checking the resolved catalog column type or whether it is 
multi-valued. For a plain catalog table with no baseTable, configuration 
validation has no dimensions spec to inspect, so numeric or array keys are 
accepted by the coordinator; the MSQ task then rejects the range partition at 
execution time and auto-compaction repeatedly submits failing tasks instead of 
rejecting the configuration up front.
   
   **Suggestion:** Validate each catalog cluster key against the resolved 
column schema as a single-valued VARCHAR during compaction-config validation, 
and return a configuration error before any task is scheduled.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to