Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Prev Previous commit
Next Next commit
chore: Created Aggregation interface
Signed-off-by: ntkathole <[email protected]>
  • Loading branch information
ntkathole committed Nov 21, 2025
commit 712d5f92a419805c0e4501745b0443bf162ae910
2 changes: 2 additions & 0 deletions sdk/python/feast/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
from feast.infra.offline_stores.redshift_source import RedshiftSource
from feast.infra.offline_stores.snowflake_source import SnowflakeSource

from .aggregation import Aggregation
from .batch_feature_view import BatchFeatureView
from .data_source import KafkaSource, KinesisSource, PushSource, RequestSource
from .dataframe import DataFrameEngine, FeastDataFrame
Expand All @@ -32,6 +33,7 @@
pass

__all__ = [
"Aggregation",
"BatchFeatureView",
"DataFrameEngine",
"Entity",
Expand Down
Original file line number Diff line number Diff line change
@@ -1,3 +1,7 @@
"""
Aggregation module for Feast.
"""

from datetime import timedelta
from typing import Optional

Expand Down Expand Up @@ -91,3 +95,6 @@ def __eq__(self, other):
return False

return True


__all__ = ["Aggregation"]
Original file line number Diff line number Diff line change
Expand Up @@ -11,9 +11,9 @@
4. Engine nodes: Convert back to engine format (e.g., from_pandas(), createDataFrame())
"""

from feast.infra.tiling.base import IRMetadata, get_ir_metadata_for_aggregation
from feast.infra.tiling.orchestrator import apply_sawtooth_window_tiling
from feast.infra.tiling.tile_subtraction import (
from feast.aggregation.tiling.base import IRMetadata, get_ir_metadata_for_aggregation
from feast.aggregation.tiling.orchestrator import apply_sawtooth_window_tiling
from feast.aggregation.tiling.tile_subtraction import (
convert_cumulative_to_windowed,
deduplicate_keep_latest,
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@
import pandas as pd

from feast.aggregation import Aggregation
from feast.infra.tiling.base import get_ir_metadata_for_aggregation
from feast.aggregation.tiling.base import get_ir_metadata_for_aggregation


def apply_sawtooth_window_tiling(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@
import pandas as pd

from feast.aggregation import Aggregation
from feast.infra.tiling.base import get_ir_metadata_for_aggregation
from feast.aggregation.tiling.base import get_ir_metadata_for_aggregation


def convert_cumulative_to_windowed(
Expand Down
4 changes: 2 additions & 2 deletions sdk/python/feast/infra/compute_engines/ray/nodes.py
Original file line number Diff line number Diff line change
Expand Up @@ -360,8 +360,8 @@ def _execute_tiled_aggregation(self, dataset: Dataset) -> DAGValue:
3. Convert to windowed aggregations
4. Convert pandas → Ray Dataset
"""
from feast.infra.tiling.orchestrator import apply_sawtooth_window_tiling
from feast.infra.tiling.tile_subtraction import (
from feast.aggregation.tiling.orchestrator import apply_sawtooth_window_tiling
from feast.aggregation.tiling.tile_subtraction import (
convert_cumulative_to_windowed,
deduplicate_keep_latest,
)
Expand Down
4 changes: 2 additions & 2 deletions sdk/python/feast/infra/compute_engines/spark/nodes.py
Original file line number Diff line number Diff line change
Expand Up @@ -146,8 +146,8 @@ def _execute_tiled_aggregation(self, input_df: DataFrame) -> DAGValue:
)
aggs_by_window[agg.time_window].append(agg)

from feast.infra.tiling.orchestrator import apply_sawtooth_window_tiling
from feast.infra.tiling.tile_subtraction import (
from feast.aggregation.tiling.orchestrator import apply_sawtooth_window_tiling
from feast.aggregation.tiling.tile_subtraction import (
convert_cumulative_to_windowed,
deduplicate_keep_latest,
)
Expand Down
Loading