Repository navigation
Newest updates break DaskOfflineStore with S3 parquets #4753
Copy link
Copy link
Closed
Labels
Description
Activity
CC @ntkathole
@bjmccotter7192 Can you please confirm if latest version fixed this issue ?
Had the same issue:
driver_stats_source = FileSource( name="driver_hourly_stats_source", path="s3://feast-data/driver_stats.parquet", timestamp_field="event_timestamp", s3_endpoint_override="https://s3-basket/....ru", created_timestamp_column="created", )File "/opt/app-root/lib64/python3.11/site-packages/pyarrow/dataset.py", line 476, in _filesystem_dataset fs, paths_or_selector = _ensure_single_source(source, filesystem) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/pyarrow/dataset.py", line 441, in _ensure_single_source raise FileNotFoundError(path) FileNotFoundError: /tmp/tmpp4_p5yjf/s3:/feast-data/driver_stats.parquet@makSSMZ
Can you please provide a longer stack trace ?@ntkathole
I think you forgot to fix it also infeast.infra.offline_stores.file_source.FileSource.get_uri_for_file_path()
That might be the issue.......:54260 - "POST /push HTTP/1.1" 500 File "/opt/app-root/lib64/python3.11/site-packages/pyarrow/parquet/core.py", line 1371, in __init__ self._dataset = ds.dataset(path_or_paths, filesystem=filesystem, ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/pyarrow/dataset.py", line 794, in dataset return _filesystem_dataset(source, **kwargs) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/pyarrow/dataset.py", line 476, in _filesystem_dataset fs, paths_or_selector = _ensure_single_source(source, filesystem) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/pyarrow/dataset.py", line 441, in _ensure_single_source raise FileNotFoundError(path) FileNotFoundError: /tmp/tmphakk65gr/s3:/feast-data/driver_stats.parquet Traceback (most recent call last): File "/opt/app-root/lib64/python3.11/site-packages/starlette/middleware/errors.py", line 165, in __call__ await self.app(scope, receive, _send) File "/opt/app-root/lib64/python3.11/site-packages/starlette/middleware/exceptions.py", line 62, in __call__ await wrap_app_handling_exceptions(self.app, conn)(scope, receive, send) File "/opt/app-root/lib64/python3.11/site-packages/starlette/_exception_handler.py", line 53, in wrapped_app raise exc File "/opt/app-root/lib64/python3.11/site-packages/starlette/_exception_handler.py", line 42, in wrapped_app await app(scope, receive, sender) File "/opt/app-root/lib64/python3.11/site-packages/starlette/routing.py", line 714, in __call__ await self.middleware_stack(scope, receive, send) File "/opt/app-root/lib64/python3.11/site-packages/starlette/routing.py", line 734, in app await route.handle(scope, receive, send) File "/opt/app-root/lib64/python3.11/site-packages/starlette/routing.py", line 288, in handle await self.app(scope, receive, send) File "/opt/app-root/lib64/python3.11/site-packages/starlette/routing.py", line 76, in app await wrap_app_handling_exceptions(app, request)(scope, receive, send) File "/opt/app-root/lib64/python3.11/site-packages/starlette/_exception_handler.py", line 53, in wrapped_app raise exc File "/opt/app-root/lib64/python3.11/site-packages/starlette/_exception_handler.py", line 42, in wrapped_app await app(scope, receive, sender) File "/opt/app-root/lib64/python3.11/site-packages/starlette/routing.py", line 73, in app response = await f(request) ^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/fastapi/routing.py", line 301, in app raw_response = await run_endpoint_function( ^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/fastapi/routing.py", line 212, in run_endpoint_function return await dependant.call(**values) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/feast/feature_server.py", line 315, in push store.push(**push_params) File "/opt/app-root/lib64/python3.11/site-packages/feast/feature_store.py", line 1485, in push self.write_to_offline_store( File "/opt/app-root/lib64/python3.11/site-packages/feast/feature_store.py", line 1705, in write_to_offline_store provider.ingest_df_to_offline_store(feature_view, table) File "/opt/app-root/lib64/python3.11/site-packages/feast/infra/passthrough_provider.py", line 418, in ingest_df_to_offline_store self.offline_write_batch(self.repo_config, feature_view, table, None) File "/opt/app-root/lib64/python3.11/site-packages/feast/infra/passthrough_provider.py", line 219, in offline_write_batch self.offline_store.__class__.offline_write_batch( File "/opt/app-root/lib64/python3.11/site-packages/feast/infra/offline_stores/dask.py", line 480, in offline_write_batch prev_table = pyarrow.parquet.read_table( ^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/pyarrow/parquet/core.py", line 1793, in read_table dataset = ParquetDataset( ^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/pyarrow/parquet/core.py", line 1371, in __init__ self._dataset = ds.dataset(path_or_paths, filesystem=filesystem, ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/pyarrow/dataset.py", line 794, in dataset return _filesystem_dataset(source, **kwargs) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/pyarrow/dataset.py", line 476, in _filesystem_dataset fs, paths_or_selector = _ensure_single_source(source, filesystem) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/pyarrow/dataset.py", line 441, in _ensure_single_source raise FileNotFoundError(path) FileNotFoundError: /tmp/tmphakk65gr/s3:/feast-data/driver_stats.parquet [2025-04-03 08:51:56 +0000] [23] [ERROR] Exception in ASGI application Traceback (most recent call last): File "/opt/app-root/lib64/python3.11/site-packages/uvicorn/protocols/http/httptools_impl.py", line 409, in run_asgi result = await app( # type: ignore[func-returns-value] ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/uvicorn/middleware/proxy_headers.py", line 60, in __call__ return await self.app(scope, receive, send) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/fastapi/applications.py", line 1054, in __call__ await super().__call__(scope, receive, send) File "/opt/app-root/lib64/python3.11/site-packages/starlette/applications.py", line 112, in __call__ await self.middleware_stack(scope, receive, send) File "/opt/app-root/lib64/python3.11/site-packages/starlette/middleware/errors.py", line 187, in __call__ raise exc File "/opt/app-root/lib64/python3.11/site-packages/starlette/middleware/errors.py", line 165, in __call__ await self.app(scope, receive, _send) File "/opt/app-root/lib64/python3.11/site-packages/starlette/middleware/exceptions.py", line 62, in __call__ await wrap_app_handling_exceptions(self.app, conn)(scope, receive, send) File "/opt/app-root/lib64/python3.11/site-packages/starlette/_exception_handler.py", line 53, in wrapped_app raise exc File "/opt/app-root/lib64/python3.11/site-packages/starlette/_exception_handler.py", line 42, in wrapped_app await app(scope, receive, sender) File "/opt/app-root/lib64/python3.11/site-packages/starlette/routing.py", line 714, in __call__ await self.middleware_stack(scope, receive, send) File "/opt/app-root/lib64/python3.11/site-packages/starlette/routing.py", line 734, in app await route.handle(scope, receive, send) File "/opt/app-root/lib64/python3.11/site-packages/starlette/routing.py", line 288, in handle await self.app(scope, receive, send) File "/opt/app-root/lib64/python3.11/site-packages/starlette/routing.py", line 76, in app await wrap_app_handling_exceptions(app, request)(scope, receive, send) File "/opt/app-root/lib64/python3.11/site-packages/starlette/_exception_handler.py", line 53, in wrapped_app raise exc File "/opt/app-root/lib64/python3.11/site-packages/starlette/_exception_handler.py", line 42, in wrapped_app await app(scope, receive, sender) File "/opt/app-root/lib64/python3.11/site-packages/starlette/routing.py", line 73, in app response = await f(request) ^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/fastapi/routing.py", line 301, in app raw_response = await run_endpoint_function( ^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/fastapi/routing.py", line 212, in run_endpoint_function return await dependant.call(**values) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/feast/feature_server.py", line 315, in push store.push(**push_params) File "/opt/app-root/lib64/python3.11/site-packages/feast/feature_store.py", line 1485, in push self.write_to_offline_store( File "/opt/app-root/lib64/python3.11/site-packages/feast/feature_store.py", line 1705, in write_to_offline_store provider.ingest_df_to_offline_store(feature_view, table) File "/opt/app-root/lib64/python3.11/site-packages/feast/infra/passthrough_provider.py", line 418, in ingest_df_to_offline_store self.offline_write_batch(self.repo_config, feature_view, table, None) File "/opt/app-root/lib64/python3.11/site-packages/feast/infra/passthrough_provider.py", line 219, in offline_write_batch self.offline_store.__class__.offline_write_batch( File "/opt/app-root/lib64/python3.11/site-packages/feast/infra/offline_stores/dask.py", line 480, in offline_write_batch prev_table = pyarrow.parquet.read_table( ^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/pyarrow/parquet/core.py", line 1793, in read_table dataset = ParquetDataset( ^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/pyarrow/parquet/core.py", line 1371, in __init__ self._dataset = ds.dataset(path_or_paths, filesystem=filesystem, ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/pyarrow/dataset.py", line 794, in dataset return _filesystem_dataset(source, **kwargs) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/pyarrow/dataset.py", line 476, in _filesystem_dataset fs, paths_or_selector = _ensure_single_source(source, filesystem) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/opt/app-root/lib64/python3.11/site-packages/pyarrow/dataset.py", line 441, in _ensure_single_source raise FileNotFoundError(path) FileNotFoundError: /tmp/tmphakk65gr/s3:/feast-data/driver_stats.parquet@tmihalac
We also tried to change the line locally just to test in the filelib/python3.11/site-packages/feast/infra/offline_stores/dask.pyif config.repo_path is not None and not Path(file_options.uri).is_absolute(): absolute_path = config.repo_path / file_options.uri else: absolute_path = Path(file_options.uri)to
if config.repo_path is not None and not Path(file_options.uri).is_absolute(): absolute_path = file_options.uri else: absolute_path = Path(file_options.uri)And the push worked. Maybe the problem is somewhere in that part.
#5208 will fix it
Thank you all for looking and working on this issue. I will follow both this issue and the #5208 for the final fix and then I'll update my environment and make sure everything is working as expected.
Expected Behavior
In version 0.40.1 the Dask Offline store was able to read the data_source.path directly from the FileSource and retrieve the data from S3 using a path like:
s3://<your-bucket>/<file-name>Current Behavior
Failing to pull data because it is now appending the repo_path to the front of the s3 url.
Example:
/tmp/feast:s3//<your-bucket>/<file-name>I believe this is because of a recent change: #4624 which is now not accepting the S3 url as a
absolute PathSteps to reproduce
0.41.3get_historical_featuresand call hung for a while then errored with the file path error not existingSpecifications
Possible Solution