| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
There was a problem hiding this comment.
On-demand feature outputs remain unvalidated, and mixed failure types can report the wrong first offending row.
Review effort: Balanced
Findings: 2
Adds schema-driven vector-length validation to online writes and Arrow-based materialization paths.
Changes:
| File | Description |
|---|---|
| sdk/python/feast/utils.py | Adds Arrow vector-length validation. |
| sdk/python/feast/feature_store.py | Vectorizes online-write validation. |
| sdk/python/tests/unit/test_vector_length_validation.py | Tests both validation implementations. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
Sorry, something went wrong.
| not_a_sequence = lengths.isna() | ||
| if not_a_sequence.any(): | ||
| i = not_a_sequence.idxmax() | ||
| raise ValueError( | ||
| f"Row {i}: Vector feature '{name}' is not a sequence. Got: {type(column[i])}" | ||
| ) | ||
|
|
||
| mismatched = lengths != expected | ||
| if mismatched.any(): | ||
| i = mismatched.idxmax() | ||
| raise ValueError( | ||
| f"Row {i}: Vector length {lengths[i]} does not match expected {expected} " | ||
| f"for feature '{name}' in feature view '{feature_view.name}'." | ||
| ) |
| f"Feature view '{feature_view.name}' has no batch_source and cannot be converted to proto." | ||
| ) | ||
|
|
||
| _validate_vector_field_lengths(table, feature_view) |
|
⚠️ Please install the Codecov Report✅ All modified and coverable lines are covered by tests. @@ Coverage Diff @@
## master #6909 +/- ##
==========================================
+ Coverage 48.45% 48.50% +0.04%
==========================================
Files 427 427
Lines 53718 53755 +37
Branches 7822 7827 +5
==========================================
+ Hits 26031 26075 +44
+ Misses 25819 25813 -6
+ Partials 1868 1867 -1
Continue to review full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
Sorry, something went wrong.
|
@ntkathole mind take a look |
Sorry, something went wrong.
| # reports None so the two failure modes stay distinguishable. | ||
| lengths = column.map(lambda v: len(v) if hasattr(v, "__len__") else None) | ||
|
|
||
| not_a_sequence = lengths.isna() |
There was a problem hiding this comment.
lengths.isna() will be True for both non-sequence values (mapped to None) and actual NaN/None values in the original column.
May be filter out rows where column itself is null/NaN ?
Sorry, something went wrong.
There was a problem hiding this comment.
Just an edge case, overall Looks good
Sorry, something went wrong.
Three findings from review of feast-dev#6909. Tolerate genuine nulls on the DataFrame path (@ntkathole). lengths.isna() was true both for values that were not sequences and for rows that were genuinely null, so a null vector was reported as "not a sequence". Nulls are now skipped, which also makes this path consistent with the Arrow path, which already tolerated null rows. Report the first offending row across both failure modes (Copilot). Checking the not-a-sequence mask fully before the length mask meant a non-sequence in a later row was reported ahead of a wrong length in an earlier one. The masks are now unioned and the first offending position selected before choosing the message, which restores the row ordering the original iterrows loop had. Value lookup is positional so a DataFrame with a duplicate index reports a scalar type rather than a Series. Validate the on-demand branch too (Copilot). _convert_arrow_to_proto dispatches to _convert_arrow_odfv_to_proto or _convert_arrow_fv_to_proto, and the check sat in the regular branch only, so a stored ODFV declaring vector_length had its transformed output unchecked. The call moved up into the dispatcher, ahead of the branch, so both paths are covered. Each finding has a test that fails against the previous commit and passes here. Signed-off-by: hao-xu5 <hxu44@apple.com>
|
Thanks both — all three were real. Fixed in b4d932a. @ntkathole, on lengths.isna() — you were right, and it also made the two code paths disagree: the Arrow validator already tolerated null rows while the DataFrame one did not. Nulls are now skipped on both: is_null = column.isna()
lengths = column.map(lambda v: len(v) if hasattr(v, "__len__") else None, na_action="ignore")
not_a_sequence = lengths.isna() & ~is_nullna_action="ignore" leaves nulls as NaN, so a NaN length that is not null means the value genuinely was not a sequence. A null vector is no longer reported as "not a sequence". On first offending row — correct, and worth noting this was a regression I introduced rather than pre-existing: the original iterrows() loop got row ordering right by construction, and splitting into two sequential mask checks lost it. Now unioned, with the position chosen before the message: offending = not_a_sequence | mismatched
position = int(offending.to_numpy().argmax())
label = column.index[position]Value lookup is column.iloc[position], so a duplicate index reports a scalar type rather than a Series — thanks for catching that second half, there's a test with index=[7, 7, 7]. On the ODFV branch — also correct. The check sat in _convert_arrow_fv_to_proto, so a stored ODFV declaring vector_length had its transformed output unvalidated. Moved up into _convert_arrow_to_proto ahead of the dispatch, so both branches are covered. I confirmed the ODFV table reaching the dispatcher is post-transformation, and that the validator no-ops when the output column is absent. Each finding has a test that fails against the previous commit and passes here — verified by splicing the old validator back in: test_null_vectors_are_tolerated_on_the_dataframe_path FAILED -> passes test_first_offending_row_wins_across_both_failure_modes FAILED -> passes test_duplicate_index_reports_a_scalar_type_not_a_series FAILED -> passes test_on_demand_feature_view_output_is_validated FAILED -> passes 21 tests in the file now. ruff clean, mypy unchanged (same 2 pre-existing errors in feature_store.py:2123), 134 passed across the surrounding unit suites. Two other things worth recording: Perf, since this adds work to the materialize path. Measured against the cost of the conversion it sits inside:
fixed_size_list is O(1) — it compares type.list_size and never touches the data, hence flat across batch sizes. Variable-length lists are ~4 vectorized pyarrow.compute passes. Feature views that do not set vector_length pay 2.5 us for the metadata lookup and return. The DataFrame path got ~60x faster by dropping iterrows() (1.508s to 0.025s on 200k rows). One known gap, unchanged: the snowflake compute engine builds its rows without going through _convert_arrow_to_proto, so it is not covered here. Happy to do it in a follow-up if you'd like — it needs a different seam. The two items I left in #6907 are still open questions rather than fixes, so not in this PR: requiring vector_length when vector_index=True (today 0 is the default and silently means "skip"), and whether Field.__eq__ should compare vector_index / vector_search_metric, which are currently commented out so a COSINE -> L2 change produces no diff in feast plan. |
Sorry, something went wrong.
…rized Addresses three of the five items in feast-dev#6907. Validation was reachable only from the online-write path. _validate_vector_features had a single call site inside _get_feature_view_and_df_for_online_write, so materialize, materialize_incremental and get_historical_features never ran it. Embedding workloads are overwhelmingly batch, so the path carrying vectors at volume was the unvalidated one. Every compute engine (local, spark, ray, flink, kubernetes, aws_lambda) and the passthrough provider funnel through utils._convert_arrow_to_proto, so a single Arrow-level check placed in _convert_arrow_fv_to_proto covers all of them. _validate_vector_field_lengths is O(1) for fixed-size lists and one vectorized pass for variable-size lists, so it is cheap enough to leave on. Null rows are tolerated; a non-list column with a declared vector_length is an error. Note the snowflake compute engine builds its rows without going through _convert_arrow_to_proto, so it is not covered by this change. The check assumed the vector was the first feature. feature_view.features[0].vector_index meant that declaring any field before the vector field silently disabled validation. Both validators now resolve the field via _get_feature_view_vector_field_metadata, which scans the schema and is already used elsewhere in feature_store.py. iterrows() did not scale. Replaced with a single pass over the column, which measured ~60x faster on 200k rows (1.51s to 0.025s) while keeping the row index and the two distinct failure messages. Not included, since feast-dev#6907 raises them as behaviour decisions rather than fixes: requiring vector_length when vector_index=True (item 3), and whether Field.__eq__ should compare vector_index and vector_search_metric (item 5). Adds sdk/python/tests/unit/test_vector_length_validation.py with 16 tests covering both validators, including regression tests for the vector field not being first. Signed-off-by: hao-xu5 <hxu44@apple.com>
Three findings from review of feast-dev#6909. Tolerate genuine nulls on the DataFrame path (@ntkathole). lengths.isna() was true both for values that were not sequences and for rows that were genuinely null, so a null vector was reported as "not a sequence". Nulls are now skipped, which also makes this path consistent with the Arrow path, which already tolerated null rows. Report the first offending row across both failure modes (Copilot). Checking the not-a-sequence mask fully before the length mask meant a non-sequence in a later row was reported ahead of a wrong length in an earlier one. The masks are now unioned and the first offending position selected before choosing the message, which restores the row ordering the original iterrows loop had. Value lookup is positional so a DataFrame with a duplicate index reports a scalar type rather than a Series. Validate the on-demand branch too (Copilot). _convert_arrow_to_proto dispatches to _convert_arrow_odfv_to_proto or _convert_arrow_fv_to_proto, and the check sat in the regular branch only, so a stored ODFV declaring vector_length had its transformed output unchecked. The call moved up into the dispatcher, ahead of the branch, so both paths are covered. Each finding has a test that fails against the previous commit and passes here. Signed-off-by: hao-xu5 <hxu44@apple.com>
Adds LanceSource and teaches the DuckDB offline store to read it, completing the read half of feast-dev#6899. LanceFormat landed in feast-dev#6925 as a format descriptor; nothing read Lance until now. Lance already worked through SparkSource, which drives its reader generically from table_format.format_type.value and table_format.properties. What was missing is a path that needs no JVM, which is also the real test of whether the DataSource abstraction is engine-agnostic rather than Spark-agnostic in name only. Placement: a new source read by the existing DuckDB store, rather than a Lance offline store or an extension of FileSource. FileSource is the wrong host. Its format axis is already taken by file_format, so adding table_format would give one source two overlapping format axes. It is also read by two stores with incompatible contracts: duckdb._read_data_source dispatches on type, while dask._read_datasource has no dispatch seam and reads file_options.uri unconditionally as Parquet, and asserts isinstance(..., FileSource) in three places. Decisively, Lance's catalog addressing has no path to put in FileSource.path, so the catalog-based layer would not fit the class even if the path-based one did. A Lance offline store would be the wrong 120 lines. duckdb.py is a binding that injects reader and writer callbacks into the engine in ibis.py, so a Lance store would be a near-copy of it plus a repo_config entry, and would force a choice between Lance and Parquet instead of mixing them in one feature service. Reading a source as ibis.memtable(arrow_table) in the DuckDB store already has two precedents, IcebergSource and MlflowDatasetSource. Following them leaves ibis.py untouched, so the point-in-time join, TTL handling, field mapping and ODFVs work unchanged, and no edit to repo_config.py or data_source.py is needed because CUSTOM_SOURCE plus data_source_class_type is self-describing. Both addressing modes work: a uri, and catalog/namespace/table through namespace_client and table_id. Pin semantics follow what was argued on feast-dev#5782 and feast-dev#6925: a pin selects data, never shape. get_table_column_names_and_types reads the pinned schema so feast apply infers what reads will actually see; a pre-flight check fails with a message naming the pin when a pinned version cannot satisfy the declared schema; and vector widths go through _validate_vector_field_lengths from feast-dev#6909 rather than a second validator. Tests use the dir namespace implementation, which exercises the same namespace_client and table_id code path as a remote catalog with no server required. 43 tests, including a demonstration that a tag pin returns earlier data after the dataset has been overwritten for the same entity and timestamp. Read-only for now: _write_data_source is untouched, so a LanceSource is not yet a persist target and there is no SavedDatasetLanceStorage. Signed-off-by: hao-xu5 <hxu44@apple.com>
* feat: Add a non-JVM read path for Lance data sources Adds LanceSource and teaches the DuckDB offline store to read it, completing the read half of #6899. LanceFormat landed in #6925 as a format descriptor; nothing read Lance until now. Lance already worked through SparkSource, which drives its reader generically from table_format.format_type.value and table_format.properties. What was missing is a path that needs no JVM, which is also the real test of whether the DataSource abstraction is engine-agnostic rather than Spark-agnostic in name only. Placement: a new source read by the existing DuckDB store, rather than a Lance offline store or an extension of FileSource. FileSource is the wrong host. Its format axis is already taken by file_format, so adding table_format would give one source two overlapping format axes. It is also read by two stores with incompatible contracts: duckdb._read_data_source dispatches on type, while dask._read_datasource has no dispatch seam and reads file_options.uri unconditionally as Parquet, and asserts isinstance(..., FileSource) in three places. Decisively, Lance's catalog addressing has no path to put in FileSource.path, so the catalog-based layer would not fit the class even if the path-based one did. A Lance offline store would be the wrong 120 lines. duckdb.py is a binding that injects reader and writer callbacks into the engine in ibis.py, so a Lance store would be a near-copy of it plus a repo_config entry, and would force a choice between Lance and Parquet instead of mixing them in one feature service. Reading a source as ibis.memtable(arrow_table) in the DuckDB store already has two precedents, IcebergSource and MlflowDatasetSource. Following them leaves ibis.py untouched, so the point-in-time join, TTL handling, field mapping and ODFVs work unchanged, and no edit to repo_config.py or data_source.py is needed because CUSTOM_SOURCE plus data_source_class_type is self-describing. Both addressing modes work: a uri, and catalog/namespace/table through namespace_client and table_id. Pin semantics follow what was argued on #5782 and #6925: a pin selects data, never shape. get_table_column_names_and_types reads the pinned schema so feast apply infers what reads will actually see; a pre-flight check fails with a message naming the pin when a pinned version cannot satisfy the declared schema; and vector widths go through _validate_vector_field_lengths from #6909 rather than a second validator. Tests use the dir namespace implementation, which exercises the same namespace_client and table_id code path as a remote catalog with no server required. 43 tests, including a demonstration that a tag pin returns earlier data after the dataset has been overwritten for the same entity and timestamp. Read-only for now: _write_data_source is untouched, so a LanceSource is not yet a persist target and there is no SavedDatasetLanceStorage. Signed-off-by: hao-xu5 <hxu44@apple.com> * fix: Address Lance read review feedback Signed-off-by: HaoXuAI <sduxuhao@gmail.com> --------- Signed-off-by: hao-xu5 <hxu44@apple.com> Signed-off-by: HaoXuAI <sduxuhao@gmail.com>
feast-dev#6943 made a `LanceSource` readable without Spark. This makes one writable, through the same two addressing modes: a uri, or a `namespace_client` plus `table_id` resolved through a Lance namespace, so a catalog-addressed dataset is committed through its catalog rather than behind its back. `_write_data_source` gains a `LanceSource` branch, following the `IcebergSource` branch already there, which makes the path reachable from `DuckDBOfflineStore.offline_write_batch`. The `isinstance` narrowing the read path open-coded is now a shared `_as_lance_source` helper used by both, rather than a second copy of the optional-import guard. Three things are settled before Lance is called. A pinned source is refused. A Lance commit always produces a new version, so `write_dataset` has no `version` argument at all; a write through a source pinned to version 1 or to a tag would succeed and then be invisible through the very source that performed it. Measured: after appending to a dataset tagged `prod` at version 1, the unpinned source reads 5 rows and the pinned source still reads 3. An absent dataset is created rather than appended to, because Lance has no create-or-append mode and rejects `create` on an existing dataset. An existing dataset's vector widths are compared against the incoming data. Lance rejects an `append` whose schema disagrees, but it accepts an `overwrite` that replaces a 8-wide embedding column with a 32-wide one, silently rewriting the declared shape and leaving a version history whose vectors are not mutually comparable, with every index built on the old width invalidated. No legitimate schema evolution changes an embedding's dimension, so this is refused. The guard is narrow on purpose: only vector widths are policed, and other type changes remain Lance's business. The declared `vector_length` is checked at the store entry point rather than in the writer callback, which is handed a `DataSource` and so cannot see the feature view. It reuses `_validate_vector_field_lengths` from feast-dev#6909 rather than adding a second implementation. Without it, creating a dataset at a width other than the declared one would succeed and fail only on the next read. 19 tests added, all using the `dir` namespace implementation for the catalog-based cases so no server is needed. Signed-off-by: hao-xu5 <hxu44@apple.com>
feast-dev#6943 made a `LanceSource` readable without Spark. This makes one writable, through the same two addressing modes: a uri, or a `namespace_client` plus `table_id` resolved through a Lance namespace, so a catalog-addressed dataset is committed through its catalog rather than behind its back. `_write_data_source` gains a `LanceSource` branch, following the `IcebergSource` branch already there, which makes the path reachable from `DuckDBOfflineStore.offline_write_batch`. The `isinstance` narrowing the read path open-coded is now a shared `_as_lance_source` helper used by both, rather than a second copy of the optional-import guard. Three things are settled before Lance is called. A pinned source is refused. A Lance commit always produces a new version, so `write_dataset` has no `version` argument at all; a write through a source pinned to version 1 or to a tag would succeed and then be invisible through the very source that performed it. Measured: after appending to a dataset tagged `prod` at version 1, the unpinned source reads 5 rows and the pinned source still reads 3. An absent dataset is created rather than appended to, because Lance has no create-or-append mode and rejects `create` on an existing dataset. An existing dataset's vector widths are compared against the incoming data. Lance rejects an `append` whose schema disagrees, but it accepts an `overwrite` that replaces a 8-wide embedding column with a 32-wide one, silently rewriting the declared shape and leaving a version history whose vectors are not mutually comparable, with every index built on the old width invalidated. No legitimate schema evolution changes an embedding's dimension, so this is refused. The guard is narrow on purpose: only vector widths are policed, and other type changes remain Lance's business. The declared `vector_length` is checked at the store entry point rather than in the writer callback, which is handed a `DataSource` and so cannot see the feature view. It reuses `_validate_vector_field_lengths` from feast-dev#6909 rather than adding a second implementation. Without it, creating a dataset at a width other than the declared one would succeed and fail only on the next read. 19 tests added, all using the `dir` namespace implementation for the catalog-based cases so no server is needed. Signed-off-by: hao-xu5 <hxu44@apple.com>
* feat: Add a non-JVM write path for Lance data sources #6943 made a `LanceSource` readable without Spark. This makes one writable, through the same two addressing modes: a uri, or a `namespace_client` plus `table_id` resolved through a Lance namespace, so a catalog-addressed dataset is committed through its catalog rather than behind its back. `_write_data_source` gains a `LanceSource` branch, following the `IcebergSource` branch already there, which makes the path reachable from `DuckDBOfflineStore.offline_write_batch`. The `isinstance` narrowing the read path open-coded is now a shared `_as_lance_source` helper used by both, rather than a second copy of the optional-import guard. Three things are settled before Lance is called. A pinned source is refused. A Lance commit always produces a new version, so `write_dataset` has no `version` argument at all; a write through a source pinned to version 1 or to a tag would succeed and then be invisible through the very source that performed it. Measured: after appending to a dataset tagged `prod` at version 1, the unpinned source reads 5 rows and the pinned source still reads 3. An absent dataset is created rather than appended to, because Lance has no create-or-append mode and rejects `create` on an existing dataset. An existing dataset's vector widths are compared against the incoming data. Lance rejects an `append` whose schema disagrees, but it accepts an `overwrite` that replaces a 8-wide embedding column with a 32-wide one, silently rewriting the declared shape and leaving a version history whose vectors are not mutually comparable, with every index built on the old width invalidated. No legitimate schema evolution changes an embedding's dimension, so this is refused. The guard is narrow on purpose: only vector widths are policed, and other type changes remain Lance's business. The declared `vector_length` is checked at the store entry point rather than in the writer callback, which is handed a `DataSource` and so cannot see the feature view. It reuses `_validate_vector_field_lengths` from #6909 rather than adding a second implementation. Without it, creating a dataset at a width other than the declared one would succeed and fail only on the next read. 19 tests added, all using the `dir` namespace implementation for the catalog-based cases so no server is needed. Signed-off-by: hao-xu5 <hxu44@apple.com> * fix: Correct what a Lance namespace does during a write The docstring claimed a catalog-addressed write is "committed through its namespace". That is only true of a create. Instrumenting the namespace client shows how far it is actually involved: create -> declare_table append -> describe_table only open -> describe_table, namespace_id So a namespace records that a table exists and where it lives, and does not track its versions: the version an append produces is never reported back to it. Worth stating precisely, because it is the reason the pin contract has to be enforced in `assert_writable` rather than left to the catalog -- no catalog is going to reject a write on a pin's behalf. Adds a test that records the namespace calls, so the claim is pinned by a test rather than asserted in prose. Signed-off-by: hao-xu5 <hxu44@apple.com> * fix: Accept an optional source in _as_lance_source offline_write_batch passes FeatureView.batch_source, which is optional, so the narrow DataSource annotation made mypy reject the call. The body already handled None, since isinstance(None, LanceSource) is False, so only the signature changes. Signed-off-by: hao-xu5 <hxu44@apple.com> --------- Signed-off-by: hao-xu5 <hxu44@apple.com>
| Back | FazBrowse Home | New Git URL |
Fixes three of the five items in #6907.
1. Validation was unreachable from the batch path
_validate_vector_features had a single call site, inside _get_feature_view_and_df_for_online_write. So:
Embedding workloads are overwhelmingly batch, so the path actually carrying vectors at volume was the unvalidated one.
Rather than patch each engine, this uses the one place they all funnel through: utils._convert_arrow_to_proto. The local, spark, ray, flink, kubernetes and aws_lambda engines plus PassthroughProvider all call it, so a single Arrow-level check in _convert_arrow_fv_to_proto covers all of them.
The new _validate_vector_field_lengths is O(1) for fixed_size_list (compares type.list_size) and one vectorized pass for list/large_list, so it is cheap enough to leave on unconditionally. Null rows are tolerated. A non-list column carrying a declared vector_length is an error.
One known gap: the snowflake compute engine builds its rows directly rather than via _convert_arrow_to_proto, so it is not covered. Happy to extend if preferred, but it needs a different seam.
2. The check assumed the vector was the first feature
Declaring any field ahead of the vector field silently disabled validation entirely. Both validators now resolve the field with _get_feature_view_vector_field_metadata(), which scans feature_view.schema, raises on more than one vector field, and is already used in three other places in feature_store.py.
There are regression tests for this on both paths.
3. iterrows() did not scale
Replaced with a single pass over the column. Measured on 200k rows:
The row index and both distinct failure messages are preserved.
Deliberately not included
#6907 items 3 and 5 are behaviour decisions rather than clear-cut fixes, so they are left for discussion there:
Testing
New sdk/python/tests/unit/test_vector_length_validation.py, 16 tests across both validators: fixed-size and variable-size lists, first offending row reporting, vector field not first, unset vector_length, null rows, missing column, non-list type, RecordBatch input, and numpy vectors.
No behaviour change for feature views that do not set vector_length, since 0 still short-circuits.
Verified no regressions against master by running the relevant suites on both sides: