diff --git a/python/ray/data/_internal/datasource/iceberg_datasource.py b/python/ray/data/_internal/datasource/iceberg_datasource.py index 41f8afe6d031..68a19de346a4 100644 --- a/python/ray/data/_internal/datasource/iceberg_datasource.py +++ b/python/ray/data/_internal/datasource/iceberg_datasource.py @@ -11,15 +11,17 @@ import pyarrow as pa from packaging import version +from ray.data._internal.arrow_block import _BATCH_SIZE_PRESERVING_STUB_COL_NAME from ray.data._internal.planner.plan_expression.expression_visitors import _ExprVisitor from ray.data._internal.util import _check_import -from ray.data.block import Block, BlockMetadata +from ray.data.block import Block, BlockAccessor, BlockMetadata from ray.data.datasource.datasource import Datasource, ReadTask from ray.data.expressions import ( AliasExpr, BinaryExpr, ColumnExpr, DownloadExpr, + Expr, LiteralExpr, MonotonicallyIncreasingIdExpr, Operation, @@ -84,6 +86,44 @@ logger = logging.getLogger(__name__) +# Field ID for the fabricated stub field, see ``_get_empty_projection_schema``. +# Iceberg matches columns by ID, so this must not be a real column's: colliding +# with a non-boolean column raises ``ResolveError``, and with a boolean column +# reads that column's data instead (the row count is still right, but the read +# is no longer free). Catalogs assign IDs sequentially from 1, so a value this +# large will not collide; Iceberg reserves 2147483546 and above for metadata +# columns, so stay below that. +_STUB_FIELD_ID = 2147483000 + + +def _get_empty_projection_schema() -> "Schema": + """Return a projected schema that reads no data but keeps the row count. + + PyIceberg rebuilds its output to match the projected schema, and a table + rebuilt from zero fields reports zero rows however many the scan matched -- + so an empty projection, which ``Dataset.count()`` legitimately requests, + would silently count zero. Project a single optional field that no data file + has instead: Iceberg fills a missing optional field with nulls, the same way + it reads a file written before a column was added, so the row count arrives + in a column that costs nothing to read. + + The field is named after Ray's stub column so the resulting block matches + what every other reader produces for an empty projection. ``boolean`` is + used because Iceberg gained a null-valued type only after 0.9.0; + ``_get_read_task`` retypes the column to ``null`` on the way out. + """ + from pyiceberg.schema import Schema + from pyiceberg.types import BooleanType, NestedField + + return Schema( + NestedField( + field_id=_STUB_FIELD_ID, + name=_BATCH_SIZE_PRESERVING_STUB_COL_NAME, + field_type=BooleanType(), + required=False, + ) + ) + class _IcebergExpressionVisitor( _ExprVisitor["BooleanExpression | UnboundTerm | Literal"] @@ -210,6 +250,7 @@ def _get_read_task( case_sensitive: bool, limit: Optional[int], schema: "Schema", + empty_projection: bool = False, ) -> Iterable[Block]: # Determine the PyIceberg version to handle backward compatibility import pyiceberg @@ -255,7 +296,14 @@ def _generate_tables() -> Iterable[pa.Table]: ) yield table - yield from _generate_tables() + for table in _generate_tables(): + if empty_projection: + # The scan read a boolean-typed stub, see + # ``_get_empty_projection_schema``. Re-derive it through Ray's own + # empty projection so the column is ``null``-typed, matching the + # schema we report and every other reader's stub. + table = BlockAccessor.for_block(table).select([]) + yield table @DeveloperAPI @@ -297,8 +345,13 @@ def __init__( _check_import(self, module="pyiceberg", package="pyiceberg") from pyiceberg.expressions import AlwaysTrue - self._scan_kwargs = scan_kwargs if scan_kwargs is not None else {} - self._catalog_kwargs = catalog_kwargs if catalog_kwargs is not None else {} + # Copy both dicts: we pop ``name`` from the catalog kwargs and write + # ``snapshot_id`` into the scan kwargs, and doing that in place would + # break a second read that reuses the caller's dict. + self._scan_kwargs = dict(scan_kwargs) if scan_kwargs is not None else {} + self._catalog_kwargs = ( + dict(catalog_kwargs) if catalog_kwargs is not None else {} + ) if "name" in self._catalog_kwargs: self._catalog_name = self._catalog_kwargs.pop("name") @@ -344,10 +397,39 @@ def plan_files(self) -> List["FileScanTask"]: # Calculate and cache the plan_files if they don't already exist if self._plan_files is None: data_scan = self._get_data_scan() - self._plan_files = data_scan.plan_files() + # Annotated ``Iterable[FileScanTask]``; a list today, but this cache + # is iterated more than once, so don't depend on that. + self._plan_files = list(data_scan.plan_files()) return self._plan_files + def apply_predicate(self, predicate_expr: Expr) -> "IcebergDatasource": + """Push a predicate down, discarding plan files planned without it. + + The base implementation shallow-copies, so the clone would inherit a + cache planned before this predicate existed: all-``AlwaysTrue`` + residuals, which ``get_read_tasks`` turns into an unfiltered row count + for a filtered read. + """ + return self._invalidate_plan_files(super().apply_predicate(predicate_expr)) + + def apply_projection( + self, projection_map: Optional[Dict[str, str]] + ) -> "IcebergDatasource": + """Push a projection down, discarding plan files planned without it. + + The projected columns are part of the scan, so a cache built for a + different projection does not belong to the clone. See ``apply_predicate``. + """ + return self._invalidate_plan_files(super().apply_projection(projection_map)) + + def _invalidate_plan_files(self, clone: "IcebergDatasource") -> "IcebergDatasource": + """Drop ``clone``'s inherited plan-file cache, so it re-plans its scan.""" + # ``self`` signals "no pushdown applied": no clone, cache still valid. + if clone is not self: + clone._plan_files = None + return clone + def _get_combined_filter(self) -> "BooleanExpression": """Get the combined filter including both row_filter and pushed-down predicates.""" combined_filter = self._row_filter @@ -367,13 +449,22 @@ def _get_combined_filter(self) -> "BooleanExpression": return combined_filter + def _is_empty_projection(self) -> bool: + """Whether the pushed-down projection selects no columns at all.""" + return self._get_data_columns() == [] + def _get_data_scan(self) -> "DataScan": # Get the combined filter combined_filter = self._get_combined_filter() - # Convert back to tuple for PyIceberg API (None -> ("*",)) + # Convert back to tuple for PyIceberg API (None -> ("*",)). An empty + # projection also scans everything: it selects its own columns through + # the projected schema instead, see ``_get_empty_projection_schema``. data_columns = self._get_data_columns() - selected_fields = ("*",) if data_columns is None else tuple(data_columns) + if not data_columns: + selected_fields = ("*",) + else: + selected_fields = tuple(data_columns) data_scan = self.table.scan( row_filter=combined_filter, @@ -433,6 +524,7 @@ def get_read_tasks( per_task_row_limit: Optional[int] = None, data_context: Optional["DataContext"] = None, ) -> List[ReadTask]: + from pyiceberg.expressions import AlwaysTrue from pyiceberg.io import pyarrow as pyi_pa_io from pyiceberg.manifest import DataFileContent @@ -447,10 +539,21 @@ def get_read_tasks( # Get the arrow schema, to set in the metadata pya_schema = pyi_pa_io.schema_to_pyarrow(projected_schema) + # An empty projection reads a fabricated stub column instead of none at + # all, see ``_get_empty_projection_schema``. Declare the placeholder so + # the reported schema matches the blocks; it is hidden from the + # user-visible schema. + empty_projection = self._is_empty_projection() + if empty_projection: + projected_schema = _get_empty_projection_schema() + pya_schema = pa.schema( + [pa.field(_BATCH_SIZE_PRESERVING_STUB_COL_NAME, pa.null())] + ) + # Set the n_chunks to the min of the number of plan files and the actual # requested n_chunks, so that there are no empty tasks - if parallelism > len(list(plan_files)): - parallelism = len(list(plan_files)) + if parallelism > len(plan_files): + parallelism = len(plan_files) logger.warning( f"Reducing the parallelism to {parallelism}, as that is the number of files" ) @@ -467,6 +570,26 @@ def get_read_tasks( case_sensitive = self._scan_kwargs.get("case_sensitive", True) limit = self._scan_kwargs.get("limit") + # Manifests record how many rows each file holds, which is only still an + # exact count if every surviving row matches the filter. PyIceberg puts + # whatever the filter leaves after file pruning in + # ``FileScanTask.residual``, so ``AlwaysTrue`` means nothing is left to + # check and the counts hold. Same test PyIceberg's own + # ``DataScan.count()`` uses. A ``limit`` stops the read early, so no + # count taken from the manifests describes what comes back. + # + # The subtraction below only accounts for position deletes, but an + # equality delete cannot reach here: ``plan_files`` rejects them in both + # of its paths, locally in ``_plan_files_local`` and over REST in + # ``FileScanTask.from_rest_response``. + counts_are_exact = limit is None and ( + isinstance(row_filter, AlwaysTrue) + or all( + isinstance(getattr(task, "residual", None), AlwaysTrue) + for task in plan_files + ) + ) + get_read_task = partial( _get_read_task, table_io=table_io, @@ -475,6 +598,7 @@ def get_read_tasks( case_sensitive=case_sensitive, limit=limit, schema=projected_schema, + empty_projection=empty_projection, ) read_tasks = [] @@ -495,9 +619,16 @@ def get_read_tasks( for delete in unique_deletes if delete.content == DataFileContent.POSITION_DELETES ) + # ``None`` means "unknown", which makes ``Dataset.count()`` do the + # read instead of trusting these manifest counts. + num_rows = None + if counts_are_exact: + num_rows = ( + sum(task.file.record_count for task in chunk_tasks) + - position_delete_count + ) metadata = BlockMetadata( - num_rows=sum(task.file.record_count for task in chunk_tasks) - - position_delete_count, + num_rows=num_rows, size_bytes=sum(task.file.file_size_in_bytes for task in chunk_tasks), input_files=[task.file.file_path for task in chunk_tasks], exec_stats=None, diff --git a/python/ray/data/tests/datasource/test_iceberg.py b/python/ray/data/tests/datasource/test_iceberg.py index a7473c71262b..cd86a5ef9596 100644 --- a/python/ray/data/tests/datasource/test_iceberg.py +++ b/python/ray/data/tests/datasource/test_iceberg.py @@ -21,6 +21,7 @@ import ray from ray.data import read_iceberg +from ray.data._internal.arrow_block import _BATCH_SIZE_PRESERVING_STUB_COL_NAME from ray.data._internal.datasource.iceberg_datasource import IcebergDatasource from ray.data._internal.logical.operators import Filter, Project from ray.data._internal.logical.optimizers import LogicalOptimizer @@ -223,6 +224,332 @@ def test_filtered_read(): assert all(len(rt.metadata.input_files) == 1 for rt in read_tasks) +def _iceberg_scan_row_count(**scan_kwargs) -> int: + """Rows PyIceberg itself returns for a scan -- the ground truth for counts.""" + sql_catalog = pyi_catalog.load_catalog(**_CATALOG_KWARGS.copy()) + table = sql_catalog.load_table(f"{_DB_NAME}.{_TABLE_NAME}") + return table.scan(**scan_kwargs).to_arrow().num_rows + + +# ``count()`` has two strategies -- read the count from plan metadata, or project +# to zero columns and count what comes back. These cases cover both, including +# the ones that defeat the metadata short-circuit. +_COUNT_CASES = { + "plain": lambda: read_iceberg( + table_identifier=f"{_DB_NAME}.{_TABLE_NAME}", + catalog_kwargs=_CATALOG_KWARGS.copy(), + ), + "select_columns": lambda: read_iceberg( + table_identifier=f"{_DB_NAME}.{_TABLE_NAME}", + catalog_kwargs=_CATALOG_KWARGS.copy(), + ).select_columns(["col_b"]), + # An expression filter is pushed into the datasource, which removes the + # ``Filter`` from the plan and leaves the zero-column projection sitting + # directly on the read. + "expr_filter": lambda: read_iceberg( + table_identifier=f"{_DB_NAME}.{_TABLE_NAME}", + catalog_kwargs=_CATALOG_KWARGS.copy(), + ).filter(expr=col("col_a") < 10), + # Same query as a Python UDF: the predicate cannot be pushed down, so the + # ``Filter`` stays in the plan. This is the control case -- it must keep + # working, so a fix cannot simply disable predicate pushdown. + "udf_filter": lambda: read_iceberg( + table_identifier=f"{_DB_NAME}.{_TABLE_NAME}", + catalog_kwargs=_CATALOG_KWARGS.copy(), + ).filter(lambda row: row["col_a"] < 10), + "expr_filter_then_select": lambda: read_iceberg( + table_identifier=f"{_DB_NAME}.{_TABLE_NAME}", + catalog_kwargs=_CATALOG_KWARGS.copy(), + ) + .filter(expr=col("col_a") < 10) + .select_columns(["col_b"]), + # A filter supplied at read time leaves no ``Filter`` in the plan at all, so + # the count is answered from the Iceberg manifest. + "read_time_row_filter": lambda: read_iceberg( + table_identifier=f"{_DB_NAME}.{_TABLE_NAME}", + row_filter=pyi_expr.LessThan("col_a", 10), + catalog_kwargs=_CATALOG_KWARGS.copy(), + ), + "select_columns_materialized": lambda: read_iceberg( + table_identifier=f"{_DB_NAME}.{_TABLE_NAME}", + catalog_kwargs=_CATALOG_KWARGS.copy(), + ) + .select_columns(["col_b"]) + .materialize(), +} + + +@pytest.mark.skipif( + get_pyarrow_version() < parse_version("14.0.0"), + reason="PyIceberg 0.7.0 fails on pyarrow <= 14.0.0", +) +@pytest.mark.parametrize("case", list(_COUNT_CASES), ids=list(_COUNT_CASES)) +def test_count_matches_rows_actually_produced(case): + """``count()`` must agree with the number of rows the query yields. + + The count is a property of the query, not of how much of it was pushed into + the reader. + """ + make_ds = _COUNT_CASES[case] + expected = len(make_ds().take_all()) + assert make_ds().count() == expected + + +@pytest.mark.skipif( + get_pyarrow_version() < parse_version("14.0.0"), + reason="PyIceberg 0.7.0 fails on pyarrow <= 14.0.0", +) +def test_empty_projection_preserves_row_count(): + """Selecting zero columns must yield N rows of no columns, not zero rows. + + ``Dataset.count()`` projects to zero columns to avoid reading column data, so + an empty projection that drops rows corrupts every count built that way. + """ + iceberg_ds = IcebergDatasource( + table_identifier=f"{_DB_NAME}.{_TABLE_NAME}", + selected_fields=(), + catalog_kwargs=_CATALOG_KWARGS.copy(), + ) + expected = _iceberg_scan_row_count() + assert expected > 0, "fixture table should not be empty" + + rows_read = sum( + block.num_rows + for read_task in iceberg_ds.get_read_tasks(1) + for block in read_task() + ) + assert rows_read == expected + + +@pytest.mark.skipif( + get_pyarrow_version() < parse_version("14.0.0"), + reason="PyIceberg 0.7.0 fails on pyarrow <= 14.0.0", +) +def test_empty_projection_reads_only_the_stub_column(): + """An empty projection must read no real column, only the stub. + + It projects a field present in no schema version, relying on PyIceberg + null-filling a missing optional field. Should an upgrade stop doing that, + the row count collapses to zero silently, so pin the resulting schema: one + ``null``-typed stub column, matching what every other reader produces. + """ + iceberg_ds = IcebergDatasource( + table_identifier=f"{_DB_NAME}.{_TABLE_NAME}", + selected_fields=(), + catalog_kwargs=_CATALOG_KWARGS.copy(), + ) + blocks = [ + block for read_task in iceberg_ds.get_read_tasks(1) for block in read_task() + ] + assert blocks, "fixture table should produce at least one block" + expected_schema = pa.schema( + [pa.field(_BATCH_SIZE_PRESERVING_STUB_COL_NAME, pa.null())] + ) + for block in blocks: + assert block.schema.equals(expected_schema), block.schema + + +@pytest.mark.skipif( + get_pyarrow_version() < parse_version("14.0.0"), + reason="PyIceberg 0.7.0 fails on pyarrow <= 14.0.0", +) +@pytest.mark.parametrize( + ("row_filter", "count_must_be_exact"), + [ + (None, True), + # Prunes whole files, so the manifest row counts stay exact and must + # keep being reported -- this is the free ``count()`` we do not want to + # regress while fixing the case below. + (pyi_expr.In("col_c", {1, 2}), True), + # Selective *within* a file: file-level pruning cannot resolve it, so a + # manifest row count would overcount and must be reported as unknown. + (pyi_expr.LessThan("col_a", 10), False), + ], + ids=["no_filter", "partition_filter", "row_level_filter"], +) +def test_reported_num_rows_matches_rows_read(row_filter, count_must_be_exact): + """Read-task metadata must never claim a row count the read does not deliver. + + ``Dataset.count()`` returns this number directly when nothing in the plan can + change the row count, so an overcount reaches the user as the answer. ``None`` + is always safe but gives up a free count, hence ``count_must_be_exact`` + pinning the cases where the shortcut is sound. + """ + kwargs = {} if row_filter is None else {"row_filter": row_filter} + iceberg_ds = IcebergDatasource( + table_identifier=f"{_DB_NAME}.{_TABLE_NAME}", + catalog_kwargs=_CATALOG_KWARGS.copy(), + **kwargs, + ) + read_tasks = iceberg_ds.get_read_tasks(2) + + claimed = [read_task.metadata.num_rows for read_task in read_tasks] + actual = sum(block.num_rows for read_task in read_tasks for block in read_task()) + + if count_must_be_exact: + assert all(count is not None for count in claimed), ( + "row counts are exact for this filter and must still be reported, " + "otherwise count() pays for a read it does not need" + ) + if all(count is not None for count in claimed): + assert sum(claimed) == actual + + +@pytest.mark.skipif( + get_pyarrow_version() < parse_version("14.0.0"), + reason="PyIceberg 0.7.0 fails on pyarrow <= 14.0.0", +) +def test_reported_num_rows_is_unknown_when_scan_stops_early(): + """A ``limit`` makes the read stop early, so no manifest count describes it. + + The manifests still say how many rows each file holds and the filter may be + fully resolved, so the exactness test would otherwise pass and report the + full total for a read that returns ``limit`` rows. + """ + iceberg_ds = IcebergDatasource( + table_identifier=f"{_DB_NAME}.{_TABLE_NAME}", + catalog_kwargs=_CATALOG_KWARGS.copy(), + # No row filter, so residuals are all ``AlwaysTrue``: the limit is the + # only reason the count cannot be trusted. + scan_kwargs={"limit": 5}, + ) + read_tasks = iceberg_ds.get_read_tasks(1) + + claimed = [read_task.metadata.num_rows for read_task in read_tasks] + actual = sum(block.num_rows for read_task in read_tasks for block in read_task()) + + assert actual == 5, f"limit should cap a single task at 5 rows, got {actual}" + assert all( + count is None for count in claimed + ), f"row counts cannot be exact under a limit, got {claimed}" + + +@pytest.mark.skipif( + get_pyarrow_version() < parse_version("14.0.0"), + reason="PyIceberg 0.7.0 fails on pyarrow <= 14.0.0", +) +def test_empty_projection_survives_schema_evolution_on_pinned_snapshot(): + """An empty projection must not depend on the table's current schema. + + Reading a pinned snapshot after a column was added is where naming a real + stand-in column breaks: picked from the *current* schema it is absent from + the snapshot, PyIceberg raises ``ValueError: Could not find column``, and + ``count()`` fails instead of counting. The fabricated stub belongs to no + schema version, so no schema change can reach it. + """ + sql_catalog = pyi_catalog.load_catalog(**_CATALOG_KWARGS.copy()) + table = sql_catalog.load_table(f"{_DB_NAME}.{_TABLE_NAME}") + old_snapshot_id = table.current_snapshot().snapshot_id + expected_rows = 101 # the fixture appends 120 rows, then deletes col_a >= 101 + + with table.update_schema() as update: + update.add_column("col_d", pyi_types.BooleanType()) + table = sql_catalog.load_table(f"{_DB_NAME}.{_TABLE_NAME}") + assert "col_d" in table.schema().column_names, "schema evolution should apply" + + iceberg_ds = IcebergDatasource( + table_identifier=f"{_DB_NAME}.{_TABLE_NAME}", + catalog_kwargs=_CATALOG_KWARGS.copy(), + snapshot_id=old_snapshot_id, + ) + read_tasks = iceberg_ds.apply_projection({}).get_read_tasks(2) + actual = sum(block.num_rows for read_task in read_tasks for block in read_task()) + assert actual == expected_rows + + +@pytest.mark.skipif( + get_pyarrow_version() < parse_version("14.0.0"), + reason="PyIceberg 0.7.0 fails on pyarrow <= 14.0.0", +) +@pytest.mark.parametrize( + "push_down", + [ + lambda ds: ds.apply_predicate(col("col_a") < 10), + lambda ds: ds.apply_predicate(col("col_a") < 10).apply_projection( + {"col_b": "col_b"} + ), + ], + ids=["predicate", "predicate_then_projection"], +) +def test_pushdown_does_not_inherit_stale_plan_files(push_down): + """A pushdown clone must re-plan its scan, not inherit the cache. + + ``apply_predicate`` and ``apply_projection`` shallow-copy the datasource, so a + cache populated beforehand -- by ``estimate_inmemory_data_size``, say, which + Ray calls to autodetect parallelism -- is shared with the clone. Those tasks + were planned without the predicate, so every residual is ``AlwaysTrue`` and + ``get_read_tasks`` would report the unfiltered manifest total. + + Reading the cache twice also has to keep working: ``plan_files`` is annotated + ``Iterable``, so a future PyIceberg returning a generator would otherwise + leave the second read empty -- zero read tasks, empty dataset. + """ + iceberg_ds = IcebergDatasource( + table_identifier=f"{_DB_NAME}.{_TABLE_NAME}", + catalog_kwargs=_CATALOG_KWARGS.copy(), + ) + + # Warm the cache before pushdown, and read it twice. + assert iceberg_ds.estimate_inmemory_data_size() > 0 + assert iceberg_ds.estimate_inmemory_data_size() > 0, "cache must be re-readable" + unfiltered_files = len(iceberg_ds.plan_files) + assert unfiltered_files > 0, "fixture should plan at least one file" + + filtered_ds = push_down(iceberg_ds) + read_tasks = filtered_ds.get_read_tasks(2) + assert read_tasks, "pushdown must not leave the clone with an exhausted cache" + + claimed = [read_task.metadata.num_rows for read_task in read_tasks] + actual = sum(block.num_rows for read_task in read_tasks for block in read_task()) + + assert actual == 10, f"filter should match 10 of the fixture's rows, got {actual}" + if all(count is not None for count in claimed): + assert ( + sum(claimed) == actual + ), "reported row count came from plan files that predate the predicate" + + +@pytest.mark.skipif( + get_pyarrow_version() < parse_version("14.0.0"), + reason="PyIceberg 0.7.0 fails on pyarrow <= 14.0.0", +) +def test_read_iceberg_does_not_mutate_caller_kwargs(): + """The caller's dicts belong to the caller. + + Both were stored by reference. ``catalog_kwargs`` is then ``pop("name")``-ed, + so a second read of the same table raises ``NoSuchTableError``. + ``scan_kwargs`` gets ``snapshot_id`` written into it, which fails silently + instead: a later read reusing the dict inherits the pin and returns stale rows. + """ + catalog_kwargs = _CATALOG_KWARGS.copy() + scan_kwargs = {} + expected_catalog_kwargs = catalog_kwargs.copy() + + # An explicit ``snapshot_id`` is required to reach the ``scan_kwargs`` write: + # the buggy line is guarded by ``if snapshot_id``, so passing ``None`` would + # leave the dict untouched whether or not the bug is present. + sql_catalog = pyi_catalog.load_catalog(**_CATALOG_KWARGS.copy()) + table = sql_catalog.load_table(f"{_DB_NAME}.{_TABLE_NAME}") + snapshot_id = table.current_snapshot().snapshot_id + + first = read_iceberg( + table_identifier=f"{_DB_NAME}.{_TABLE_NAME}", + catalog_kwargs=catalog_kwargs, + scan_kwargs=scan_kwargs, + snapshot_id=snapshot_id, + ) + assert catalog_kwargs == expected_catalog_kwargs + assert scan_kwargs == {} + + # Reusing the same dicts must work exactly like the first read. + second = read_iceberg( + table_identifier=f"{_DB_NAME}.{_TABLE_NAME}", + catalog_kwargs=catalog_kwargs, + scan_kwargs=scan_kwargs, + ) + assert second.count() == first.count() + + @pytest.mark.skipif( get_pyarrow_version() < parse_version("14.0.0"), reason="PyIceberg 0.7.0 fails on pyarrow <= 14.0.0",