Skip to content

Commit 11ac31a

Browse files
fix(test): create ModelTests before concurrent pool workers (#6040)
Signed-off-by: devtechedge <devtechedge@users.noreply.github.com> Signed-off-by: Dev M <devtechedge@gmail.com> Co-authored-by: devtechedge <devtechedge@users.noreply.github.com> Co-authored-by: Cortland Goffena <30168413+cmgoffena13@users.noreply.github.com>
1 parent 847de8c commit 11ac31a

1 file changed

Lines changed: 24 additions & 23 deletions

File tree

‎sqlmesh/core/test/runner.py‎

Lines changed: 24 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -125,25 +125,7 @@ def run_tests(
125125
# Ensure workers are not greater than the number of tests
126126
num_workers = min(len(model_test_metadata) or 1, default_test_connection.concurrent_tasks)
127127

128-
def _run_single_test(
129-
metadata: ModelTestMetadata, engine_adapter: EngineAdapter
130-
) -> t.Optional[ModelTextTestResult]:
131-
test = ModelTest.create_test(
132-
body=metadata.body,
133-
test_name=metadata.test_name,
134-
models=models,
135-
engine_adapter=engine_adapter,
136-
dialect=dialect,
137-
path=metadata.path,
138-
default_catalog=default_catalog,
139-
preserve_fixtures=preserve_fixtures,
140-
concurrency=num_workers > 1,
141-
verbosity=verbosity,
142-
)
143-
144-
if not test:
145-
return None
146-
128+
def _run_single_test(test: ModelTest) -> ModelTextTestResult:
147129
result = t.cast(
148130
ModelTextTestResult,
149131
ModelTextTestRunner().run(t.cast(unittest.TestCase, test)),
@@ -158,11 +140,30 @@ def _run_single_test(
158140

159141
start_time = time.perf_counter()
160142
try:
143+
# Build ModelTest instances on the calling thread before workers start. create_test()
144+
# can call to_datetime() / ttl_cache (time.time()), which races with another worker's
145+
# time_machine freeze when execution_time is set under concurrent_tasks > 1.
146+
# NOTE: We can run create_tests in a separate parallel stage for a future optimization.
147+
# We just can't overlap runs/creations.
148+
tests: list[ModelTest] = []
149+
for metadata, engine_adapter in metadata_to_adapter.items():
150+
test = ModelTest.create_test(
151+
body=metadata.body,
152+
test_name=metadata.test_name,
153+
models=models,
154+
engine_adapter=engine_adapter,
155+
dialect=dialect,
156+
path=metadata.path,
157+
default_catalog=default_catalog,
158+
preserve_fixtures=preserve_fixtures,
159+
concurrency=num_workers > 1,
160+
verbosity=verbosity,
161+
)
162+
if test:
163+
tests.append(test)
164+
161165
with ThreadPoolExecutor(max_workers=num_workers) as pool:
162-
futures = [
163-
pool.submit(_run_single_test, metadata=metadata, engine_adapter=engine_adapter)
164-
for metadata, engine_adapter in metadata_to_adapter.items()
165-
]
166+
futures = [pool.submit(_run_single_test, test) for test in tests]
166167

167168
for future in concurrent.futures.as_completed(futures):
168169
test_results.append(future.result())

0 commit comments

Comments
 (0)