sadpandajoe commented on code in PR #43870:
URL: https://github.com/apache/superset/pull/43870#discussion_r4175105710
##########
superset/models/helpers.py:
##########
@@ -5602,6 +5831,11 @@ def get_sqla_query( # pylint:
disable=too-many-arguments,too-many-locals,too-ma
# col_obj is None and sqla_col is None - column not found!
# Silently skip - this can happen for removed columns or
invalid filters
pass
+ if partition_mapping is not None:
+ where_clause_and += self._build_partition_mirror_predicates(
Review Comment:
These outer-window partition predicates also enter the series-limit subquery
through `where_clause_and`, but relative time comparisons rank series using
different `inner_from_dttm`/`inner_to_dttm` bounds. A January 2025 comparison
can therefore rank January 2026 timestamps while requiring January 2025
partition keys, producing an empty ranking and dropping the comparison series;
could each query window build its own mirrors?
##########
superset/connectors/sqla/partition_mapping.py:
##########
@@ -0,0 +1,1314 @@
+# 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.
+"""
+Partition filter mapping.
+
+Datasets on Hadoop-family engines are commonly partitioned on a *technical*
+column -- an epoch integer, a lowercased region key -- that no analyst would
+filter on. Unless a query carries a predicate on that column the engine scans
+every partition.
+
+A dataset owner names one partition column ``p``, one business column that
+filters are mirrored from, and a value transform ``T`` (a SQL expression
+containing a ``:value`` placeholder). Superset then appends an equivalent
+predicate on ``p`` to every query, so chart authors change nothing and queries
+prune.
+
+The load-bearing assumption
+---------------------------
+Everything here reasons about ``T(col) op T(v)``, but what is emitted is
+``p op T(v)`` -- a predicate on a *physically different column*. The step from
+one to the other is::
+
+ p = T(mapped_col) for every row in the table
+
+Superset cannot verify that; it is a property of whatever ETL populates the
+partition column. If that job lags, backfills with different logic, or writes
+the partition key in a different timezone, mirrored predicates silently drop
+real rows. The mapping is only as trustworthy as the pipeline behind it.
+
+Operator safety
+---------------
+A mirrored predicate ``P2`` may only be ``AND``-ed onto a query when the
+original predicate ``P1`` *implies* it:
+
+=========================================== =============================
+Original Safe when
+=========================================== =============================
+``col = v``, ``col IN (...)`` always -- ``T`` is a function
+``col >=|>|<|<= v``, ``TEMPORAL_RANGE`` only if ``T`` is monotonic
+``col != v``, ``NOT IN``, ``LIKE``, ... never
+=========================================== =============================
+
+Negations are never safe because ``T`` need not be injective:
+``lower(:value)`` with ``country != 'US'`` mirrors to ``region_key != 'us'``,
+which wrongly excludes rows whose ``country`` is already lowercase ``'us'`` --
+rows the original filter *keeps*.
+
+Monotonicity is a property of the transform, not of the column's data type:
+``hour(:value)``, ``date_format(:value, 'dd')`` and ``dayofweek(:value)`` are
+all reasonable transforms on a ``TIMESTAMP`` column and none of them preserve
+ordering. It is therefore declared by the owner, not inferred.
+
+Time grains
+-----------
+A grained filter -- drill-to-detail, mostly -- compares the *truncated* column,
+``trunc(col) op v``, so the raw bounds it carries do not describe the rows it
+keeps. A row in the final partial bucket satisfies ``trunc(col) < until`` while
+``col < until`` excludes it.
+
+Which direction a grain rounds is not knowable from the duration alone, and the
+obvious guess is wrong: ``WEEK_ENDING_SATURDAY`` rounds *forward* on Hive and
+Presto, and Ocient's grains are ``ROUND``, i.e. to nearest. What every grain
+does satisfy is a bound on the displacement::
+
+ |trunc(ts) - ts| < width(grain)
+
+so widening *both* bounds by one bucket width is no narrower than the real
+predicate whichever way the grain rounds -- and it needs no monotonicity of
+``trunc`` itself, which is what rescues Hive's oddly-anchored ``P1W``. Width is
+a property of the grain, not of the engine, which is what makes this
+maintainable; see `grain_bucket_width`. A grain whose SQL an operator supplied
+has no known width, and does not mirror at all.
+"""
+
+from __future__ import annotations
+
+import hashlib
+import logging
+import re
+from dataclasses import dataclass
+from functools import lru_cache
+from typing import Any, cast, TYPE_CHECKING
+
+import numpy as np
+import pandas as pd
+import sqlalchemy as sa
+from dateutil.relativedelta import relativedelta
+from flask import current_app as app
+from flask_babel import lazy_gettext as _
+from sqlalchemy.engine.interfaces import Dialect
+from sqlalchemy.sql.elements import ColumnElement
+
+from superset.constants import LRU_CACHE_MAX_SIZE, TimeGrain
+from superset.exceptions import (
+ QueryClauseValidationException,
+ SupersetParseError,
+ SupersetSecurityException,
+)
+from superset.extensions import cache_manager, feature_flag_manager
+from superset.sql.parse import SQLStatement
+from superset.utils import json
+from superset.utils.core import FilterOperator
+
+if TYPE_CHECKING:
+ from superset.connectors.sqla.models import SqlaTable, TableColumn
+ from superset.models.core import Database
+
+logger = logging.getLogger(__name__)
+
+FEATURE_FLAG = "PARTITION_FILTER_MAPPING"
+
+#: Longest transform accepted anywhere. A transform is one SQL expression an
+#: owner types by hand, so this is generous; the point is that it is bounded.
+#: Every entry point enforces it -- the typed column field, the import schema
+#: and the preview request -- because `is_transform_active` parses the stored
+#: value on each Explore load, and an unbounded string would make that parse
+#: the expensive part of rendering a chart.
+MAX_TRANSFORM_LENGTH = 1024
+
+#: Placeholder the owner writes in the transform, e.g.
``unix_timestamp(:value)``.
+#: Matched with word boundaries so ``:values`` is not mistaken for it.
+VALUE_PLACEHOLDER_RE = re.compile(r":value\b")
+
+#: Balanced Jinja blocks. The probe would render these in a different context
+#: at a different time from the chart query, so they are rejected at save time.
+JINJA_BLOCK_RE = re.compile(r"\{\{.*?\}\}|\{%.*?%\}|\{#.*?#\}", re.DOTALL)
+
+#: Substituted for ``:value`` before parsing -- sqlglot rejects a bare
``:value``
+#: on most dialects. Mirrors the ``_JINJA_BLOCK_RE`` -> ``NULL`` trick used by
+#: ``validate_stored_expression``.
+_PARSE_STANDIN = "NULL"
+
+#: Functions whose value depends on wall-clock time or randomness. The probe
+#: runs in a different session at a different moment from the chart query and
+#: its result is then cached, so any of these freezes a snapshot of probe time
+#: into the emitted predicate.
+NON_DETERMINISTIC_FUNCTIONS = {
+ "CURRENT_DATE",
+ "CURRENT_TIME",
+ "CURRENT_TIMESTAMP",
+ "NOW",
+ "RAND",
+ "RANDOM",
+ "UUID",
+}
+
+#: Functions that mean "now" only in their zero-argument form. On Hive and
+#: Impala ``unix_timestamp()`` is the current time while ``unix_timestamp(x)``
+#: -- the canonical transform for this feature -- is pure.
+NON_DETERMINISTIC_WHEN_NILADIC = {"UNIX_TIMESTAMP"}
+
+#: Safe for any function ``T``.
+MIRRORABLE_ALWAYS = {FilterOperator.EQUALS, FilterOperator.IN}
+
+#: Safe only when ``T`` preserves ordering.
+MIRRORABLE_IF_MONOTONIC = {
+ FilterOperator.GREATER_THAN,
+ FilterOperator.GREATER_THAN_OR_EQUALS,
+ FilterOperator.LESS_THAN,
+ FilterOperator.LESS_THAN_OR_EQUALS,
+ FilterOperator.TEMPORAL_RANGE,
+}
+
+
+#: Every operator that can be mirrored under *some* transform. The preview
+#: endpoint accepts these; whether a given one actually mirrors still depends
on
+#: the monotonicity declaration.
+MIRRORABLE_OPERATORS = MIRRORABLE_ALWAYS | MIRRORABLE_IF_MONOTONIC
+
+#: Operators the preview endpoint can construct. `TEMPORAL_RANGE` is mirrored
by
+#: the query path, but as *two* bounds that `_collect_partition_mirror_range`
+#: decomposes a time range into -- the preview request carries sample values,
not
+#: a since/until pair, so there is no range for it to build. The editor never
+#: asks for one either: `previewOperatorFor` sends `>=`, `=` or `IN`.
+PREVIEWABLE_OPERATORS = MIRRORABLE_OPERATORS - {FilterOperator.TEMPORAL_RANGE}
+
+
+def mirrorable_operators(is_monotonic: bool) -> set[FilterOperator]:
+ """
+ The operators whose predicates may be mirrored onto the partition column.
+
+ :param is_monotonic: whether the owner declared the transform
+ order-preserving
+ """
+ if is_monotonic:
+ return MIRRORABLE_ALWAYS | MIRRORABLE_IF_MONOTONIC
+ return set(MIRRORABLE_ALWAYS)
+
+
+#: How wide one bucket of each built-in time grain is.
+#:
+#: Keyed on the ISO duration a filter carries in its ``grain``. Written out
+#: rather than parsed: four of the week grains are ISO *intervals* with an
+#: anchor (``P1W/1970-01-03T00:00:00Z``) that ``isodate.parse_duration``
+#: rejects outright, and ``PT0.5H`` / ``P0.25Y`` are fractional. A literal
+#: table is also the thing a reviewer can check a line at a time.
+#:
+#: ``relativedelta`` rather than ``timedelta`` so the calendar grains stay
+#: calendar arithmetic: a month is not 30 days.
+GRAIN_BUCKET_WIDTHS: dict[str, relativedelta] = {
+ TimeGrain.SECOND: relativedelta(seconds=1),
+ TimeGrain.FIVE_SECONDS: relativedelta(seconds=5),
+ TimeGrain.THIRTY_SECONDS: relativedelta(seconds=30),
+ TimeGrain.MINUTE: relativedelta(minutes=1),
+ TimeGrain.FIVE_MINUTES: relativedelta(minutes=5),
+ TimeGrain.TEN_MINUTES: relativedelta(minutes=10),
+ TimeGrain.FIFTEEN_MINUTES: relativedelta(minutes=15),
+ TimeGrain.THIRTY_MINUTES: relativedelta(minutes=30),
+ TimeGrain.HALF_HOUR: relativedelta(minutes=30),
+ TimeGrain.HOUR: relativedelta(hours=1),
+ TimeGrain.SIX_HOURS: relativedelta(hours=6),
+ TimeGrain.DAY: relativedelta(days=1),
+ TimeGrain.WEEK: relativedelta(days=7),
+ TimeGrain.WEEK_STARTING_SUNDAY: relativedelta(days=7),
+ TimeGrain.WEEK_STARTING_MONDAY: relativedelta(days=7),
+ TimeGrain.WEEK_ENDING_SATURDAY: relativedelta(days=7),
+ TimeGrain.WEEK_ENDING_SUNDAY: relativedelta(days=7),
+ TimeGrain.MONTH: relativedelta(months=1),
+ TimeGrain.QUARTER: relativedelta(months=3),
+ TimeGrain.QUARTER_YEAR: relativedelta(months=3),
+ TimeGrain.YEAR: relativedelta(years=1),
+}
+
+
+def grain_bucket_width(grain: str | None, engine: str) -> relativedelta | None:
+ """
+ How far a grain's truncation can move a timestamp, or ``None`` if unknown.
+
+ A grained filter compares the *truncated* column, so the raw bounds it
+ carries do not describe the rows it keeps. Widening both bounds by one
+ bucket recovers a predicate that is no narrower than the real one -- see
+ `_collect_partition_mirror_range` for the argument. That only works for a
+ grain whose bucket width Superset knows, which excludes anything an
+ operator supplied.
+
+ :param grain: the ISO duration from the filter's ``grain``
+ :param engine: the engine spec's ``engine``, to check per-engine overrides
+ """
+ if not grain:
+ return None
+
+ # An operator-declared grain carries whatever duration string they typed,
+ # and `TIME_GRAIN_ADDON_EXPRESSIONS` can redefine a *built-in* grain's SQL
+ # per engine -- `P1D` could be anything at all. Neither has a width we can
+ # claim to know, so both fall back to not mirroring.
+ if grain in app.config["TIME_GRAIN_ADDONS"]:
+ return None
+ if grain in app.config["TIME_GRAIN_ADDON_EXPRESSIONS"].get(engine, {}):
+ return None
+
+ return GRAIN_BUCKET_WIDTHS.get(grain)
+
+
+@dataclass(frozen=True)
+class PartitionMapping:
+ """A resolved, usable partition filter mapping."""
+
+ partition_column: str
+ mapped_column: str
+ value_transform: str
+ is_monotonic: bool
+
+ def mirrors(self, operator: FilterOperator) -> bool:
+ return operator in mirrorable_operators(self.is_monotonic)
+
+
+def contains_value_placeholder(transform: str | None) -> bool:
+ """Whether the transform contains the ``:value`` placeholder."""
+ return bool(transform) and VALUE_PLACEHOLDER_RE.search(transform or "") is
not None
+
+
+def contains_jinja(transform: str | None) -> bool:
+ """Whether the transform contains a balanced Jinja block."""
+ return bool(transform) and JINJA_BLOCK_RE.search(transform or "") is not
None
+
+
+def parse_skeleton(transform: str) -> str:
+ """
+ The transform with ``:value`` substituted out, ready for a SQL parser.
+
+ ``sanitize_clause`` / sqlglot choke on a bare ``:value`` on most dialects,
+ so the placeholder is swapped for a benign literal first -- the same trick
+ ``validate_stored_expression`` uses for Jinja blocks.
+ """
+ return VALUE_PLACEHOLDER_RE.sub(_PARSE_STANDIN, transform)
+
+
+#: Prefix the transform is wrapped in before parsing. Its length is subtracted
+#: from any reported column so positions refer to what the owner actually
typed.
+_SELECT_PREFIX = "SELECT "
+
+
+def _parse_skeleton(transform: str, engine: str) -> SQLStatement | None:
+ """
+ Parse ``SELECT <transform>`` with the placeholder substituted out.
+
+ Returns ``None`` when the transform does not parse.
+ """
+ try:
+ return SQLStatement(f"{_SELECT_PREFIX}{parse_skeleton(transform)}",
engine)
+ except SupersetParseError:
+ return None
+
+
+def parse_error_detail(transform: str, engine: str) -> str | None:
+ """
+ Where the parser gave up on the transform.
+
+ Returns ``None`` when it parses, or when the parser offered no position.
+ Note that sqlglot parses unknown functions happily -- a misspelled function
+ name is not a parse error, it is an engine error, and surfaces only when
the
+ transform is evaluated.
+
+ The parser's own ``highlight`` is deliberately dropped: it would name the
+ ``NULL`` we substituted for ``:value``, which is not a token the owner
typed.
+ """
+ try:
+ SQLStatement(f"{_SELECT_PREFIX}{parse_skeleton(transform)}", engine)
+ except SupersetParseError as ex:
+ column = (ex.error.extra or {}).get("column")
+ if not isinstance(column, int):
+ return None
+ return str(
+ _(
+ "syntax error at position %(position)d.",
+ position=_position_in_transform(transform, column),
+ )
+ )
+ return None
+
+
+def _position_in_transform(transform: str, parsed_column: int) -> int:
+ """
+ Map a column in the parsed skeleton back to the transform as typed.
+
+ Two substitutions stand between them: the ``SELECT`` prefix, and every
+ ``:value`` that became a shorter ``NULL``. Without unwinding both, a
+ reported position drifts left by two characters per placeholder ahead of it
+ -- which is worst exactly where transforms usually break, at the end.
+ """
+ position = max(parsed_column - len(_SELECT_PREFIX), 0)
+ shift = len(":value") - len(_PARSE_STANDIN)
+ preceding = sum(
+ 1
+ for index, match in enumerate(VALUE_PLACEHOLDER_RE.finditer(transform))
+ if match.start() - index * shift < position
+ )
+ return position + preceding * shift
+
+
+def is_parseable(transform: str | None, engine: str) -> bool:
+ """
+ Whether the transform parses as a single select expression.
+
+ "Single" is the load-bearing word. `SELECT lower(:value), 'x'` parses just
+ as happily as `SELECT lower(:value)`, but it returns two columns per input
+ value, and the probe reads one column per input -- so an `IN` filter would
+ get a predicate built from the wrong halves of the wrong rows. Rejecting
the
+ list here is what lets the probe trust its own column count.
+ """
+ if not transform or not transform.strip():
+ return False
+ statement = _parse_skeleton(transform, engine)
+ return statement is not None and statement.count_select_expressions() == 1
Review Comment:
A transform such as `secret || :value FROM vault` passes the
single-expression and subquery checks even with `ALLOW_ADHOC_SUBQUERY`
disabled, then executes as `SELECT secret || 'us' FROM vault AS v0`, exposing
another table's contents through the preview. Could this require an
expression-only AST with no FROM or other statement clauses before running the
probe?
##########
superset-frontend/src/components/Datasource/components/DatasourceEditor/DatasourceEditor.tsx:
##########
@@ -1250,6 +1356,124 @@ function DatasourceEditor({
[],
);
+ // Which column's row expand to open, for the "map a different column" links.
+ // Consumed by the Columns tab, which scrolls the row into view and expands
it.
+ //
+ // Carries a nonce because the request is an event, not a state: clicking the
+ // same link twice -- after collapsing the row by hand in between -- asks for
+ // the same column name, and a bare string would make the second
+ // `setColumnToReveal` a no-op. React would bail out of the render, nothing
+ // downstream would see a change, and the row would stay shut.
+ const [columnToReveal, setColumnToReveal] = useState<{
+ name: string;
+ nonce: number;
+ } | null>(null);
+
+ const handlePartitionColumnChange = useCallback(
+ (columnName: string | null) => {
+ setDatasource(prev => ({
+ ...prev,
+ partition_column: columnName,
+ partition_mapped_column: nextMappedColumnOverride(
+ prev.partition_mapped_column,
+ columnName,
+ ),
+ }));
+ if (columnName) {
+ setDatabaseColumns(prev =>
+ applyPartitionColumnDefaults(prev, columnName),
Review Comment:
These partition defaults and the edited transform live in `databaseColumns`,
but `syncMetadata` merges source metadata with the older local
`datasource.columns`. Configuring the mapping and then syncing unchanged source
columns before saving restores the old transform, monotonicity and picker
flags, so the save loses those edits; could sync merge against the current
column state?
##########
superset/commands/dataset/importers/v1/utils.py:
##########
@@ -267,6 +272,72 @@ def _get_template_params(dataset: SqlaTable) -> dict[str,
Any]:
return params if isinstance(params, dict) else {}
+def drop_unusable_partition_transforms(config: dict[str, Any]) -> None:
+ """
+ Drop a partition value transform the save path would have rejected.
+
+ Import is the one door into this field that does not go through
+ `UpdateDatasetCommand`, so a bundle can carry a transform holding Jinja or
+ calling a non-deterministic function -- one that will never mirror a
filter,
+ and that the editor can only report as broken after the fact. Dropping it
on
+ the way in leaves the column with no transform, a state an owner can see
and
+ fix, rather than a stored expression that looks configured and is not.
+
+ Deliberately sanitizes rather than raises. A bundle is imported as a whole,
+ and failing someone's entire dataset over one unusable expression is a
worse
+ trade than importing the dataset without it. This is the same bargain
+ `DatasetDAO.clear_dangling_partition_mapping` already strikes for a mapping
+ whose column went away.
+
+ Two kinds of check, and only one of them answers to the feature flag. The
+ usability checks ask whether a transform will ever mirror a filter, which
+ is a question about a live feature; with the flag off nothing mirrors, so
+ rewriting the bundle would discard configuration for no gain. The
+ structural check asks whether the expression is one this product stores at
+ all, and that answer does not change with a flag -- an import written while
+ the flag was off would otherwise sit in the metadata DB fully armed,
+ waiting for an operator to turn the flag on.
+ """
+ columns = config.get("columns")
+ if not columns:
+ return
+ if not any(column.get("partition_value_transform") for column in columns):
+ return
+
+ database =
db.session.query(Database).filter_by(id=config["database_id"]).first()
+ if database is None:
+ return
+
+ check_usability = is_feature_enabled(PARTITION_FILTER_MAPPING)
+ catalog = config.get("catalog")
+ schema = config.get("schema")
+
+ for column in columns:
+ transform = column.get("partition_value_transform")
+ if not transform:
+ continue
+
+ reasons: list[str] = []
+ if reason := stored_expression_error(database, catalog, schema,
transform):
Review Comment:
PUT deliberately saves an unfinished transform such as
`unix_timestamp(:value` as inactive, but exporting and importing that unchanged
dataset hits this parse rejection and silently replaces the transform with
null, even with the feature disabled. Could import preserve the same
non-blocking inactive state instead of discarding configuration that PUT
accepts?
##########
superset-frontend/src/components/Datasource/components/DatasourceEditor/components/PartitionFilterMapping/useDebouncedCommit.ts:
##########
@@ -0,0 +1,100 @@
+/**
+ * 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.
+ */
+import { useCallback, useEffect, useRef, useState } from 'react';
+import { debounce } from 'lodash-es';
+import { Constants } from '@superset-ui/core/components';
+
+/**
+ * Local text state that commits upward on a debounce.
+ *
+ * The dataset editor's commit path is asynchronous and snapshot-based: a value
+ * handed to `onChange` travels Field -> Fieldset -> CollectionTable ->
+ * DatasourceEditor -> DatasourceModal and arrives back on the `value` prop
+ * several renders later, and the callback that commits it was built during an
+ * earlier render. An input driven straight off that prop therefore loses any
+ * keystroke that lands while a commit is in flight -- which is why every other
+ * control in the editor types into local state and commits on a debounce, and
+ * why Fieldset's own comment says the editor assumes exactly that. This is
+ * TextControl's contract (src/explore/components/controls/TextControl), minus
+ * the ControlHeader and the number parsing this field does not want.
+ */
+export function useDebouncedCommit(
+ value: string | null | undefined,
+ commit: (next: string) => void,
+ delay: number = Constants.FAST_DEBOUNCE,
+) {
+ const [localValue, setLocalValue] = useState(value ?? '');
+ const [prevValue, setPrevValue] = useState(value);
+ // What was last handed to `commit`. A commit's own echo arrives as a prop
+ // change like any other, and the re-seed below cannot tell the two apart --
+ // so a commit landing while typing continues (a blur or Enter flushes
+ // immediately) used to reset the input to the committed text and drop every
+ // keystroke since. State rather than a ref because the re-seed reads it
+ // during render, which is exactly what a ref is not for.
+ const [lastCommitted, setLastCommitted] = useState<string | undefined>(
+ undefined,
+ );
+
+ // The commit fires from a timer, so the callback is handed in at call time
+ // rather than captured when the debounce was built -- otherwise it would
+ // close over whichever render happened to create it, and fire a stale
+ // callback that writes to a superseded owner.
+ const debouncedCommit = useRef(
+ debounce((next: string, commitFn: (value: string) => void) => {
+ setLastCommitted(next);
+ commitFn(next);
+ }, delay),
+ );
+
+ useEffect(() => () => debouncedCommit.current.cancel(), []);
+
+ const onChange = useCallback(
+ (next: string) => {
+ setLocalValue(next);
+ debouncedCommit.current(next, commit);
Review Comment:
This queues the commit callback from the input event, whose
`CollectionTable.onFieldsetChange` still captures the whole column collection.
With two rows expanded, a description edit committing before a queued transform
edit can be reverted when that transform commits through the older collection
snapshot; could the delayed commit merge into the latest collection, with both
rows' edits covered together?
##########
superset/daos/dataset.py:
##########
@@ -476,10 +546,98 @@ def _validate_column_date_formats(
"python_date_format is an invalid date/timestamp format."
)
+ @staticmethod
+ def clear_unmapped_partition_transforms(model: SqlaTable) -> None:
+ """
+ Drop the value transform from every column the mapping does not mirror.
+
+ A mapping has exactly one mirrored column, so at most one column may
+ carry a transform. A transform parked on any other column is invisible
+ -- no row but the mapped one renders one -- yet it is still stored, and
+ it goes live the moment the mapped column resolves back to it. Clearing
+ an override is enough to do that: a null `partition_mapped_column`
means
+ "follow `main_dttm_col`", so dropping a mapping would otherwise
activate
+ whatever the default datetime column happened to be holding, turning a
+ request to remove a mapping into a request to add a different one.
+
+ Enforced here rather than in the editor alone because the editor is
only
+ one writer: a PUT, an `override_columns=true` metadata sync and an
+ import all reach the columns directly. The same argument
+ `clear_dangling_partition_mapping` makes about dangling columns.
+
+ `DatasetDAO.update` calls this twice, before and after it applies the
+ request, and the reason is here rather than only at the call site: this
+ reads the mapping as it stands, so a single call can only enforce the
+ invariant against one of the two resolutions a mapping change has.
+
+ Gated on the feature flag, unlike its sibling, because this discards
+ stored configuration rather than repairing a broken reference. With the
+ flag off nothing mirrors, so there is no armed mapping to disarm and
+ clearing would be pure loss.
+ """
+ if not
feature_flag_manager.is_feature_enabled(PARTITION_FILTER_MAPPING_FLAG):
+ return
+
+ mapped_column = (
+ (model.partition_mapped_column or model.main_dttm_col)
+ if model.partition_column
+ else None
+ )
+ # `_upsert_columns` and `_override_columns` insert a new column with
+ # `db.session.add(TableColumn(..., table_id=model.id))` rather than
+ # appending to this relationship, and `BaseDAO.update` does not flush
--
+ # so a column created *and* given a transform in the same request was
+ # invisible here and kept it. That is exactly the stray transform this
+ # method exists to clear, waiting for the mapping to resolve back to
it.
+ db.session.flush()
+ db.session.expire(model, ["columns"])
+ for column in model.columns:
+ if column.column_name == mapped_column:
+ continue
+ if (
+ column.partition_value_transform
+ or column.partition_transform_is_monotonic
+ ):
+ column.partition_value_transform = None
+ column.partition_transform_is_monotonic = False
+
+ @staticmethod
+ def clear_dangling_partition_mapping(
+ model: SqlaTable, surviving_column_names: set[str]
+ ) -> None:
+ """
+ Drop parts of the partition filter mapping whose columns no longer
exist.
+
+ A metadata sync can remove the partition column at the source, which
+ would otherwise leave the dataset pointing at a column that isn't
there.
+ The query layer bails out defensively on a dangling mapping, so this is
+ about the dataset's stored state being honest rather than about
+ correctness of the SQL.
+
+ Called from `update_columns` for both write paths: the editor clears
the
+ mapping client-side too, but `override_columns=true` (an API-driven
+ metadata sync) bypasses the editor entirely, and the upsert path drops
+ every column the payload omits, so either one can take the mapped
+ column out from under the mapping.
+ """
+ if (
Review Comment:
When one PUT renames the partition column from `old_part` to `new_part` and
updates `partition_column`, this child-update cleanup still sees the old
reference and clears the surviving `partition_mapped_column='country'`. Scalar
updates restore only the partition column, and the subsequent cleanup deletes
country's transform, silently switching the mapping to the main datetime
column; could cleanup use the final requested references instead?
##########
superset/models/helpers.py:
##########
@@ -4290,36 +4297,153 @@ def dttm_sql_literal(self, dttm: datetime, col:
"TableColumn") -> str:
return f"""'{dttm.strftime("%Y-%m-%d %H:%M:%S.%f")}'"""
- def get_time_filter( # pylint: disable=too-many-arguments # noqa: C901
+ def _collect_partition_mirror_range(
self,
- time_col: "TableColumn",
+ mapping: Optional["PartitionMapping"],
+ sink: list[tuple[utils.FilterOperator, Any]],
+ column_name: str,
start_dttm: Optional[sa.DateTime],
end_dttm: Optional[sa.DateTime],
- time_grain: Optional[str] = None,
- label: Optional[str] = "__time",
- template_processor: Optional[BaseTemplateProcessor] = None,
- ) -> Optional[ColumnElement]:
- col = (
- time_col.get_timestamp_expression(
- time_grain=time_grain,
- label=label,
- template_processor=template_processor,
- )
- if time_grain
- else self.convert_tbl_column_to_sqla_col(
- time_col, label=label, template_processor=template_processor
- )
+ widen_bounds_by: Optional[relativedelta] = None,
+ ) -> None:
+ """
+ Record a time range for mirroring onto the partition column.
+
+ The bounds are adjusted first, by the same helper `get_time_filter`
uses:
+ the mirrored predicate has to describe the same instants as the
+ timestamp predicate it stands in for, or the pruning is wrong by
exactly
+ the dataset's timezone offset -- silently.
+
+ Either bound may be `None` for an open-ended range, in which case only
+ the bound that exists is mirrored.
+
+ `widen_bounds_by` is one grain bucket width, passed when the filter
this
+ stands in for compares a *truncated* column. The real predicate then
+ keeps rows the raw bounds exclude, and widening both ends by one bucket
+ is the smallest range guaranteed to contain all of them -- see the
"Time
+ grains" section of `superset.connectors.sqla.partition_mapping`. Note
+ the widened upper bound is `<=` rather than `<`: `ts < until + width`
+ only gives `T(ts) <= T(until + width)` for a transform that is
monotonic
+ but not strictly so, such as `unix_timestamp` on a sub-second column.
+
+ A widened range can also be wide enough to fail `_bounds_are_ordered`
+ against a tighter filter on the same column. That fails open -- no
+ mirroring, no pruning, no dropped rows.
+ """
+ if mapping is None or column_name != mapping.mapped_column:
+ return
+ if not mapping.mirrors(utils.FilterOperator.TEMPORAL_RANGE):
+ return
+
+ # Resolve the column so the offset adjustment sees the same type
+ # `get_time_filter` does -- an offset the column's type cannot
represent
+ # is dropped, and the mirrored bounds have to agree with the originals.
+ mapped_col = next(
+ (col for col in self.columns if col.column_name == column_name),
None
)
+ start_dttm, end_dttm = self.adjust_time_bounds(start_dttm, end_dttm,
mapped_col)
+
+ # Adjust first, then widen. `adjust_time_bounds` moves the naive UI
+ # bounds into the frame the column is *stored* in, and the grain
+ # truncation in the real predicate happens in that same frame, so the
+ # bucket width has to be added there too. The two commute for the
+ # hour-offset branch but not for the `ZoneInfo` one, where widening
+ # across a DST boundary first lands an hour out.
+ upper_operator = utils.FilterOperator.LESS_THAN
+ if widen_bounds_by is not None:
+ if start_dttm is not None:
+ start_dttm -= widen_bounds_by
+ if end_dttm is not None:
+ end_dttm += widen_bounds_by
+ upper_operator = utils.FilterOperator.LESS_THAN_OR_EQUALS
+
+ if start_dttm is not None:
+ sink.append((utils.FilterOperator.GREATER_THAN_OR_EQUALS,
start_dttm))
Review Comment:
The probe uses this Python datetime, while the original filter renders it
through `dttm_sql_literal` and the engine's `convert_dttm`, so the bounds can
differ: SQLite truncates a `.500000` lower bound to whole seconds, but a
`julianday(:value)` mirror keeps `.500000` and drops a matching `.250000` row.
Could mirroring transform the same typed temporal literal used by the original
predicate?
##########
superset/datasets/api.py:
##########
@@ -323,6 +401,8 @@ class DatasetRestApi(SoftDeleteApiMixin,
BaseSupersetModelRestApi):
]
show_columns = show_select_columns + [
"columns.type_generic",
+ # Engine-supplied pre-fill for the editor's value transform input.
+ "partition_value_transform_default",
Review Comment:
Saving the dataset from Explore reloads this GET response and replaces the
Explore datasource with it, but the response omits `partition_filter_mapping`.
The pruning indicators consequently disappear after an unrelated save even
though queries still mirror filters; could this expose the same mapping summary
as `SqlaTable.data` and cover that save/reload boundary?
##########
superset/commands/dataset/create.py:
##########
@@ -248,5 +253,52 @@ def validate(self) -> None: # noqa: C901
# so a ``viewers`` key would be dropped by the DAO's ``setattr`` loop.
populate_subjects(self._properties, exceptions, include_viewers=False)
+ self._validate_partition_mapping(exceptions)
+
if exceptions:
raise DatasetInvalidError(exceptions=exceptions)
+
+ def _validate_partition_mapping(self, exceptions: list[ValidationError])
-> None:
+ """
+ Reject a partition mapping that cannot work, at create time too.
+
+ `UpdateDatasetCommand` has validated this since the mapping landed, but
+ create did not, so the same mapping accepted on POST and rejected on
PUT
+ was reachable -- and once stored it is only reported as broken by the
+ editor, which an API-only caller never opens.
+
+ What create can check is narrower than what update can. The dataset's
+ columns are synced by `fetch_metadata` *after* this runs, so "is that a
+ real column?" has no answer yet; both submitted names are therefore
+ passed as known so the existence checks pass trivially and the checks
+ that do not need a column list still run. There is no column payload on
+ POST either, so no transform can arrive here to validate.
+ """
+ if not is_feature_enabled(PARTITION_FILTER_MAPPING):
+ return
+
+ partition_column = self._properties.get("partition_column")
+ database = self._properties.get("database")
+ if not partition_column or not database:
+ # No database means the caller already has a
+ # `DatabaseNotFoundValidationError`; there is no engine to validate
+ # a mapping against and no value in a second error about it.
+ return
+
+ partition_mapped_column =
self._properties.get("partition_mapped_column")
+ for issue in validate_partition_mapping(
+ column_names={
+ name
+ for name in (partition_column, partition_mapped_column)
+ if name is not None
+ },
+ partition_column=partition_column,
+ partition_mapped_column=partition_mapped_column,
+ main_dttm_col=None,
Review Comment:
Passing `main_dttm_col=None` lets POST accept
`partition_column='event_time'` without an override when metadata subsequently
selects that same column as the default datetime column. The resulting
self-mapping then makes even a description-only PUT fail validation; could
create validate the effective mapping after metadata supplies the default?
##########
superset/connectors/sqla/partition_mapping.py:
##########
@@ -0,0 +1,1314 @@
+# 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.
+"""
+Partition filter mapping.
+
+Datasets on Hadoop-family engines are commonly partitioned on a *technical*
+column -- an epoch integer, a lowercased region key -- that no analyst would
+filter on. Unless a query carries a predicate on that column the engine scans
+every partition.
+
+A dataset owner names one partition column ``p``, one business column that
+filters are mirrored from, and a value transform ``T`` (a SQL expression
+containing a ``:value`` placeholder). Superset then appends an equivalent
+predicate on ``p`` to every query, so chart authors change nothing and queries
+prune.
+
+The load-bearing assumption
+---------------------------
+Everything here reasons about ``T(col) op T(v)``, but what is emitted is
+``p op T(v)`` -- a predicate on a *physically different column*. The step from
+one to the other is::
+
+ p = T(mapped_col) for every row in the table
+
+Superset cannot verify that; it is a property of whatever ETL populates the
+partition column. If that job lags, backfills with different logic, or writes
+the partition key in a different timezone, mirrored predicates silently drop
+real rows. The mapping is only as trustworthy as the pipeline behind it.
+
+Operator safety
+---------------
+A mirrored predicate ``P2`` may only be ``AND``-ed onto a query when the
+original predicate ``P1`` *implies* it:
+
+=========================================== =============================
+Original Safe when
+=========================================== =============================
+``col = v``, ``col IN (...)`` always -- ``T`` is a function
+``col >=|>|<|<= v``, ``TEMPORAL_RANGE`` only if ``T`` is monotonic
+``col != v``, ``NOT IN``, ``LIKE``, ... never
+=========================================== =============================
+
+Negations are never safe because ``T`` need not be injective:
+``lower(:value)`` with ``country != 'US'`` mirrors to ``region_key != 'us'``,
+which wrongly excludes rows whose ``country`` is already lowercase ``'us'`` --
+rows the original filter *keeps*.
+
+Monotonicity is a property of the transform, not of the column's data type:
+``hour(:value)``, ``date_format(:value, 'dd')`` and ``dayofweek(:value)`` are
+all reasonable transforms on a ``TIMESTAMP`` column and none of them preserve
+ordering. It is therefore declared by the owner, not inferred.
+
+Time grains
+-----------
+A grained filter -- drill-to-detail, mostly -- compares the *truncated* column,
+``trunc(col) op v``, so the raw bounds it carries do not describe the rows it
+keeps. A row in the final partial bucket satisfies ``trunc(col) < until`` while
+``col < until`` excludes it.
+
+Which direction a grain rounds is not knowable from the duration alone, and the
+obvious guess is wrong: ``WEEK_ENDING_SATURDAY`` rounds *forward* on Hive and
+Presto, and Ocient's grains are ``ROUND``, i.e. to nearest. What every grain
+does satisfy is a bound on the displacement::
+
+ |trunc(ts) - ts| < width(grain)
+
+so widening *both* bounds by one bucket width is no narrower than the real
+predicate whichever way the grain rounds -- and it needs no monotonicity of
+``trunc`` itself, which is what rescues Hive's oddly-anchored ``P1W``. Width is
+a property of the grain, not of the engine, which is what makes this
+maintainable; see `grain_bucket_width`. A grain whose SQL an operator supplied
+has no known width, and does not mirror at all.
+"""
+
+from __future__ import annotations
+
+import hashlib
+import logging
+import re
+from dataclasses import dataclass
+from functools import lru_cache
+from typing import Any, cast, TYPE_CHECKING
+
+import numpy as np
+import pandas as pd
+import sqlalchemy as sa
+from dateutil.relativedelta import relativedelta
+from flask import current_app as app
+from flask_babel import lazy_gettext as _
+from sqlalchemy.engine.interfaces import Dialect
+from sqlalchemy.sql.elements import ColumnElement
+
+from superset.constants import LRU_CACHE_MAX_SIZE, TimeGrain
+from superset.exceptions import (
+ QueryClauseValidationException,
+ SupersetParseError,
+ SupersetSecurityException,
+)
+from superset.extensions import cache_manager, feature_flag_manager
+from superset.sql.parse import SQLStatement
+from superset.utils import json
+from superset.utils.core import FilterOperator
+
+if TYPE_CHECKING:
+ from superset.connectors.sqla.models import SqlaTable, TableColumn
+ from superset.models.core import Database
+
+logger = logging.getLogger(__name__)
+
+FEATURE_FLAG = "PARTITION_FILTER_MAPPING"
+
+#: Longest transform accepted anywhere. A transform is one SQL expression an
+#: owner types by hand, so this is generous; the point is that it is bounded.
+#: Every entry point enforces it -- the typed column field, the import schema
+#: and the preview request -- because `is_transform_active` parses the stored
+#: value on each Explore load, and an unbounded string would make that parse
+#: the expensive part of rendering a chart.
+MAX_TRANSFORM_LENGTH = 1024
+
+#: Placeholder the owner writes in the transform, e.g.
``unix_timestamp(:value)``.
+#: Matched with word boundaries so ``:values`` is not mistaken for it.
+VALUE_PLACEHOLDER_RE = re.compile(r":value\b")
+
+#: Balanced Jinja blocks. The probe would render these in a different context
+#: at a different time from the chart query, so they are rejected at save time.
+JINJA_BLOCK_RE = re.compile(r"\{\{.*?\}\}|\{%.*?%\}|\{#.*?#\}", re.DOTALL)
+
+#: Substituted for ``:value`` before parsing -- sqlglot rejects a bare
``:value``
+#: on most dialects. Mirrors the ``_JINJA_BLOCK_RE`` -> ``NULL`` trick used by
+#: ``validate_stored_expression``.
+_PARSE_STANDIN = "NULL"
+
+#: Functions whose value depends on wall-clock time or randomness. The probe
+#: runs in a different session at a different moment from the chart query and
+#: its result is then cached, so any of these freezes a snapshot of probe time
+#: into the emitted predicate.
+NON_DETERMINISTIC_FUNCTIONS = {
+ "CURRENT_DATE",
+ "CURRENT_TIME",
+ "CURRENT_TIMESTAMP",
+ "NOW",
+ "RAND",
+ "RANDOM",
+ "UUID",
+}
+
+#: Functions that mean "now" only in their zero-argument form. On Hive and
+#: Impala ``unix_timestamp()`` is the current time while ``unix_timestamp(x)``
+#: -- the canonical transform for this feature -- is pure.
+NON_DETERMINISTIC_WHEN_NILADIC = {"UNIX_TIMESTAMP"}
+
+#: Safe for any function ``T``.
+MIRRORABLE_ALWAYS = {FilterOperator.EQUALS, FilterOperator.IN}
+
+#: Safe only when ``T`` preserves ordering.
+MIRRORABLE_IF_MONOTONIC = {
+ FilterOperator.GREATER_THAN,
+ FilterOperator.GREATER_THAN_OR_EQUALS,
+ FilterOperator.LESS_THAN,
+ FilterOperator.LESS_THAN_OR_EQUALS,
+ FilterOperator.TEMPORAL_RANGE,
+}
+
+
+#: Every operator that can be mirrored under *some* transform. The preview
+#: endpoint accepts these; whether a given one actually mirrors still depends
on
+#: the monotonicity declaration.
+MIRRORABLE_OPERATORS = MIRRORABLE_ALWAYS | MIRRORABLE_IF_MONOTONIC
+
+#: Operators the preview endpoint can construct. `TEMPORAL_RANGE` is mirrored
by
+#: the query path, but as *two* bounds that `_collect_partition_mirror_range`
+#: decomposes a time range into -- the preview request carries sample values,
not
+#: a since/until pair, so there is no range for it to build. The editor never
+#: asks for one either: `previewOperatorFor` sends `>=`, `=` or `IN`.
+PREVIEWABLE_OPERATORS = MIRRORABLE_OPERATORS - {FilterOperator.TEMPORAL_RANGE}
+
+
+def mirrorable_operators(is_monotonic: bool) -> set[FilterOperator]:
+ """
+ The operators whose predicates may be mirrored onto the partition column.
+
+ :param is_monotonic: whether the owner declared the transform
+ order-preserving
+ """
+ if is_monotonic:
+ return MIRRORABLE_ALWAYS | MIRRORABLE_IF_MONOTONIC
+ return set(MIRRORABLE_ALWAYS)
+
+
+#: How wide one bucket of each built-in time grain is.
+#:
+#: Keyed on the ISO duration a filter carries in its ``grain``. Written out
+#: rather than parsed: four of the week grains are ISO *intervals* with an
+#: anchor (``P1W/1970-01-03T00:00:00Z``) that ``isodate.parse_duration``
+#: rejects outright, and ``PT0.5H`` / ``P0.25Y`` are fractional. A literal
+#: table is also the thing a reviewer can check a line at a time.
+#:
+#: ``relativedelta`` rather than ``timedelta`` so the calendar grains stay
+#: calendar arithmetic: a month is not 30 days.
+GRAIN_BUCKET_WIDTHS: dict[str, relativedelta] = {
+ TimeGrain.SECOND: relativedelta(seconds=1),
+ TimeGrain.FIVE_SECONDS: relativedelta(seconds=5),
+ TimeGrain.THIRTY_SECONDS: relativedelta(seconds=30),
+ TimeGrain.MINUTE: relativedelta(minutes=1),
+ TimeGrain.FIVE_MINUTES: relativedelta(minutes=5),
+ TimeGrain.TEN_MINUTES: relativedelta(minutes=10),
+ TimeGrain.FIFTEEN_MINUTES: relativedelta(minutes=15),
+ TimeGrain.THIRTY_MINUTES: relativedelta(minutes=30),
+ TimeGrain.HALF_HOUR: relativedelta(minutes=30),
+ TimeGrain.HOUR: relativedelta(hours=1),
+ TimeGrain.SIX_HOURS: relativedelta(hours=6),
+ TimeGrain.DAY: relativedelta(days=1),
+ TimeGrain.WEEK: relativedelta(days=7),
+ TimeGrain.WEEK_STARTING_SUNDAY: relativedelta(days=7),
+ TimeGrain.WEEK_STARTING_MONDAY: relativedelta(days=7),
+ TimeGrain.WEEK_ENDING_SATURDAY: relativedelta(days=7),
+ TimeGrain.WEEK_ENDING_SUNDAY: relativedelta(days=7),
+ TimeGrain.MONTH: relativedelta(months=1),
+ TimeGrain.QUARTER: relativedelta(months=3),
+ TimeGrain.QUARTER_YEAR: relativedelta(months=3),
+ TimeGrain.YEAR: relativedelta(years=1),
+}
+
+
+def grain_bucket_width(grain: str | None, engine: str) -> relativedelta | None:
+ """
+ How far a grain's truncation can move a timestamp, or ``None`` if unknown.
+
+ A grained filter compares the *truncated* column, so the raw bounds it
+ carries do not describe the rows it keeps. Widening both bounds by one
+ bucket recovers a predicate that is no narrower than the real one -- see
+ `_collect_partition_mirror_range` for the argument. That only works for a
+ grain whose bucket width Superset knows, which excludes anything an
+ operator supplied.
+
+ :param grain: the ISO duration from the filter's ``grain``
+ :param engine: the engine spec's ``engine``, to check per-engine overrides
+ """
+ if not grain:
+ return None
+
+ # An operator-declared grain carries whatever duration string they typed,
+ # and `TIME_GRAIN_ADDON_EXPRESSIONS` can redefine a *built-in* grain's SQL
+ # per engine -- `P1D` could be anything at all. Neither has a width we can
+ # claim to know, so both fall back to not mirroring.
+ if grain in app.config["TIME_GRAIN_ADDONS"]:
+ return None
+ if grain in app.config["TIME_GRAIN_ADDON_EXPRESSIONS"].get(engine, {}):
+ return None
+
+ return GRAIN_BUCKET_WIDTHS.get(grain)
+
+
+@dataclass(frozen=True)
+class PartitionMapping:
+ """A resolved, usable partition filter mapping."""
+
+ partition_column: str
+ mapped_column: str
+ value_transform: str
+ is_monotonic: bool
+
+ def mirrors(self, operator: FilterOperator) -> bool:
+ return operator in mirrorable_operators(self.is_monotonic)
+
+
+def contains_value_placeholder(transform: str | None) -> bool:
+ """Whether the transform contains the ``:value`` placeholder."""
+ return bool(transform) and VALUE_PLACEHOLDER_RE.search(transform or "") is
not None
+
+
+def contains_jinja(transform: str | None) -> bool:
+ """Whether the transform contains a balanced Jinja block."""
+ return bool(transform) and JINJA_BLOCK_RE.search(transform or "") is not
None
+
+
+def parse_skeleton(transform: str) -> str:
+ """
+ The transform with ``:value`` substituted out, ready for a SQL parser.
+
+ ``sanitize_clause`` / sqlglot choke on a bare ``:value`` on most dialects,
+ so the placeholder is swapped for a benign literal first -- the same trick
+ ``validate_stored_expression`` uses for Jinja blocks.
+ """
+ return VALUE_PLACEHOLDER_RE.sub(_PARSE_STANDIN, transform)
+
+
+#: Prefix the transform is wrapped in before parsing. Its length is subtracted
+#: from any reported column so positions refer to what the owner actually
typed.
+_SELECT_PREFIX = "SELECT "
+
+
+def _parse_skeleton(transform: str, engine: str) -> SQLStatement | None:
+ """
+ Parse ``SELECT <transform>`` with the placeholder substituted out.
+
+ Returns ``None`` when the transform does not parse.
+ """
+ try:
+ return SQLStatement(f"{_SELECT_PREFIX}{parse_skeleton(transform)}",
engine)
+ except SupersetParseError:
+ return None
+
+
+def parse_error_detail(transform: str, engine: str) -> str | None:
+ """
+ Where the parser gave up on the transform.
+
+ Returns ``None`` when it parses, or when the parser offered no position.
+ Note that sqlglot parses unknown functions happily -- a misspelled function
+ name is not a parse error, it is an engine error, and surfaces only when
the
+ transform is evaluated.
+
+ The parser's own ``highlight`` is deliberately dropped: it would name the
+ ``NULL`` we substituted for ``:value``, which is not a token the owner
typed.
+ """
+ try:
+ SQLStatement(f"{_SELECT_PREFIX}{parse_skeleton(transform)}", engine)
+ except SupersetParseError as ex:
+ column = (ex.error.extra or {}).get("column")
+ if not isinstance(column, int):
+ return None
+ return str(
+ _(
+ "syntax error at position %(position)d.",
+ position=_position_in_transform(transform, column),
+ )
+ )
+ return None
+
+
+def _position_in_transform(transform: str, parsed_column: int) -> int:
+ """
+ Map a column in the parsed skeleton back to the transform as typed.
+
+ Two substitutions stand between them: the ``SELECT`` prefix, and every
+ ``:value`` that became a shorter ``NULL``. Without unwinding both, a
+ reported position drifts left by two characters per placeholder ahead of it
+ -- which is worst exactly where transforms usually break, at the end.
+ """
+ position = max(parsed_column - len(_SELECT_PREFIX), 0)
+ shift = len(":value") - len(_PARSE_STANDIN)
+ preceding = sum(
+ 1
+ for index, match in enumerate(VALUE_PLACEHOLDER_RE.finditer(transform))
+ if match.start() - index * shift < position
+ )
+ return position + preceding * shift
+
+
+def is_parseable(transform: str | None, engine: str) -> bool:
+ """
+ Whether the transform parses as a single select expression.
+
+ "Single" is the load-bearing word. `SELECT lower(:value), 'x'` parses just
+ as happily as `SELECT lower(:value)`, but it returns two columns per input
+ value, and the probe reads one column per input -- so an `IN` filter would
+ get a predicate built from the wrong halves of the wrong rows. Rejecting
the
+ list here is what lets the probe trust its own column count.
+ """
+ if not transform or not transform.strip():
+ return False
+ statement = _parse_skeleton(transform, engine)
+ return statement is not None and statement.count_select_expressions() == 1
+
+
+def find_non_deterministic_functions(transform: str, engine: str) -> set[str]:
+ """
+ Names of non-deterministic functions the transform calls.
+
+ ``UNIX_TIMESTAMP`` is only reported in its zero-argument form, which means
+ "now" on Hive and Impala; the one-argument form is the canonical temporal
+ transform and stays allowed.
+ """
+ statement = _parse_skeleton(transform, engine)
+ if statement is None:
+ return set()
+
+ found = {
+ name
+ for name in NON_DETERMINISTIC_FUNCTIONS
+ if statement.check_functions_present({name})
+ }
+ return found | _find_niladic_calls(statement)
+
+
+def _find_niladic_calls(statement: SQLStatement) -> set[str]:
+ """
+ Names from ``NON_DETERMINISTIC_WHEN_NILADIC`` called with no arguments.
+
+ Note some dialects resolve the zero-argument form themselves -- Hive parses
+ ``unix_timestamp()`` straight to ``CURRENT_TIMESTAMP`` -- in which case the
+ name-based check above has already caught it. This is the backstop for the
+ dialects that do not.
+ """
+ return NON_DETERMINISTIC_WHEN_NILADIC & statement.get_niladic_functions()
+
+
+def resolve_partition_mapping(datasource: SqlaTable) -> PartitionMapping |
None:
+ """
+ Resolve the dataset's mapping, or ``None`` when nothing may be mirrored.
+
+ Every bail-out here is defensive as well as functional: save-time
validation
+ rejects most of these, but rows predating the validation can still violate
+ the invariants, and a column sync can invalidate a mapping that was fine
+ when it was written.
+
+ The transform gate is `is_transform_active`, the same function the Explore
+ indicator reads, which is `validate_transform` with the messages discarded.
+ Anything narrower here would be a second, weaker statement of the same
rule:
+ a transform calling `now()` that reached storage without passing
+ `UpdateDatasetCommand` -- through import, or through a bundle written by
hand
+ -- would be reported inactive by the editor and still mirrored by this
+ function, freezing a snapshot of probe time into the predicate with nothing
+ on screen to say so.
+ """
+ if not feature_flag_manager.is_feature_enabled(FEATURE_FLAG):
+ return None
+
+ partition_column = getattr(datasource, "partition_column", None)
+ if not partition_column:
+ return None
+
+ columns_by_name = {column.column_name: column for column in
datasource.columns}
+ if partition_column not in columns_by_name:
+ # The partition column was dropped by a column sync or at the source.
+ return None
+
+ mapped_column_name = (
+ getattr(datasource, "partition_mapped_column", None) or
datasource.main_dttm_col
+ )
+ if not mapped_column_name or mapped_column_name not in columns_by_name:
+ return None
+
+ if mapped_column_name == partition_column:
+ # Self-mapping: the mirrored predicate would duplicate the original.
+ return None
+
+ mapped_column = columns_by_name[mapped_column_name]
+ transform = getattr(mapped_column, "partition_value_transform", None)
+ if not is_transform_active(transform, datasource.database.backend):
+ return None
+
+ if has_active_advanced_data_type(mapped_column):
+ # `translate_filter` builds its own predicate shape from *translated*
+ # values, so the `(operator, value)` pair the operator matrix reasons
+ # about does not exist and mirroring would apply the wrong values.
+ return None
+
+ return PartitionMapping(
+ partition_column=str(partition_column),
+ mapped_column=str(mapped_column_name),
+ value_transform=cast(str, transform),
+ is_monotonic=bool(
+ getattr(mapped_column, "partition_transform_is_monotonic", False)
+ ),
+ )
+
+
+def has_active_advanced_data_type(column: TableColumn) -> bool:
+ """
+ Whether the column's advanced data type is configured and switched on.
+
+ Such a column is never mirrored: ``translate_filter`` builds its own
+ predicate shape from translated values, so the ``(operator, value)`` pair
+ the operator matrix reasons about does not exist.
+ """
+ advanced_data_type = getattr(column, "advanced_data_type", None)
+ if not advanced_data_type:
+ return False
+ if not
feature_flag_manager.is_feature_enabled("ENABLE_ADVANCED_DATA_TYPES"):
+ return False
+ return advanced_data_type in app.config.get("ADVANCED_DATA_TYPES", {})
+
+
+def stored_expression_error(
+ database: "Database",
+ catalog: str | None,
+ schema: str | None,
+ transform: str,
+) -> str | None:
+ """
+ Why this transform may not be stored or run, if there is a reason.
+
+ The same subquery, function-denylist and RLS policy every other stored
+ expression goes through, applied to a value transform. It matters more
+ here than the name suggests: `build_probe_sql` binds only `:value` and
+ splices the rest of the transform in as SQL text, which the engine then
+ executes, so an ungated transform is arbitrary SQL. A dataset editor
+ without SQL Lab could store `(SELECT secret FROM protected_table LIMIT 1)
+ || :value` and read the answer back out of the emitted predicate.
+
+ Returns the engine-agnostic reason as a string rather than raising,
+ because its four callers disagree about what to do with it: the preview
+ reports it at the field, a PUT and the legacy datasource save refuse the
+ write, the importer drops the transform and keeps the dataset, and the
+ probe simply declines to run. Raising would make three of those four
+ write a `try` around a question.
+
+ Imported inside the function rather than at module scope:
+ `connectors.sqla.models` imports this module, so the dependency only runs
+ one way at import time.
+ """
+ from superset.connectors.sqla.models import ( # pylint:
disable=import-outside-toplevel,cyclic-import
+ validate_stored_expression,
+ )
+
+ try:
+ validate_stored_expression(database, catalog, schema,
parse_skeleton(transform))
+ except SupersetSecurityException as ex:
+ return str(ex.error.message)
+ except QueryClauseValidationException as ex:
+ return str(ex.message)
+ return None
+
+
+def build_probe_sql(
+ transform: str,
+ values: list[Any],
+ dialect: Dialect | None = None,
+ from_suffix: str = "",
+) -> str:
+ """
+ Compile a single ``SELECT`` that evaluates the transform at every value.
+
+ Values are attacker-controlled (a Gamma user picks filter values), so they
+ are bound as parameters and rendered by the dialect's own literal processor
+ rather than interpolated into the SQL text.
+
+ ``from_suffix`` comes from the engine spec's
``select_without_from_suffix``.
+ A ``SELECT`` with no ``FROM`` is not universal SQL: Oracle and Db2 need a
+ one-row table to select from, and without it every probe on those engines
+ raises -- which `_run_probe` swallows, so the only symptom is a correctly
+ configured mapping that silently never prunes.
+
+ Note this deliberately does *not* go through ``BaseEngineSpec``'s text
+ helper, which escapes ``:`` on every engine but Athena and would destroy
the
+ ``:value`` placeholder before it can be bound.
+ """
+ selections = []
+ for index, value in enumerate(values):
+ clause = sa.text(transform).bindparams(sa.bindparam("value",
value=value))
+ compiled = clause.compile(
+ dialect=dialect,
+ compile_kwargs={"literal_binds": True},
+ )
+ selections.append(f"{compiled} AS v{index}")
+ return "SELECT " + ", ".join(selections) + from_suffix
+
+
+def evaluate_transform(
+ database: Database,
+ catalog: str | None,
+ schema: str | None,
+ transform: str,
+ values: list[Any],
+ *,
+ errors: list[str] | None = None,
+) -> list[Any] | None:
+ """
+ Evaluate ``transform`` against the engine once per distinct value.
+
+ Returns one result per input value, positionally aligned with ``values``,
or
+ ``None`` if anything at all goes wrong. Failing open costs pruning, never
+ correctness: the chart query still runs, it just scans more partitions.
+
+ The probe is pinned to the dataset's catalog and schema so session settings
+ match the chart query as closely as the connection pool allows. It still
+ runs in a *different* session, which is why transforms calling
+ session-dependent functions are rejected at save time.
+
+ :param errors: optional sink for the engine's own account of a failure. The
+ query path passes nothing and stays silent; the editor's preview passes
+ a list so it can tell the owner *why* the transform did not evaluate --
+ a misspelled function is the common case and sqlglot parses it happily.
+ """
+ if not values:
+ return None
+
+ # Dedupe so a 200-value `IN` list costs one column, not 200.
+ distinct: list[Any] = []
+ seen: set[Any] = set()
+ for value in values:
+ key = _hashable(value)
+ if key not in seen:
+ seen.add(key)
+ distinct.append(value)
+
+ cache_key = _probe_cache_key(database, catalog, schema, transform,
distinct)
+ cached = _cache_get(cache_key)
+ if cached is None:
+ cached = _run_probe(
+ database, catalog, schema, transform, distinct, errors=errors
+ )
+ if cached is None:
+ # Deliberately not cached: a transient engine blip would otherwise
+ # keep the dataset pruning-free for the whole cache timeout.
+ return None
+ _cache_set(cache_key, cached)
+
+ evaluated = dict(
+ zip((_hashable(value) for value in distinct), cached, strict=False)
+ )
+ return [evaluated[_hashable(value)] for value in values]
+
+
+def _run_probe(
+ database: Database,
+ catalog: str | None,
+ schema: str | None,
+ transform: str,
+ distinct: list[Any],
+ *,
+ errors: list[str] | None = None,
+) -> list[Any] | None:
+ # The last gate before the transform becomes SQL the engine runs, and the
+ # only one that covers a transform already in storage. The write-side
+ # checks can only speak for rows written after they existed; this speaks
+ # for every row, including ones a pre-fix release stored and ones a door
+ # that forgets to ask still lets in.
+ #
+ # Declining costs pruning, never correctness -- see `evaluate_transform`.
+ if reason := stored_expression_error(database, catalog, schema, transform):
+ logger.warning(
+ "Refusing to probe a partition transform that is not a storable "
+ "expression; queries will not prune: %s",
+ reason,
+ )
+ return None
+
+ try:
+ sql = build_probe_sql(
+ transform,
+ distinct,
+ _dialect_for(database),
+ database.db_engine_spec.select_without_from_suffix,
+ )
+ frame = database.get_df(sql=sql, catalog=catalog, schema=schema)
+ if frame is None or frame.empty:
+ logger.warning(
+ "Partition transform probe returned no rows; skipping
mirroring"
+ )
+ return None
+ row = frame.iloc[0]
Review Comment:
Reading the probe through `frame.iloc[0]` coerces mixed numeric columns to a
common float dtype. On SQLite, `CAST(:value AS NUMERIC)` for an IN filter
containing `'9007199254740993'` and `'1.5'` changes the first exact result to
`9007199254740992.0`, so the mirror drops a matching row; could each probe cell
be read without row-wide dtype coercion?
--
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]