This is an automated email from the ASF dual-hosted git repository. yiguolei pushed a commit to branch branch-4.2 in repository https://gitbox.apache.org/repos/asf/doris.git
commit 6640614be23c99cfd2b193809629de6dc645ef3d Author: morrySnow <[email protected]> AuthorDate: Wed Sep 16 21:55:44 2026 +0800 branch-4.1: [fix](statistics) Treat invalid column statistics as UNKNOWN instead of aborting analyze job (#66756) (#68061) ### What problem does this PR solve? Issue Number: close #64122 Related PR: #66756 Problem Summary: This backports #66756 to `branch-4.1`. Sample analysis can produce an invalid combination such as `ndv = 0`, non-null min/max values, and `null_count != row_count` for an almost-all-NULL column. On `branch-4.1`, this causes `BaseAnalysisTask` to throw and abort the entire analyze job. The backport keeps the analyze job running, converts invalid column statistics to `ColumnStatistic.UNKNOWN` when read, and publishes UNKNOWN into the statistics cache to avoid exposing invalid or stale values to the optimizer. The only cherry-pick conflict was in `BaseAnalysisTaskTest`; it was resolved using the existing 4.1 test structure while preserving the assertions introduced by #66756. ### Release note Treat invalid collected column statistics as UNKNOWN instead of aborting the analyze job. ### Check List (For Author) - Test: Unit Test - `./run-fe-ut.sh --run org.apache.doris.statistics.BaseAnalysisTaskTest` - Tests run: 3, failures: 0, errors: 0, skipped: 0 - Behavior changed: Yes. Invalid collected column statistics no longer abort the whole analyze job and are exposed to the optimizer as UNKNOWN. - Does this need documentation: No Co-authored-by: csding <[email protected]> --- .../apache/doris/statistics/BaseAnalysisTask.java | 9 +- .../org/apache/doris/statistics/ColStatsData.java | 4 + .../apache/doris/statistics/ColumnStatistic.java | 5 + .../apache/doris/statistics/StatisticsCache.java | 11 +- .../doris/statistics/BaseAnalysisTaskTest.java | 47 +------ .../stats/invalid_stats/invalid_stats.out | 2 +- .../suites/statistics/analyze_stats.groovy | 4 +- .../suites/statistics/test_analyze_all_null.groovy | 4 +- .../test_analyze_sample_almost_all_null.groovy | 138 +++++++++++++++++++++ 9 files changed, 173 insertions(+), 51 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/statistics/BaseAnalysisTask.java b/fe/fe-core/src/main/java/org/apache/doris/statistics/BaseAnalysisTask.java index 47a2cb98f05..2707d28ca6d 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/statistics/BaseAnalysisTask.java +++ b/fe/fe-core/src/main/java/org/apache/doris/statistics/BaseAnalysisTask.java @@ -631,9 +631,12 @@ public abstract class BaseAnalysisTask { if (MetricRepo.isInit) { MetricRepo.COUNTER_STATISTICS_INVALID_STATS.increase(1L); } - String message = String.format("ColStatsData is invalid, skip analyzing. %s", colStatsData.toSQL(true)); - LOG.warn(message); - throw new RuntimeException(message); + // Don't throw: keep writing the row into the statistics table so that the + // whole job still finishes. toColumnStatistic() will defensively convert + // this pattern to ColumnStatistic.UNKNOWN at read time, so the optimizer + // never sees the bogus numbers. See issue #64122. + LOG.warn("ColStatsData is invalid, will write to table but be treated as UNKNOWN at read time. {}", + colStatsData.toSQL(true)); } // Update index row count after analyze. if (this instanceof OlapAnalysisTask) { diff --git a/fe/fe-core/src/main/java/org/apache/doris/statistics/ColStatsData.java b/fe/fe-core/src/main/java/org/apache/doris/statistics/ColStatsData.java index 232ae506d71..f9cbacb86c5 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/statistics/ColStatsData.java +++ b/fe/fe-core/src/main/java/org/apache/doris/statistics/ColStatsData.java @@ -139,6 +139,10 @@ public class ColStatsData { public ColumnStatistic toColumnStatistic() { try { + if (!isValid()) { + return ColumnStatistic.UNKNOWN; + } + ColumnStatisticBuilder columnStatisticBuilder = new ColumnStatisticBuilder(count); columnStatisticBuilder.setNdv(ndv); columnStatisticBuilder.setNumNulls(nullCount); diff --git a/fe/fe-core/src/main/java/org/apache/doris/statistics/ColumnStatistic.java b/fe/fe-core/src/main/java/org/apache/doris/statistics/ColumnStatistic.java index 7cdde54dce8..8dfd8b86f1f 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/statistics/ColumnStatistic.java +++ b/fe/fe-core/src/main/java/org/apache/doris/statistics/ColumnStatistic.java @@ -141,6 +141,11 @@ public class ColumnStatistic { * this function is used by analyze job and cbo job. */ public static ColumnStatistic fromResultRow(ResultRow row) { + ColStatsData statsData = new ColStatsData(row); + if (!statsData.isValid()) { + return ColumnStatistic.UNKNOWN; + } + double count = Double.parseDouble(row.get(7)); ColumnStatisticBuilder columnStatisticBuilder = new ColumnStatisticBuilder(count); double ndv = Double.parseDouble(row.getWithDefault(8, "0")); diff --git a/fe/fe-core/src/main/java/org/apache/doris/statistics/StatisticsCache.java b/fe/fe-core/src/main/java/org/apache/doris/statistics/StatisticsCache.java index 08532976def..4e7ea91f945 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/statistics/StatisticsCache.java +++ b/fe/fe-core/src/main/java/org/apache/doris/statistics/StatisticsCache.java @@ -295,7 +295,7 @@ public class StatisticsCache { } /** - * Refresh stats cache, invalidate cache if the new data is unknown. + * Refresh stats cache, publish UNKNOWN if the new data is unknown. */ public void syncColStats(ColStatsData data) { StatsId statsId = data.statsId; @@ -303,7 +303,14 @@ public class StatisticsCache { statsId.idxId, statsId.colId); ColumnStatistic columnStatistic = data.toColumnStatistic(); if (columnStatistic == ColumnStatistic.UNKNOWN) { - invalidateColumnStatsCache(k.catalogId, k.dbId, k.tableId, k.idxId, k.colName); + // Publish a blocking UNKNOWN instead of invalidating. Invalidation leaves the + // entry absent, so a concurrent get() can trigger the async loader which reads + // the previous (stale) row from the statistics table before the buffered + // insert commits, and that stale value would then survive until the next + // refresh. A put closes this window: readers hit UNKNOWN directly and no + // reload races with the flush. + updateColStatsCache(k.catalogId, k.dbId, k.tableId, k.idxId, k.colName, + ColumnStatistic.UNKNOWN); } else { putCache(k, columnStatistic); } diff --git a/fe/fe-core/src/test/java/org/apache/doris/statistics/BaseAnalysisTaskTest.java b/fe/fe-core/src/test/java/org/apache/doris/statistics/BaseAnalysisTaskTest.java index 343407dcaa6..9e2ea6d3510 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/statistics/BaseAnalysisTaskTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/statistics/BaseAnalysisTaskTest.java @@ -20,11 +20,8 @@ package org.apache.doris.statistics; import org.apache.doris.analysis.TableSample; import org.apache.doris.catalog.Column; import org.apache.doris.catalog.PrimitiveType; -import org.apache.doris.qe.StmtExecutor; import com.google.common.collect.Lists; -import mockit.Mock; -import mockit.MockUp; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; @@ -84,25 +81,9 @@ public class BaseAnalysisTaskTest { values.add("500"); values.add(null); ResultRow row = new ResultRow(values); - List<ResultRow> result = Lists.newArrayList(); - result.add(row); - - new MockUp<StmtExecutor>() { - @Mock - public List<ResultRow> executeInternalQuery() { - return result; - } - }; - BaseAnalysisTask task = new OlapAnalysisTask(); - try { - task.runQuery("test"); - } catch (Exception e) { - Assertions.assertEquals(e.getMessage(), - "ColStatsData is invalid, skip analyzing. " - + "('id',10000,20000,30000,0,'col',null,100,1100,300,'min','max',400,'500',NULL)"); - return; - } - Assertions.fail(); + ColStatsData data = new ColStatsData(row); + Assertions.assertFalse(data.isValid()); + Assertions.assertEquals(ColumnStatistic.UNKNOWN, data.toColumnStatistic()); } @Test @@ -124,24 +105,8 @@ public class BaseAnalysisTaskTest { values.add("500"); values.add(null); ResultRow row = new ResultRow(values); - List<ResultRow> result = Lists.newArrayList(); - result.add(row); - - new MockUp<StmtExecutor>() { - @Mock - public List<ResultRow> executeInternalQuery() { - return result; - } - }; - BaseAnalysisTask task = new OlapAnalysisTask(); - try { - task.runQuery("test"); - } catch (Exception e) { - Assertions.assertEquals(e.getMessage(), - "ColStatsData is invalid, skip analyzing. " - + "('id',10000,20000,30000,0,'col',null,500,0,300,'min','max',400,'500',NULL)"); - return; - } - Assertions.fail(); + ColStatsData data = new ColStatsData(row); + Assertions.assertFalse(data.isValid()); + Assertions.assertEquals(ColumnStatistic.UNKNOWN, data.toColumnStatistic()); } } diff --git a/regression-test/data/nereids_p0/stats/invalid_stats/invalid_stats.out b/regression-test/data/nereids_p0/stats/invalid_stats/invalid_stats.out index 9b1b2e2aa97..84eab1a26a6 100644 --- a/regression-test/data/nereids_p0/stats/invalid_stats/invalid_stats.out +++ b/regression-test/data/nereids_p0/stats/invalid_stats/invalid_stats.out @@ -26,6 +26,6 @@ PhysicalResultSink -- !ndv_row_invalid -- PhysicalResultSink --hashJoin[INNER_JOIN broadcast] hashCondition=((region.r_regionkey = nation.n_regionkey)) otherCondition=() -----PhysicalOlapScan[region] ----PhysicalOlapScan[nation] +----PhysicalOlapScan[region] diff --git a/regression-test/suites/statistics/analyze_stats.groovy b/regression-test/suites/statistics/analyze_stats.groovy index 2f715c92f55..7814f2924d1 100644 --- a/regression-test/suites/statistics/analyze_stats.groovy +++ b/regression-test/suites/statistics/analyze_stats.groovy @@ -2788,9 +2788,9 @@ PARTITION `p599` VALUES IN (599) sql """alter table alter_test modify column id set stats ('row_count'='100', 'ndv'='0', 'num_nulls'='0.0', 'data_size'='2.69975443E8', 'min_value'='1', 'max_value'='2');""" alter_result = sql """show column stats alter_test(id)""" logger.info("show column alter_test(id) stats: " + alter_result) - assertEquals(1, alter_result.size()) + assertEquals(0, alter_result.size()) alter_result = sql """show column cached stats alter_test(id)""" - assertEquals(1, alter_result.size()) + assertEquals(0, alter_result.size()) sql """alter table alter_test modify column id set stats ('row_count'='100', 'ndv'='0', 'num_nulls'='100', 'data_size'='2.69975443E8', 'min_value'='1', 'max_value'='2');""" alter_result = sql """show column stats alter_test(id)""" logger.info("show column alter_test(id) stats: " + alter_result) diff --git a/regression-test/suites/statistics/test_analyze_all_null.groovy b/regression-test/suites/statistics/test_analyze_all_null.groovy index 44d2f3a6c1f..0ef536afcd6 100644 --- a/regression-test/suites/statistics/test_analyze_all_null.groovy +++ b/regression-test/suites/statistics/test_analyze_all_null.groovy @@ -95,12 +95,12 @@ suite("test_analyze_all_null") { sql """alter table invalidTest modify column col2 set stats ('row_count'='100', 'ndv'='0', 'num_nulls'='0.0', 'data_size'='3.2E8', 'min_value'='min', 'max_value'='max');""" sql """alter table invalidTest modify column col3 set stats ('row_count'='100', 'ndv'='0', 'num_nulls'='100', 'data_size'='3.2E8', 'min_value'='min', 'max_value'='max');""" result = sql """show column cached stats invalidTest""" - assertEquals(3, result.size()) + assertEquals(2, result.size()) explain { sql("memo plan select * from invalidTest") contains "col1#0 -> ndv=100.0000" - contains "col2#1 -> ndv=0.0000" + contains "col2#1 -> unknown(" contains "col3#2 -> ndv=0.0000" } diff --git a/regression-test/suites/statistics/test_analyze_sample_almost_all_null.groovy b/regression-test/suites/statistics/test_analyze_sample_almost_all_null.groovy new file mode 100644 index 00000000000..13f7f6a3479 --- /dev/null +++ b/regression-test/suites/statistics/test_analyze_sample_almost_all_null.groovy @@ -0,0 +1,138 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +// Regression for issue #64122: ColStatsData.isValid() falsely rejects +// sampled column statistics when a column is (almost) all NULL. +// +// On a Unique-Key MoW table where column v is almost entirely NULL but +// has one surviving non-null value, sample analyze produces +// ndv=0 (estimated), min=max='x' (full-scan), nullCount != count +// which trips the second isValid() guard. Before the fix, runQuery() +// threw and aborted the whole analyze job; after the fix the row is +// written and toColumnStatistic() falls back to UNKNOWN at read time. +suite("test_analyze_sample_almost_all_null", "nonConcurrent") { + + def wait_row_count_at_least = { db, table, threshold -> + // For Unique MoW the post-DELETE row count is non-trivial to predict + // exactly, so we just gate on "row count is reported and large enough", + // which is what we need to trigger the isValid() guard. count=0 would + // short-circuit isValid() and the issue would not reproduce. + def result = sql """show frontends;""" + def host + def port + for (int i = 0; i < result.size(); i++) { + if (result[i][8] == "true") { + host = result[i][1] + port = result[i][4] + } + } + def tokens = context.config.jdbcUrl.split('/') + def url = tokens[0] + "//" + host + ":" + port + connect(context.config.jdbcUser, context.config.jdbcPassword, url) { + sql """use ${db}""" + for (int i = 0; i < 120; i++) { + Thread.sleep(5000) + result = sql """SHOW DATA FROM ${table};""" + logger.info("SHOW DATA FROM ${table}: " + result) + // Sum the row-count column across all rows returned by SHOW DATA. + // Layout: rows are per-partition + a Total row at the end. The + // row-count column index is 4 (same assumption as the existing + // test_analyze_all_null suite). + def total = 0L + for (int r = 0; r < result.size(); r++) { + def v = result[r][4] + if (v == null) { + continue + } + try { + total += Long.parseLong(v.toString()) + } catch (NumberFormatException ignored) { + // "Total" row may already be a formatted string; skip. + } + } + if (total >= threshold) { + return + } + } + throw new Exception("Row count report timeout for ${db}.${table}, " + + "threshold=" + threshold + ", last result=" + result) + } + } + + sql """drop database if exists regression_test_analyze_sample_almost_all_null""" + sql """create database regression_test_analyze_sample_almost_all_null""" + sql """use regression_test_analyze_sample_almost_all_null""" + + // The suite mutates the cluster-global enable_auto_analyze, so it runs in the + // nonConcurrent group and restores the variable via setGlobalVarTemporary to + // avoid leaking the disabled state into other suites on failure. + setGlobalVarTemporary([enable_auto_analyze: false], { + sql """CREATE TABLE tbl_del_big ( + k INT NOT NULL, + v VARCHAR(64) NULL + ) + UNIQUE KEY(k) + DISTRIBUTED BY HASH(k) BUCKETS 64 + PROPERTIES ( + "replication_num" = "1", + "enable_unique_key_merge_on_write" = "true" + ) + """ + + // 2M rows with v=NULL; then 1 row with v='x' overwriting (k=1, NULL). + // DELETE half of the NULL rows so the surviving data set is ~1M rows, + // with exactly one non-null v value. + sql """INSERT INTO tbl_del_big SELECT number, NULL FROM numbers("number"="2000000")""" + sql """INSERT INTO tbl_del_big SELECT number * 64 + 1, 'x' FROM numbers("number"="1")""" + sql """DELETE FROM tbl_del_big WHERE k % 2 = 0 AND v IS NULL""" + + wait_row_count_at_least("regression_test_analyze_sample_almost_all_null", + "tbl_del_big", 500000L) + + sql """ANALYZE TABLE tbl_del_big WITH SAMPLE PERCENT 1 WITH SYNC""" + + def result = sql """show column stats tbl_del_big""" + + // k (NOT NULL) always produces valid sampled stats. Whether v also survives + // isValid() depends on whether the single 'x' row lands in one of the randomly + // chosen sample tablets (sampled -> ndv ~ 1, valid; not sampled -> ndv = 0 with + // full-scan min/max = 'x', invalid). So only assert on k and a loose row count. + assertTrue(result.size() >= 1) + assertTrue(result.any { it[0] == "k" }) + + // Deterministically construct the issue #64122 invalid pattern. SET STATS writes + // the row into the statistics table directly (no isValid check on that path) and + // syncColStats publishes UNKNOWN into the cache. Any later read also goes through + // ColumnStatistic.fromResultRow, whose isValid() guard returns UNKNOWN for + // ndv=0 + min/max!=null + nullCount!=count, so the optimizer must see unknown. + sql """ALTER TABLE tbl_del_big MODIFY COLUMN v SET STATS ( + 'row_count'='1000000', 'ndv'='0', 'num_nulls'='999999', + 'data_size'='8000000', 'min_value'='x', 'max_value'='x')""" + + explain { + sql("select * from tbl_del_big") + contains("planned with unknown column statistics") + } + + explain { + sql("memo plan select * from tbl_del_big") + contains("v#1 -> unknown(") + } + }) + + sql """drop database if exists regression_test_analyze_sample_almost_all_null""" +} --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
