Skip to content

Commit 839c7a4

Browse files
committed
chore: Created Aggregation interface
Signed-off-by: ntkathole <[email protected]>
1 parent ecee4f6 commit 839c7a4

8 files changed

Lines changed: 18 additions & 9 deletions

File tree

sdk/python/feast/__init__.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@
99
from feast.infra.offline_stores.redshift_source import RedshiftSource
1010
from feast.infra.offline_stores.snowflake_source import SnowflakeSource
1111

12+
from .aggregation import Aggregation
1213
from .batch_feature_view import BatchFeatureView
1314
from .data_source import KafkaSource, KinesisSource, PushSource, RequestSource
1415
from .dataframe import DataFrameEngine, FeastDataFrame
@@ -32,6 +33,7 @@
3233
pass
3334

3435
__all__ = [
36+
"Aggregation",
3537
"BatchFeatureView",
3638
"DataFrameEngine",
3739
"Entity",
Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,7 @@
1+
"""
2+
Aggregation module for Feast.
3+
"""
4+
15
from datetime import timedelta
26
from typing import Optional
37

@@ -91,3 +95,6 @@ def __eq__(self, other):
9195
return False
9296

9397
return True
98+
99+
100+
__all__ = ["Aggregation"]

sdk/python/feast/infra/tiling/__init__.py renamed to sdk/python/feast/aggregation/tiling/__init__.py

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -11,9 +11,9 @@
1111
4. Engine nodes: Convert back to engine format (e.g., from_pandas(), createDataFrame())
1212
"""
1313

14-
from feast.infra.tiling.base import IRMetadata, get_ir_metadata_for_aggregation
15-
from feast.infra.tiling.orchestrator import apply_sawtooth_window_tiling
16-
from feast.infra.tiling.tile_subtraction import (
14+
from feast.aggregation.tiling.base import IRMetadata, get_ir_metadata_for_aggregation
15+
from feast.aggregation.tiling.orchestrator import apply_sawtooth_window_tiling
16+
from feast.aggregation.tiling.tile_subtraction import (
1717
convert_cumulative_to_windowed,
1818
deduplicate_keep_latest,
1919
)
File renamed without changes.

sdk/python/feast/infra/tiling/orchestrator.py renamed to sdk/python/feast/aggregation/tiling/orchestrator.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@
1212
import pandas as pd
1313

1414
from feast.aggregation import Aggregation
15-
from feast.infra.tiling.base import get_ir_metadata_for_aggregation
15+
from feast.aggregation.tiling.base import get_ir_metadata_for_aggregation
1616

1717

1818
def apply_sawtooth_window_tiling(

sdk/python/feast/infra/tiling/tile_subtraction.py renamed to sdk/python/feast/aggregation/tiling/tile_subtraction.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,7 @@
1111
import pandas as pd
1212

1313
from feast.aggregation import Aggregation
14-
from feast.infra.tiling.base import get_ir_metadata_for_aggregation
14+
from feast.aggregation.tiling.base import get_ir_metadata_for_aggregation
1515

1616

1717
def convert_cumulative_to_windowed(

sdk/python/feast/infra/compute_engines/ray/nodes.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -360,8 +360,8 @@ def _execute_tiled_aggregation(self, dataset: Dataset) -> DAGValue:
360360
3. Convert to windowed aggregations
361361
4. Convert pandas → Ray Dataset
362362
"""
363-
from feast.infra.tiling.orchestrator import apply_sawtooth_window_tiling
364-
from feast.infra.tiling.tile_subtraction import (
363+
from feast.aggregation.tiling.orchestrator import apply_sawtooth_window_tiling
364+
from feast.aggregation.tiling.tile_subtraction import (
365365
convert_cumulative_to_windowed,
366366
deduplicate_keep_latest,
367367
)

sdk/python/feast/infra/compute_engines/spark/nodes.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -146,8 +146,8 @@ def _execute_tiled_aggregation(self, input_df: DataFrame) -> DAGValue:
146146
)
147147
aggs_by_window[agg.time_window].append(agg)
148148

149-
from feast.infra.tiling.orchestrator import apply_sawtooth_window_tiling
150-
from feast.infra.tiling.tile_subtraction import (
149+
from feast.aggregation.tiling.orchestrator import apply_sawtooth_window_tiling
150+
from feast.aggregation.tiling.tile_subtraction import (
151151
convert_cumulative_to_windowed,
152152
deduplicate_keep_latest,
153153
)

0 commit comments

Comments
 (0)