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 89fa7feecc2434f3e64e9cf493eaf0e311090fee Author: daidai <[email protected]> AuthorDate: Sat Oct 10 09:20:27 2026 +0800 branch-4.1: [fix](test)(dynamic-partition) fix some unstable test cases and dynamic-partition logic (#63551) (#68820) pick #63551 --- .../apache/doris/alter/SchemaChangeHandler.java | 33 +++++++++++++++++++--- .../org/apache/doris/load/loadv2/LoadManager.java | 20 +++++++++++++ .../org/apache/doris/qe/runtime/LoadProcessor.java | 6 ++++ regression-test/pipeline/p0/conf/fe.conf | 1 + .../jdbc/test_doris_jdbc_catalog.groovy | 4 +++ .../suites/manager/test_manager_interface_1.groovy | 4 +-- 6 files changed, 62 insertions(+), 6 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/alter/SchemaChangeHandler.java b/fe/fe-core/src/main/java/org/apache/doris/alter/SchemaChangeHandler.java index 1431aeae35b..c5e8f3ea976 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/alter/SchemaChangeHandler.java +++ b/fe/fe-core/src/main/java/org/apache/doris/alter/SchemaChangeHandler.java @@ -2508,10 +2508,35 @@ public class SchemaChangeHandler extends AlterHandler { skip = Boolean.parseBoolean(skipWriteIndexOnLoad) ? 1 : 0; } - for (Partition partition : partitions) { - updatePartitionProperties(db, olapTable.getName(), partition.getName(), storagePolicyId, isInMemory, - null, compactionPolicy, timeSeriesCompactionConfig, skip, - disableAutoCompaction, verticalCompactionNumColumnsPerGroup); + // Only iterate partitions when there are properties that actually need to be + // dispatched to each partition's tablets. Pure catalog-level metadata properties + // such as partition.retention_count do not require per-partition updates, and + // iterating over a stale partition snapshot can race with concurrent partition + // drops (e.g., by DynamicPartitionScheduler when retention_count or dynamic_partition + // is enabled) and fail with "Partition does not exist". + boolean needPerPartitionUpdate = isInMemory >= 0 || storagePolicyId >= 0 + || compactionPolicy != null || !timeSeriesCompactionConfig.isEmpty() + || skip >= 0 || disableAutoCompaction >= 0 + || verticalCompactionNumColumnsPerGroup >= 0; + if (needPerPartitionUpdate) { + for (Partition partition : partitions) { + try { + updatePartitionProperties(db, olapTable.getName(), partition.getName(), + storagePolicyId, isInMemory, null, compactionPolicy, timeSeriesCompactionConfig, + skip, disableAutoCompaction, + verticalCompactionNumColumnsPerGroup); + } catch (DdlException e) { + // The partition may have been dropped concurrently (e.g., by + // DynamicPartitionScheduler). It is safe to skip the meta dispatch + // for a partition that no longer exists. + if (olapTable.getPartition(partition.getName()) == null) { + LOG.info("partition {} of table {} was dropped concurrently, " + + "skip updating its properties", partition.getName(), olapTable.getName()); + continue; + } + throw e; + } + } } olapTable.writeLockOrDdlException(); diff --git a/fe/fe-core/src/main/java/org/apache/doris/load/loadv2/LoadManager.java b/fe/fe-core/src/main/java/org/apache/doris/load/loadv2/LoadManager.java index 019e06893a0..9fc0ca13fdc 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/load/loadv2/LoadManager.java +++ b/fe/fe-core/src/main/java/org/apache/doris/load/loadv2/LoadManager.java @@ -216,11 +216,25 @@ public class LoadManager implements Writable { LoadJob loadJob; if (idToLoadJob.containsKey(jobId)) { loadJob = idToLoadJob.get(jobId); + if (LOG.isDebugEnabled()) { + LOG.debug("recordFinishedLoadJob: reuse existing load job, jobId={}, label={}, dbId={}, jobType={}", + jobId, label, db.getId(), jobType); + } if (loadJob instanceof InsertLoadJob) { ((InsertLoadJob) loadJob).setJobProperties(transactionId, tableId, createTimestamp, failMsg, trackingUrl, firstErrorMsg, userInfo); } } else { + // The jobId received here does not exist in idToLoadJob. This means the InsertLoadJob + // that was registered during executor construction (and that accumulated BE-reported + // load statistics via updateJobProgress) is NOT the one we are about to snapshot here. + // A brand-new InsertLoadJob will be created below with an empty LoadStatistic, so + // SHOW LOAD's JobDetails (ScannedRows / LoadBytes / All backends) will be all zero. + // Logging this at WARN so CI failures of the form + // "test_insert_statistic: expected:<N> but was:<0>" can be diagnosed directly. + LOG.warn("recordFinishedLoadJob: jobId={} not found in idToLoadJob, creating a new " + + "{} load job for label={}, dbId={}. JobDetails statistics will be empty.", + jobId, jobType, label, db.getId()); switch (jobType) { case INSERT: loadJob = new InsertLoadJob(label, transactionId, db.getId(), tableId, createTimestamp, failMsg, @@ -813,6 +827,12 @@ public class LoadManager implements Writable { public void updateJobProgress(Long jobId, Long beId, TUniqueId loadId, TUniqueId fragmentId, long scannedRows, long scannedBytes, boolean isDone) { LoadJob job = idToLoadJob.get(jobId); + if (LOG.isDebugEnabled()) { + LOG.debug("updateJobProgress: jobId={}, beId={}, scannedRows={}, scannedBytes={}, isDone={}, " + + "found={}, jobIdMatched={}", + jobId, beId, scannedRows, scannedBytes, isDone, + idToLoadJob.containsKey(jobId), job == null ? -1L : job.getId()); + } if (job != null) { job.updateProgress(beId, loadId, fragmentId, scannedRows, scannedBytes, isDone); } diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/LoadProcessor.java b/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/LoadProcessor.java index c2b769196af..b805fea5312 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/LoadProcessor.java +++ b/fe/fe-core/src/main/java/org/apache/doris/qe/runtime/LoadProcessor.java @@ -165,6 +165,12 @@ public class LoadProcessor extends AbstractJobProcessor { @Override protected void doProcessReportExecStatus(TReportExecStatusParams params, SingleFragmentPipelineTask fragmentTask) { if (params.isSetLoadedRows() && jobId != -1) { + if (LOG.isDebugEnabled()) { + LOG.debug("doProcessReportExecStatus: forwarding load progress to LoadManager, " + + "jobId={}, beId={}, queryId={}, loadedRows={}, loadedBytes={}, isDone={}", + jobId, params.getBackendId(), DebugUtil.printId(params.getQueryId()), + params.getLoadedRows(), params.getLoadedBytes(), params.isDone()); + } if (params.isSetFragmentInstanceReports()) { for (TFragmentInstanceReport report : params.getFragmentInstanceReports()) { Env.getCurrentEnv().getLoadManager().updateJobProgress( diff --git a/regression-test/pipeline/p0/conf/fe.conf b/regression-test/pipeline/p0/conf/fe.conf index c79ec24eb5f..3d86c5c95e6 100644 --- a/regression-test/pipeline/p0/conf/fe.conf +++ b/regression-test/pipeline/p0/conf/fe.conf @@ -93,3 +93,4 @@ max_spilled_profile_num = 2000 check_table_lock_leaky=true max_bucket_num_per_partition=512 +max_remote_file_system_cache_num=1000 diff --git a/regression-test/suites/external_table_p0/jdbc/test_doris_jdbc_catalog.groovy b/regression-test/suites/external_table_p0/jdbc/test_doris_jdbc_catalog.groovy index 0057c58956a..743cc79d5de 100644 --- a/regression-test/suites/external_table_p0/jdbc/test_doris_jdbc_catalog.groovy +++ b/regression-test/suites/external_table_p0/jdbc/test_doris_jdbc_catalog.groovy @@ -262,6 +262,10 @@ suite("test_doris_jdbc_catalog", "p0,external,doris,external_docker,external_doc order_qt_keywork_table_name """ select * from `order` order by test_col; """ + // cleanup reserved-keyword table so downstream infra suites (e.g. check_hash_bucket_table) + // do not trip over an unquoted `order` table name + sql """ drop table if exists internal.regression_test_jdbc_catalog_p0.`order` """ + // //clean // qt_sql """select current_catalog()""" // sql "switch internal" diff --git a/regression-test/suites/manager/test_manager_interface_1.groovy b/regression-test/suites/manager/test_manager_interface_1.groovy index 1c9aa226200..1426306e9cc 100644 --- a/regression-test/suites/manager/test_manager_interface_1.groovy +++ b/regression-test/suites/manager/test_manager_interface_1.groovy @@ -594,7 +594,7 @@ suite('test_manager_interface_1',"p0") { assertTrue(x == 1); - sql """ admin set frontend config("query_metadata_name_ids_timeout"= "${val}")""" + sql """ admin set all frontends config("query_metadata_name_ids_timeout"= "${val}")""" result = sql """ admin show frontend config """ @@ -613,7 +613,7 @@ suite('test_manager_interface_1',"p0") { assertTrue(x == 1); val -= 2 - sql """ admin set frontend config("query_metadata_name_ids_timeout"= "${val}")""" + sql """ admin set all frontends config("query_metadata_name_ids_timeout"= "${val}")""" logger.info("result = ${result}" ) --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
