sadpandajoe commented on code in PR #43870:
URL: https://github.com/apache/superset/pull/43870#discussion_r4186147475


##########
superset/connectors/sqla/partition_mapping.py:
##########
@@ -0,0 +1,1913 @@
+# 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 (...)``                ``T`` is a function *and* the
+                                             engine compares ``col``
+                                             byte-exactly
+``col >=|>|<|<= v``, ``TEMPORAL_RANGE``      only if ``T`` is monotonic
+``col != v``, ``NOT IN``, ``LIKE``, ...      never
+===========================================  =============================
+
+Equality under the engine's own rules
+-------------------------------------
+"``T`` is a function" gives ``col = v`` implies ``T(col) = T(v)`` for *value*
+equality. The engine compares with *SQL* equality, and the two part company
+wherever that comparison is not byte-exact: under a case-insensitive collation
+a stored ``country`` of ``'us'`` satisfies a filter for ``'US'``, while the
+mirror ``region_key = hex('US')`` excludes the row and the chart loses it with
+nothing to indicate why. Trailing-space-insensitive ``CHAR`` comparison says
+the same thing about padding.
+
+So equality and ``IN`` are withdrawn for a *string* mapped column unless the
+engine spec declares ``binary_string_comparison`` -- `equality_mirrors_safely`.
+Numeric and temporal comparison is exact everywhere and is untouched, and range
+mirroring is untouched too: it already rests on the owner declaring ``T``
+order-preserving *with respect to the column's own order*, which is an
+assertion about the very semantics this is checking for.
+
+What the gate cannot see is a column declaring its own collation
+(``country COLLATE NOCASE``) on an otherwise byte-exact engine. Nothing in
+SQLAlchemy's reflection or the engine specs exposes that, so there it remains
+the owner's assumption, like ``p = T(mapped_col)`` itself.
+
+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 collections.abc import Sequence
+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 SQLScript, SQLStatement
+from superset.superset_typing import FilterValues
+from superset.utils import core as utils, json
+from superset.utils.core import FilterOperator
+
+if TYPE_CHECKING:
+    from superset.connectors.sqla.models import SqlaTable, TableColumn
+    from superset.db_engine_specs.base import BaseEngineSpec
+    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.
+#: The bound matters 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.
+#:
+#: The typed column field, the import schema and the preview request each
+#: enforce it through a marshmallow `Length` validator, which is what produces
+#: the per-field message those three callers want. `stored_expression_error`
+#: enforces it too, for the paths that load no schema -- the deprecated
+#: `POST /datasource/save/` is one, and it writes straight through
+#: `update_from_object`.
+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"
+
+#: Substituted for ``:value`` to find out *where* the placeholder landed, which
+#: `_PARSE_STANDIN` cannot answer: ``NULL`` is a keyword, so after substitution
+#: the parse tree no longer records that a placeholder was ever there. An
+#: identifier does leave a trace -- a column reference -- and one that is
+#: missing is a placeholder the engine will never evaluate. Named so it cannot
+#: plausibly collide with a real column; a collision only makes the count
+#: disagree, which fails closed.
+_PLACEHOLDER_STANDIN = "superset_pfm_value_standin"
+
+#: 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, *, equality_is_safe: bool = True
+) -> set[FilterOperator]:
+    """
+    The operators whose predicates may be mirrored onto the partition column.
+
+    :param is_monotonic: whether the owner declared the transform
+        order-preserving
+    :param equality_is_safe: whether the engine's ``=`` on the mapped column
+        compares the way ``T`` was written against -- see
+        `equality_mirrors_safely`
+
+    Both call sites -- the query path through `PartitionMapping.mirrors` and 
the
+    Explore indicator through `partition_filter_mapping_summary` -- come 
through
+    here, which is what keeps the glyph from promising pruning the SQL will not
+    do.
+    """
+    operators = MIRRORABLE_ALWAYS if equality_is_safe else set()
+    if is_monotonic:
+        operators = operators | MIRRORABLE_IF_MONOTONIC
+    return set(operators)
+
+
+def equality_mirrors_safely(
+    column: TableColumn | None,
+    db_engine_spec: type[BaseEngineSpec],
+) -> bool:
+    """
+    Whether ``col = v`` on this column implies ``T(col) = T(v)`` in the engine.
+
+    The operator matrix calls equality safe "for any function ``T``" because
+    ``T`` is a function, so ``col = v`` gives ``T(col) = T(v)``. That reasons
+    about *value* equality while the engine reasons about *SQL* equality, and
+    the two part company under any comparison that is not byte-exact: with a
+    case-insensitive collation a stored ``country`` of ``'us'`` satisfies a
+    filter for ``'US'``, but the mirror ``region_key = hex('US')`` excludes the
+    row and the chart loses it with no indication why. 
Trailing-space-insensitive
+    ``CHAR`` comparison says the same thing about padding.
+
+    Only string columns are at risk; numeric and temporal comparison is exact
+    everywhere. And only equality: range mirroring already requires the owner
+    to declare ``T`` order-preserving *with respect to the column's own order*,
+    which is an assertion about the same comparison semantics this is checking
+    for -- so the declaration covers it where equality has nothing to cover it.
+
+    What this cannot see is a column that declares its own collation
+    (``country COLLATE NOCASE``) on an otherwise byte-exact engine. Neither
+    SQLAlchemy's reflection nor the engine specs expose it, so on such a column
+    the assumption stays where ``p = T(mapped_col)`` already is: with the 
owner.
+    """
+    if column is None:
+        return True
+    try:
+        is_string = column.type_generic == utils.GenericDataType.STRING
+    except Exception:  # pylint: disable=broad-except  # noqa: BLE001
+        # An unresolvable type is treated as a string: the gate exists to stop
+        # a silent wrong answer, so it fails closed.
+        return False
+    if not is_string:
+        return True
+    return bool(db_engine_spec.binary_string_comparison)
+
+
+#: 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
+    #: Whether the engine compares the mapped column the way ``T`` was written
+    #: against. Resolved once where the column and the engine are both in hand,
+    #: rather than re-derived at each `mirrors` call.
+    equality_is_safe: bool = True
+
+    def mirrors(self, operator: FilterOperator) -> bool:
+        return operator in mirrorable_operators(
+            self.is_monotonic, equality_is_safe=self.equality_is_safe
+        )
+
+
+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, standin: str = _PARSE_STANDIN) -> 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.
+
+    :param standin: what to substitute. `placeholder_is_executable` passes
+        `_PLACEHOLDER_STANDIN` to find out where the placeholder landed.
+    """
+    return VALUE_PLACEHOLDER_RE.sub(standin, transform)

Review Comment:
   For an array filter whose rendered literal contains escaped backslashes, 
this regex replacement changes `array('a\\nb')` into `array('a\nb')`, so 
ClickHouse probes a newline instead of the literal backslash sequence and the 
resulting partition predicate can exclude matching rows. Could the rendered SQL 
be substituted literally rather than interpreted as a replacement string?



##########
superset/models/helpers.py:
##########
@@ -5317,6 +5717,74 @@ def get_sqla_query(  # pylint: 
disable=too-many-arguments,too-many-locals,too-ma
                     db_extra=self.db_extra,
                 )
 
+                # Hoisted above the mirror collection below, which has to
+                # record the values the real predicate *binds*. SQLAlchemy
+                # infers an `IN` bind's type from the first element, so a mixed
+                # int/float list has to be widened before binding or every
+                # later float is truncated (see #33206). The `IN` branch did
+                # this itself, 100-odd lines down and by rebinding `eq` -- so
+                # the mirror had already copied the unwidened values and probed
+                # the transform at a number the query never compares.
+                #
+                # Scoped to the operators whose branch did it, so that moving
+                # the call changes when it runs and not what it applies to:
+                # `CONTAINS_ANY`/`CONTAINS_ALL` are list targets too and were
+                # never widened.
+                if (
+                    target_generic_type == utils.GenericDataType.NUMERIC
+                    and op in {utils.FilterOperator.IN, 
utils.FilterOperator.NOT_IN}
+                    and isinstance(eq, (list, tuple))
+                ):
+                    eq = normalize_mixed_numbers(eq)
+
+                # Mirror onto the partition column. This single site covers
+                # ad-hoc filters, dashboard native filters and cross-filters,
+                # because they all arrive as entries in `filter`. `eq` rather
+                # than the raw `val`: it is the value the real predicate binds,
+                # which is only true because the numeric widening above now
+                # runs before this rather than inside the `IN` branch below.
+                # `TEMPORAL_RANGE` is collected in its own branch below, where
+                # the range has been resolved into a pair of bounds.
+                #
+                # A grain makes the real predicate compare the *truncated*
+                # column, so the raw value no longer describes the rows it
+                # matches: drill-to-detail sends `==` on a bucket start, which
+                # every row in the bucket satisfies once truncated. Mirroring 
it
+                # raw would keep only the bucket's first instant.
+                #
+                # Grained ranges do mirror, by widening (see the TEMPORAL_RANGE
+                # branch below). The same trick is deferred rather than unsafe
+                # here: `==` on a bucket would have to become a two-sided 
range,
+                # which needs a monotonic transform -- `EQUALS` mirrors without
+                # one today -- and a parse back into a datetime; and a grained
+                # `IN` would need a union of buckets, which this sink cannot
+                # express because its entries are AND-ed.
+                if (
+                    col_obj is not None
+                    and op != utils.FilterOperator.TEMPORAL_RANGE
+                    and not filter_grain
+                ):
+                    mirror_value: Any = eq

Review Comment:
   For TIMESTAMP equality/IN filters supplied as epoch milliseconds, the 
handler returns SQLAlchemy expressions, but these are forwarded as bind values: 
probe compilation fails and silently disables partition pruning, while preview 
wraps the same inputs as `RawProbeValue`. Could runtime normalize these 
SQL-valued inputs the same way before collecting the mirror?



##########
superset/commands/dataset/importers/v1/utils.py:
##########
@@ -582,6 +668,12 @@ def import_dataset(  # noqa: C901
     if dataset.id is None:
         db.session.flush()
 
+    # A bundle can name a mapped column and still carry transforms on other
+    # columns; the editor's client-side guard never runs here. Left in place, a
+    # transform on an unmirrored column is invisible and still stored, ready to
+    # go live the moment the mapped column resolves back to it.
+    DatasetDAO.clear_unmapped_partition_transforms(dataset)

Review Comment:
   An overwrite import that clears the mapped-column override can activate an 
old transform parked on the default datetime column: omitted column fields 
retain their stored values, and this post-import cleanup now treats that column 
as effective and skips it. Could import disarm parked transforms against the 
old mapping before applying the new references, as `DatasetDAO.update` does?



##########
superset/datasets/schemas.py:
##########
@@ -347,6 +373,18 @@ def fix_extra(self, data: dict[str, Any], **kwargs: Any) 
-> dict[str, Any]:
     datetime_format = fields.String(
         allow_none=True, validate=[Length(1, 100), validate_python_date_format]
     )
+    partition_value_transform = fields.String(
+        allow_none=True, validate=Length(1, MAX_TRANSFORM_LENGTH)
+    )
+    # `load_default` so bundles predating the field do not claim their 
transform
+    # preserves ordering, which would silently enable range mirroring on 
import.
+    # `allow_none` because the column is nullable on purpose -- the legacy
+    # datasource editor writes NULL for any field its payload omits -- and 
export
+    # emits every field unconditionally, so an untouched export of such a 
dataset
+    # carries an explicit null that import would otherwise refuse outright.
+    partition_transform_is_monotonic = fields.Boolean(
+        allow_none=True, load_default=False
+    )
     uuid = fields.UUID(allow_none=True)
 
 

Review Comment:
   On overwrite import, replacing a declared-monotonic transform with 
`hour(:value)` while omitting this field retains the old `True`: the importer 
discards `schema.load()`'s output and updates only supplied keys, so this safe 
default never reaches storage. Could the omitted declaration become false 
before persistence, rather than enabling range mirrors that can exclude valid 
rows?



-- 
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]

Reply via email to