Skip to content
Merged
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
47 changes: 24 additions & 23 deletions sqlmesh/core/test/runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -125,25 +125,7 @@ def run_tests(
# Ensure workers are not greater than the number of tests
num_workers = min(len(model_test_metadata) or 1, default_test_connection.concurrent_tasks)

def _run_single_test(
metadata: ModelTestMetadata, engine_adapter: EngineAdapter
) -> t.Optional[ModelTextTestResult]:
test = ModelTest.create_test(
body=metadata.body,
test_name=metadata.test_name,
models=models,
engine_adapter=engine_adapter,
dialect=dialect,
path=metadata.path,
default_catalog=default_catalog,
preserve_fixtures=preserve_fixtures,
concurrency=num_workers > 1,
verbosity=verbosity,
)

if not test:
return None

def _run_single_test(test: ModelTest) -> ModelTextTestResult:
result = t.cast(
ModelTextTestResult,
ModelTextTestRunner().run(t.cast(unittest.TestCase, test)),
Expand All @@ -158,11 +140,30 @@ def _run_single_test(

start_time = time.perf_counter()
try:
Comment thread
cmgoffena13 marked this conversation as resolved.
# Build ModelTest instances on the calling thread before workers start. create_test()
# can call to_datetime() / ttl_cache (time.time()), which races with another worker's
# time_machine freeze when execution_time is set under concurrent_tasks > 1.
# NOTE: We can run create_tests in a separate parallel stage for a future optimization.
# We just can't overlap runs/creations.
tests: list[ModelTest] = []
for metadata, engine_adapter in metadata_to_adapter.items():
test = ModelTest.create_test(
body=metadata.body,
test_name=metadata.test_name,
models=models,
engine_adapter=engine_adapter,
dialect=dialect,
path=metadata.path,
default_catalog=default_catalog,
preserve_fixtures=preserve_fixtures,
concurrency=num_workers > 1,
verbosity=verbosity,
)
if test:
tests.append(test)

with ThreadPoolExecutor(max_workers=num_workers) as pool:
futures = [
pool.submit(_run_single_test, metadata=metadata, engine_adapter=engine_adapter)
for metadata, engine_adapter in metadata_to_adapter.items()
]
futures = [pool.submit(_run_single_test, test) for test in tests]

for future in concurrent.futures.as_completed(futures):
test_results.append(future.result())
Expand Down
Loading