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]