Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
Show all changes
16 commits
Select commit Hold shift + click to select a range
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
fix: lazy-init K8s client for spark_application engine
Defer kubeconfig load until materialize/cleanup so feast apply can
construct the engine without a cluster.

Signed-off-by: Aniket Paluskar <apaluska@redhat.com>
  • Loading branch information
aniketpalu committed Jul 16, 2026
commit cef6fa63289d0396b7598add2b74c5d0f3332a7c
Original file line number Diff line number Diff line change
Expand Up @@ -91,12 +91,33 @@ def __init__(
"snowflake.registry, etc."
Comment thread
aniketpalu marked this conversation as resolved.
)

k8s_config.load_config()
self.k8s_client = client.ApiClient()
self.core_v1 = client.CoreV1Api(self.k8s_client)
self.custom_api = client.CustomObjectsApi(self.k8s_client)
# Defer kubeconfig load until materialize/cleanup — feast apply only
# constructs the engine and calls update() (a no-op), so it must not
# require a cluster.
self._k8s_client = None
self._core_v1 = None
self._custom_api = None
self._server_id = uuid.uuid4().hex[:8]

def _ensure_k8s(self) -> None:
"""Load kubeconfig and create API clients on first K8s use."""
if self._custom_api is not None:
return
k8s_config.load_config()
self._k8s_client = client.ApiClient()
self._core_v1 = client.CoreV1Api(self._k8s_client)
self._custom_api = client.CustomObjectsApi(self._k8s_client)

@property
def core_v1(self):
self._ensure_k8s()
return self._core_v1

@property
def custom_api(self):
self._ensure_k8s()
return self._custom_api

@property
def supports_batch(self) -> bool:
return True
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -55,15 +55,22 @@ def _make_repo_config(
@patch("feast.infra.compute_engines.spark_application.compute.k8s_config")
@patch("feast.infra.compute_engines.spark_application.compute.client")
def _make_engine(mock_client, mock_k8s_config, **kwargs):
"""Create engine with mocked K8s client."""
"""Create engine with mocked K8s client.

Pre-seed API clients so later property access does not call real
load_config() after this helper's patches have exited.
"""
from feast.infra.compute_engines.spark_application.compute import (
SparkApplicationComputeEngine,
)

repo_config = _make_repo_config(**kwargs)
return SparkApplicationComputeEngine(
engine = SparkApplicationComputeEngine(
repo_config=repo_config, offline_store=None, online_store=None
)
engine._core_v1 = MagicMock()
engine._custom_api = MagicMock()
return engine


# ── Test 1: Config defaults + required field ──
Expand Down Expand Up @@ -124,6 +131,27 @@ def test_accepts_snowflake_registry():
assert engine is not None


@patch("feast.infra.compute_engines.spark_application.compute.k8s_config")
@patch("feast.infra.compute_engines.spark_application.compute.client")
def test_init_does_not_load_kubeconfig(mock_client, mock_k8s_config):
"""feast apply constructs the engine but never materializes — no kubeconfig needed."""
from feast.infra.compute_engines.spark_application.compute import (
SparkApplicationComputeEngine,
)

engine = SparkApplicationComputeEngine(
repo_config=_make_repo_config(), offline_store=None, online_store=None
)
mock_k8s_config.load_config.assert_not_called()

_ = engine.core_v1
mock_k8s_config.load_config.assert_called_once()
assert engine.custom_api is not None
# Second access must not reload
_ = engine.custom_api
mock_k8s_config.load_config.assert_called_once()


# ── Test 5: _build_driver_repo_config — one rewrite (batch_engine only) ──


Expand Down
Loading