Configure Kubernetes executor step pods with the shared Dagster home and persistent volume used by the code location. Changelog: fixed
81 lines
2.7 KiB
Python
81 lines
2.7 KiB
Python
"""Verifies that execution-target linkage is what the guide claims it is."""
|
|
|
|
import pytest
|
|
|
|
from distributed_execution.preflight import check_env_vars, describe_pod_identity
|
|
from distributed_execution.repository import defs
|
|
from distributed_execution.tightly_coupled.jobs import (
|
|
STEP_K8S_CONFIG,
|
|
tightly_coupled_in_process_job,
|
|
tightly_coupled_k8s_job,
|
|
tightly_coupled_local_job,
|
|
)
|
|
|
|
EXPECTED_JOBS = {
|
|
"tightly_coupled_in_process_job",
|
|
"tightly_coupled_local_job",
|
|
"tightly_coupled_k8s_job",
|
|
}
|
|
|
|
|
|
def test_repository_registers_tightly_coupled_jobs():
|
|
names = {job.name for job in defs.jobs}
|
|
assert EXPECTED_JOBS <= names
|
|
|
|
|
|
@pytest.mark.parametrize(
|
|
("job", "expected_executor"),
|
|
[
|
|
(tightly_coupled_in_process_job, "in_process"),
|
|
(tightly_coupled_local_job, "multiprocess"),
|
|
(tightly_coupled_k8s_job, "k8s_job"),
|
|
],
|
|
)
|
|
def test_each_job_declares_its_execution_target(job, expected_executor):
|
|
assert job.tags["execution_target"] == "tightly_coupled"
|
|
assert job.tags["executor"] == expected_executor
|
|
|
|
|
|
def test_k8s_step_pods_mount_shared_dagster_storage():
|
|
container_config = STEP_K8S_CONFIG["container_config"]
|
|
pod_spec_config = STEP_K8S_CONFIG["pod_spec_config"]
|
|
|
|
assert {"name": "DAGSTER_HOME", "value": "/dagster/shared/distributed-execution"} in container_config["env"]
|
|
assert {"name": "dagster-shared-storage", "mount_path": "/dagster/shared"} in container_config["volume_mounts"]
|
|
assert {
|
|
"name": "dagster-shared-storage",
|
|
"persistent_volume_claim": {"claim_name": "dagster-shared-pvc"},
|
|
} in pod_spec_config["volumes"]
|
|
|
|
|
|
def test_in_process_job_runs_end_to_end():
|
|
result = tightly_coupled_in_process_job.execute_in_process()
|
|
|
|
assert result.success
|
|
summary = result.output_for_node("summarise_results")
|
|
assert summary["units"] == 4
|
|
assert summary["total"] == 14 # 0 + 1 + 4 + 9
|
|
# In-process execution keeps every step in the run worker's own process.
|
|
assert len(summary["contributing_workers"]) == 1
|
|
|
|
|
|
def test_work_units_fan_out_into_one_step_each():
|
|
result = tightly_coupled_in_process_job.execute_in_process()
|
|
|
|
mapped = result.output_for_node("process_work_unit")
|
|
# Without this the executor choice would be meaningless: one step cannot span pods.
|
|
assert set(mapped) == {"unit_0", "unit_1", "unit_2", "unit_3"}
|
|
|
|
|
|
def test_env_var_check_reports_missing_names():
|
|
result = check_env_vars(("DEFINITELY_NOT_SET_12345",))
|
|
|
|
assert result["passed"] is False
|
|
assert result["missing"] == ["DEFINITELY_NOT_SET_12345"]
|
|
|
|
|
|
def test_pod_identity_falls_back_outside_kubernetes(monkeypatch):
|
|
monkeypatch.delenv("DAGSTER_K8S_PIPELINE_RUN_NAMESPACE", raising=False)
|
|
|
|
assert describe_pod_identity()["namespace"] == "<not-in-kubernetes>"
|