diff --git a/datavines-common/src/main/java/io/datavines/common/ConfigConstants.java b/datavines-common/src/main/java/io/datavines/common/ConfigConstants.java index 6bf2817a0..367eaf349 100644 --- a/datavines-common/src/main/java/io/datavines/common/ConfigConstants.java +++ b/datavines-common/src/main/java/io/datavines/common/ConfigConstants.java @@ -38,6 +38,7 @@ public class ConfigConstants { public static final String ACTUAL_NAME = "actual_name"; public static final String ACTUAL_EXECUTE_SQL = "actual_execute_sql"; public static final String ACTUAL_AGGREGATE_SQL = "actual_aggregate_sql"; + public static final String INVALIDATE_ITEMS_SQL = "invalidate_items_sql"; public static final String ACTUAL_CUSTOM_SQL = "actual_custom_sql"; public static final String EXPECTED_NAME = "expected_name"; public static final String EXPECTED_TYPE = "expected_type"; diff --git a/datavines-engine/datavines-engine-plugins/datavines-engine-flink/datavines-engine-flink-config/src/main/java/io/datavines/engine/flink/config/BaseFlinkConfigurationBuilder.java b/datavines-engine/datavines-engine-plugins/datavines-engine-flink/datavines-engine-flink-config/src/main/java/io/datavines/engine/flink/config/BaseFlinkConfigurationBuilder.java index 5e747510b..bdaa49c57 100644 --- a/datavines-engine/datavines-engine-plugins/datavines-engine-flink/datavines-engine-flink-config/src/main/java/io/datavines/engine/flink/config/BaseFlinkConfigurationBuilder.java +++ b/datavines-engine/datavines-engine-plugins/datavines-engine-flink/datavines-engine-flink-config/src/main/java/io/datavines/engine/flink/config/BaseFlinkConfigurationBuilder.java @@ -76,7 +76,7 @@ protected List getSourceConfigs() throws DataVinesException { .getNewPlugin(metricType); if (sqlMetric.isCustomSql()) { - List tables = SqlUtils.extractTablesFromSelect(metricInputParameter.get(ACTUAL_AGGREGATE_SQL)); + List tables = SqlUtils.extractTablesFromSelect(sqlMetric.getTableDiscoverySql(metricInputParameter)); if (CollectionUtils.isEmpty(tables)) { throw new DataVinesException("custom sql must have table"); } @@ -138,13 +138,13 @@ protected List getSourceConfigs() throws DataVinesException { sourceConnectorSet.add(connectorUUID); } - if (StringUtils.isNotEmpty(metricInputParameter.get(ACTUAL_AGGREGATE_SQL))) { - String sql = metricInputParameter.get(ACTUAL_AGGREGATE_SQL); + if (StringUtils.isNotEmpty(sqlMetric.getTableDiscoverySql(metricInputParameter))) { + String sql = sqlMetric.getTableDiscoverySql(metricInputParameter); for (Map.Entry entry : table2OutputTable.entrySet()) { sql = sql.replaceAll(entry.getKey(), entry.getValue()); } - metricInputParameter.put(ACTUAL_AGGREGATE_SQL, sql); + sqlMetric.setTableDiscoverySql(metricInputParameter, sql); } } else { ConnectorParameter connectorParameter = jobExecutionParameter.getConnectorParameter(); diff --git a/datavines-engine/datavines-engine-plugins/datavines-engine-local/datavines-engine-local-config/src/main/java/io/datavines/engine/local/config/BaseLocalConfigurationBuilder.java b/datavines-engine/datavines-engine-plugins/datavines-engine-local/datavines-engine-local-config/src/main/java/io/datavines/engine/local/config/BaseLocalConfigurationBuilder.java index adeeb2f5f..930b01921 100644 --- a/datavines-engine/datavines-engine-plugins/datavines-engine-local/datavines-engine-local-config/src/main/java/io/datavines/engine/local/config/BaseLocalConfigurationBuilder.java +++ b/datavines-engine/datavines-engine-plugins/datavines-engine-local/datavines-engine-local-config/src/main/java/io/datavines/engine/local/config/BaseLocalConfigurationBuilder.java @@ -73,7 +73,7 @@ protected List getSourceConfigs() throws DataVinesException { .getNewPlugin(connectorParameter.getType()); - List tables = SqlUtils.extractTablesFromSelect(metricInputParameter.get(ACTUAL_AGGREGATE_SQL)); + List tables = SqlUtils.extractTablesFromSelect(sqlMetric.getTableDiscoverySql(metricInputParameter)); if (CollectionUtils.isEmpty(tables)) { throw new DataVinesException("custom sql must have table"); } diff --git a/datavines-engine/datavines-engine-plugins/datavines-engine-spark/datavines-engine-spark-config/src/main/java/io/datavines/engine/spark/config/BaseSparkConfigurationBuilder.java b/datavines-engine/datavines-engine-plugins/datavines-engine-spark/datavines-engine-spark-config/src/main/java/io/datavines/engine/spark/config/BaseSparkConfigurationBuilder.java index 80df939ad..eafb45a34 100644 --- a/datavines-engine/datavines-engine-plugins/datavines-engine-spark/datavines-engine-spark-config/src/main/java/io/datavines/engine/spark/config/BaseSparkConfigurationBuilder.java +++ b/datavines-engine/datavines-engine-plugins/datavines-engine-spark/datavines-engine-spark-config/src/main/java/io/datavines/engine/spark/config/BaseSparkConfigurationBuilder.java @@ -102,7 +102,7 @@ protected List getSourceConfigs() throws DataVinesException { .getNewPlugin(metricType); if (sqlMetric.isCustomSql()) { - List tables = SqlUtils.extractTablesFromSelect(metricInputParameter.get(ACTUAL_AGGREGATE_SQL)); + List tables = SqlUtils.extractTablesFromSelect(sqlMetric.getTableDiscoverySql(metricInputParameter)); if (CollectionUtils.isEmpty(tables)) { throw new DataVinesException("custom sql must have table"); } @@ -164,13 +164,13 @@ protected List getSourceConfigs() throws DataVinesException { sourceConnectorSet.add(connectorUUID); } - if (StringUtils.isNotEmpty(metricInputParameter.get(ACTUAL_AGGREGATE_SQL))) { - String sql = metricInputParameter.get(ACTUAL_AGGREGATE_SQL); + if (StringUtils.isNotEmpty(sqlMetric.getTableDiscoverySql(metricInputParameter))) { + String sql = sqlMetric.getTableDiscoverySql(metricInputParameter); for (Map.Entry entry : table2OutputTable.entrySet()) { sql = sql.replaceAll(entry.getKey(), entry.getValue()); } - metricInputParameter.put(ACTUAL_AGGREGATE_SQL, sql); + sqlMetric.setTableDiscoverySql(metricInputParameter, sql); } } else { ConnectorParameter connectorParameter = jobExecutionParameter.getConnectorParameter(); diff --git a/datavines-metric/datavines-metric-api/src/main/java/io/datavines/metric/api/SqlMetric.java b/datavines-metric/datavines-metric-api/src/main/java/io/datavines/metric/api/SqlMetric.java index d7376091e..cf5706f8c 100644 --- a/datavines-metric/datavines-metric-api/src/main/java/io/datavines/metric/api/SqlMetric.java +++ b/datavines-metric/datavines-metric-api/src/main/java/io/datavines/metric/api/SqlMetric.java @@ -25,6 +25,8 @@ import io.datavines.common.entity.ExecuteSql; import io.datavines.common.enums.DataVinesDataType; +import static io.datavines.common.ConfigConstants.ACTUAL_AGGREGATE_SQL; + public interface SqlMetric { String getName(); @@ -115,4 +117,12 @@ default MetricDirectionType getDirectionType() { default boolean isCustomSql() { return false; } + + default String getTableDiscoverySql(Map inputParameter) { + return inputParameter.get(ACTUAL_AGGREGATE_SQL); + } + + default void setTableDiscoverySql(Map inputParameter, String sql) { + inputParameter.put(ACTUAL_AGGREGATE_SQL, sql); + } } diff --git a/datavines-metric/datavines-metric-plugins/datavines-metric-all/pom.xml b/datavines-metric/datavines-metric-plugins/datavines-metric-all/pom.xml index a5a9b59b2..c8c9786c1 100644 --- a/datavines-metric/datavines-metric-plugins/datavines-metric-all/pom.xml +++ b/datavines-metric/datavines-metric-plugins/datavines-metric-all/pom.xml @@ -180,6 +180,12 @@ ${project.version} + + io.datavines + datavines-metric-custom-count-sql + ${project.version} + + io.datavines datavines-metric-multi-table-accuracy diff --git a/datavines-metric/datavines-metric-plugins/datavines-metric-custom-count-sql/pom.xml b/datavines-metric/datavines-metric-plugins/datavines-metric-custom-count-sql/pom.xml new file mode 100644 index 000000000..e81b424e0 --- /dev/null +++ b/datavines-metric/datavines-metric-plugins/datavines-metric-custom-count-sql/pom.xml @@ -0,0 +1,34 @@ + + + + + + datavines-metric-plugins + io.datavines + 1.0.0-SNAPSHOT + + 4.0.0 + + datavines-metric-custom-count-sql + + + diff --git a/datavines-metric/datavines-metric-plugins/datavines-metric-custom-count-sql/src/main/java/io/datavines/metric/plugin/CustomCountSql.java b/datavines-metric/datavines-metric-plugins/datavines-metric-custom-count-sql/src/main/java/io/datavines/metric/plugin/CustomCountSql.java new file mode 100644 index 000000000..960081c90 --- /dev/null +++ b/datavines-metric/datavines-metric-plugins/datavines-metric-custom-count-sql/src/main/java/io/datavines/metric/plugin/CustomCountSql.java @@ -0,0 +1,173 @@ +/* + * 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. + */ +package io.datavines.metric.plugin; + +import java.util.*; + +import io.datavines.common.config.CheckResult; +import io.datavines.common.config.ConfigChecker; +import io.datavines.common.entity.ExecuteSql; +import io.datavines.common.enums.DataVinesDataType; +import io.datavines.common.utils.StringUtils; +import io.datavines.metric.api.ConfigItem; +import io.datavines.metric.api.MetricDimension; +import io.datavines.metric.api.MetricType; +import io.datavines.metric.api.SqlMetric; + +import static io.datavines.common.CommonConstants.TABLE; +import static io.datavines.common.ConfigConstants.*; + +public class CustomCountSql implements SqlMetric { + + private final Set requiredOptions = new HashSet<>(); + + private final HashMap configMap = new HashMap<>(); + + private String invalidateItemsSql = null; + + public CustomCountSql() { + configMap.put(INVALIDATE_ITEMS_SQL, new ConfigItem(INVALIDATE_ITEMS_SQL, "错误数据SQL", INVALIDATE_ITEMS_SQL)); + + requiredOptions.add(INVALIDATE_ITEMS_SQL); + } + + @Override + public String getName() { + return "custom_count_sql"; + } + + @Override + public String getZhName() { + return "自定义统计SQL"; + } + + @Override + public MetricDimension getDimension() { + return MetricDimension.ACCURACY; + } + + @Override + public MetricType getType() { + return MetricType.SINGLE_TABLE; + } + + @Override + public boolean isInvalidateItemsCanOutput() { + return true; + } + + @Override + public CheckResult validateConfig(Map config) { + CheckResult basicCheck = ConfigChecker.checkConfig(config, requiredOptions); + if (!basicCheck.isSuccess()) { + return basicCheck; + } + + Object sqlValue = config.get(INVALIDATE_ITEMS_SQL); + if (sqlValue == null || StringUtils.isEmpty(String.valueOf(sqlValue).trim())) { + return new CheckResult(false, INVALIDATE_ITEMS_SQL + " cannot be empty"); + } + + return new CheckResult(true, ""); + } + + @Override + public void prepare(Map config) { + if (config.containsKey(INVALIDATE_ITEMS_SQL) && StringUtils.isNotEmpty(config.get(INVALIDATE_ITEMS_SQL))) { + this.invalidateItemsSql = config.get(INVALIDATE_ITEMS_SQL); + } + } + + @Override + public Map getConfigMap() { + return configMap; + } + + @Override + public ExecuteSql getInvalidateItems(Map inputParameter) { + if (StringUtils.isEmpty(invalidateItemsSql)) { + throw new IllegalStateException("invalidate_items_sql is not configured or empty"); + } + + String uniqueKey = inputParameter.get(METRIC_UNIQUE_KEY); + if (StringUtils.isEmpty(uniqueKey)) { + throw new IllegalStateException("metric_unique_key is missing in input parameters"); + } + + ExecuteSql executeSql = new ExecuteSql(); + executeSql.setResultTable("invalidate_items_" + uniqueKey); + executeSql.setSql(invalidateItemsSql); + executeSql.setErrorOutput(true); + return executeSql; + } + + @Override + public ExecuteSql getActualValue(Map inputParameter) { + if (StringUtils.isEmpty(invalidateItemsSql)) { + throw new IllegalStateException("invalidate_items_sql is not configured or empty"); + } + + String uniqueKey = inputParameter.get(METRIC_UNIQUE_KEY); + if (StringUtils.isEmpty(uniqueKey)) { + throw new IllegalStateException("metric_unique_key is missing in input parameters"); + } + + inputParameter.put(ACTUAL_TABLE, inputParameter.get(TABLE)); + + String actualAggregateSql = "select count(1) as actual_value_" + uniqueKey + " from ${invalidate_items_table}"; + + return new ExecuteSql(actualAggregateSql, "invalidate_count_" + uniqueKey); + } + + @Override + public ExecuteSql getDirectActualValue(Map inputParameter) { + if (StringUtils.isEmpty(invalidateItemsSql)) { + throw new IllegalStateException("invalidate_items_sql is not configured or empty"); + } + + String uniqueKey = inputParameter.get(METRIC_UNIQUE_KEY); + if (StringUtils.isEmpty(uniqueKey)) { + throw new IllegalStateException("metric_unique_key is missing in input parameters"); + } + + String actualAggregateSql = "select count(1) as actual_value_" + uniqueKey + + " from ( " + invalidateItemsSql + " ) t"; + + return new ExecuteSql(actualAggregateSql, "invalidate_count_" + uniqueKey); + } + + @Override + public List suitableType() { + return Collections.emptyList(); + } + + @Override + public boolean isCustomSql() { + return true; + } + + @Override + public String getTableDiscoverySql(Map inputParameter) { + return inputParameter.get(INVALIDATE_ITEMS_SQL); + } + + @Override + public void setTableDiscoverySql(Map inputParameter, String sql) { + inputParameter.put(INVALIDATE_ITEMS_SQL, sql); + this.invalidateItemsSql = sql; + } +} diff --git a/datavines-metric/datavines-metric-plugins/datavines-metric-custom-count-sql/src/main/resources/META-INF/services/io.datavines.metric.api.SqlMetric b/datavines-metric/datavines-metric-plugins/datavines-metric-custom-count-sql/src/main/resources/META-INF/services/io.datavines.metric.api.SqlMetric new file mode 100644 index 000000000..defbc0b06 --- /dev/null +++ b/datavines-metric/datavines-metric-plugins/datavines-metric-custom-count-sql/src/main/resources/META-INF/services/io.datavines.metric.api.SqlMetric @@ -0,0 +1 @@ +io.datavines.metric.plugin.CustomCountSql diff --git a/datavines-metric/datavines-metric-plugins/pom.xml b/datavines-metric/datavines-metric-plugins/pom.xml index fa98824fb..34aee0458 100644 --- a/datavines-metric/datavines-metric-plugins/pom.xml +++ b/datavines-metric/datavines-metric-plugins/pom.xml @@ -33,6 +33,7 @@ datavines-metric-all datavines-metric-base datavines-metric-custom-aggregate-sql + datavines-metric-custom-count-sql datavines-metric-multi-table-accuracy datavines-metric-table-freshness datavines-metric-table-row-count diff --git a/datavines-server/src/main/java/io/datavines/server/repository/service/impl/JobServiceImpl.java b/datavines-server/src/main/java/io/datavines/server/repository/service/impl/JobServiceImpl.java index f57a28c0c..b998d09f6 100644 --- a/datavines-server/src/main/java/io/datavines/server/repository/service/impl/JobServiceImpl.java +++ b/datavines-server/src/main/java/io/datavines/server/repository/service/impl/JobServiceImpl.java @@ -394,7 +394,11 @@ private String getFQN(BaseJobParameter jobParameter) { } if (StringUtils.isEmpty(table)) { - List tables = SqlUtils.extractTablesFromSelect((String) jobParameter.getMetricParameter().get(ACTUAL_AGGREGATE_SQL)); + SqlMetric sqlMetric = PluginDiscovery.getMultiKeyPluginDiscovery(SqlMetric.class, SqlMetric::getPluginNames) + .getOrCreatePlugin(jobParameter.getMetricType()); + Map metricParameter = new HashMap<>(); + jobParameter.getMetricParameter().forEach((key, value) -> metricParameter.put(key, String.valueOf(value))); + List tables = SqlUtils.extractTablesFromSelect(sqlMetric.getTableDiscoverySql(metricParameter)); if (CollectionUtils.isEmpty(tables)) { throw new DataVinesException("custom sql must have table"); } diff --git a/datavines-ui/Editor/components/MetricModal/MetricSelect/index.tsx b/datavines-ui/Editor/components/MetricModal/MetricSelect/index.tsx index 1330e2c4a..b7fc5573e 100644 --- a/datavines-ui/Editor/components/MetricModal/MetricSelect/index.tsx +++ b/datavines-ui/Editor/components/MetricModal/MetricSelect/index.tsx @@ -153,16 +153,43 @@ const Index = ({ ); }; - const dynamicRender = (item: dynamicConfigItem) => ( - - - - ); + // SQL field keys whitelist for TextArea rendering + const SQL_FIELD_KEYS = ['actual_aggregate_sql', 'invalidate_items_sql', 'actual_execute_sql', 'expected_execute_sql', 'actual_custom_sql']; + + const dynamicRender = (item: dynamicConfigItem) => { + // Use whitelist first, then fallback to pattern matching for unknown SQL fields + const isSqlField = SQL_FIELD_KEYS.includes(item.key) || item.key.toLowerCase().endsWith('_sql'); + if (isSqlField) { + const placeholderId = `${item.key}_placeholder` as any; + const placeholderText = intl.formatMessage({ id: placeholderId, defaultMessage: '' }); + const defaultPlaceholder = `${intl.formatMessage({ id: 'dv_metric_input' })} SQL`; + return ( + + + + ); + } + return ( + + + + ); + }; return ( <div> diff --git a/datavines-ui/Editor/locale/en_US.ts b/datavines-ui/Editor/locale/en_US.ts index 5f67e08c5..a25942e3b 100644 --- a/datavines-ui/Editor/locale/en_US.ts +++ b/datavines-ui/Editor/locale/en_US.ts @@ -78,6 +78,7 @@ export default { editor_dv_search_table: 'please enter table', editor_dv_search_column: 'please enter column', editor_dv_metric_name: 'name', + invalidate_items_sql_placeholder: 'e.g. SELECT * FROM users WHERE age < 0', dashboard_execution: 'Execution Dashboard', dashboard_quality_report: 'Quality Report Dashboard', diff --git a/datavines-ui/Editor/locale/zh_CN.ts b/datavines-ui/Editor/locale/zh_CN.ts index 91954187d..e900ac7c8 100644 --- a/datavines-ui/Editor/locale/zh_CN.ts +++ b/datavines-ui/Editor/locale/zh_CN.ts @@ -78,6 +78,7 @@ export default { editor_dv_search_table: '请输入表名', editor_dv_search_column: '请输入列名', editor_dv_metric_name: '名称', + invalidate_items_sql_placeholder: '如: SELECT * FROM users WHERE age < 0', dashboard_execution: '运行概况', dashboard_quality_report: '质量报告',