Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
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
4 changes: 2 additions & 2 deletions .github/workflows/acceptance.yml
Original file line number Diff line number Diff line change
Expand Up @@ -58,8 +58,8 @@ jobs:
test:
name: Testing (${{ matrix.component }})
runs-on:
group: databrickslabs-protected-runner-group
labels: linux-ubuntu-latest
group: larger-runners
labels: larger
needs: [ not-a-fork, lint ]
permissions:
id-token: write
Expand Down
12 changes: 10 additions & 2 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,13 @@ UV_RUN := uv run --exact --all-extras
# single component's tests in parallel, e.g. `make test TEST_PATH=tests/impulse_query_engine`.
TEST_PATH ?= tests/

# Extra args passed to pytest, primarily xdist parallelism. Each worker starts its own
# Spark JVM (see the worker-isolated `spark` fixture in tests/conftest.py), so `-n auto`
# is capped to bound memory on many-core machines. `worksteal` lets idle workers pull
# queued tests off busy ones, shortening the slow-test tail. Override with
# `make test PYTEST_XARGS=-n0` to run serially in-process for debugging.
PYTEST_XARGS ?= -n auto --maxprocesses=4 --dist worksteal

clean:
rm -fr .venv htmlcov .pytest_cache .ruff_cache .coverage coverage.xml test-results.xml
find . -name '__pycache__' -print0 | xargs -0 rm -fr
Expand All @@ -32,7 +39,7 @@ fmt:
$(UV_RUN) ruff check src/ tests/ --fix

test:
$(UV_RUN) pytest $(TEST_PATH) --cov=src --cov-branch --cov-report=xml
$(UV_RUN) pytest $(TEST_PATH) $(PYTEST_XARGS) --cov=src --cov-branch --cov-report=xml

coverage:
$(UV_RUN) pytest tests/ --cov=src --cov-branch --cov-report=html
Expand All @@ -49,7 +56,8 @@ lock-dependencies:
uv lock
printf 'setuptools>=61.0\nwheel\n' | uv pip compile --generate-hashes --universal --no-header --quiet - > .build-constraints.txt
@perl -pi -e 's|registry = "https://[^"]*"|registry = "https://pypi.org/simple"|g' uv.lock
@printf 'Stripped registry references from uv.lock.\n'
@perl -pi -e 's|https://pypi-proxy\.dev\.databricks\.com/|https://files.pythonhosted.org/|g' uv.lock
@printf 'Stripped registry and proxy URLs from uv.lock.\n'

# Mirror a fork PR onto a fork-test/pr-<N> branch in the main repo and open a test PR,
# so CI (which is skipped for fork PRs) runs with JFrog/OIDC. Review the fork code first.
Expand Down
1 change: 1 addition & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@ dev = [
"pytest-benchmark==5.1.0",
"pytest-cov==6.0.0",
"pytest-mock==3.15.1",
"pytest-xdist==3.6.1",
"ruff==0.9.10",
"setuptools==75.8.2",
"black==26.3.1"
Expand Down
39 changes: 34 additions & 5 deletions tests/conftest.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import fcntl
import os
from unittest.mock import create_autospec

Expand All @@ -14,17 +15,45 @@


@pytest.fixture(scope="session")
def spark() -> SparkSession:
spark = configure_spark_with_delta_pip(
def spark(tmp_path_factory, worker_id) -> SparkSession:
# Isolate the warehouse + Derby metastore per pytest-xdist worker. Each worker is a
# separate process with its own SparkSession/JVM, and a Derby-backed metastore cannot
# be shared across processes. ``getbasetemp()`` is already per-worker under xdist, and
# ``worker_id`` is "master" in serial runs, so this is a safe no-op without ``-n``.
base = tmp_path_factory.getbasetemp()
warehouse_dir = base / "spark-warehouse"
metastore_dir = base / "metastore_db"
builder = configure_spark_with_delta_pip(
SparkSession.builder.master("local")
.appName(f"impulse-tests-{worker_id}")
.config("spark.sql.warehouse.dir", str(warehouse_dir))
.config(
"spark.hadoop.javax.jdo.option.ConnectionURL",
f"jdbc:derby:;databaseName={metastore_dir};create=true",
)
.config(
"spark.sql.catalog.spark_catalog",
"org.apache.spark.sql.delta.catalog.DeltaCatalog",
)
.config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")
.config("spark.databricks.delta.retentionDurationCheck.enabled ", "false")
.config("spark.shuffle.partitions", 1)
).getOrCreate()
.config("spark.databricks.delta.retentionDurationCheck.enabled", "false")
# Local tuning for tiny test data: avoid 200-way shuffles and skip the Spark UI
# so per-worker sessions start fast and stay lean.
.config("spark.sql.shuffle.partitions", 1)
.config("spark.default.parallelism", 1)
.config("spark.ui.enabled", "false")
)
# configure_spark_with_delta_pip resolves the Delta jars via ivy at JVM launch. Under
# xdist, workers sharing the default ~/.ivy2 race on a cold cache and some fail with
# JAVA_GATEWAY_EXITED / unresolved dependency. Serialize session startup across workers
# with a cross-run file lock (base.parent is shared by all workers of a run): the first
# worker populates the shared cache, the rest reuse it. Cheap once the cache is warm.
with open(base.parent / "impulse-spark-startup.lock", "w") as lock_fh:
fcntl.flock(lock_fh, fcntl.LOCK_EX)
try:
spark = builder.getOrCreate()
finally:
fcntl.flock(lock_fh, fcntl.LOCK_UN)
spark.sql("CREATE SCHEMA IF NOT EXISTS spark_catalog.silver")
spark.sql("CREATE SCHEMA IF NOT EXISTS spark_catalog.silver_narrow_db")
spark.sql("CREATE SCHEMA IF NOT EXISTS spark_catalog.silver_key_value_store")
Expand Down
24 changes: 24 additions & 0 deletions uv.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading