Skip to content
Open
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
31 changes: 23 additions & 8 deletions pyterrier/_evaluation/_grid.py
Original file line number Diff line number Diff line change
Expand Up @@ -248,11 +248,9 @@ def _evaluate_several_settings(inputs : List[Tuple]):
else:
import itertools
import more_itertools
try:
from pyterrier_alpha.parallel import parallel_lambda # type: ignore
except ImportError as ie:
raise ImportError("pyterrier-alpha[parallel] must be installed for jobs>1") from ie

from concurrent.futures import ThreadPoolExecutor
from jnius import detach

all_inputs = [(keys, values) for values in combinations]

# how many jobs to distribute this to
Expand All @@ -261,9 +259,26 @@ def _evaluate_several_settings(inputs : List[Tuple]):
# built the batches to distribute
batched_inputs = list(more_itertools.chunked(all_inputs, num_batches))
assert len(batched_inputs) > 0, "No inputs identified for parallel_lambda"
eval_list = parallel_lambda(_evaluate_several_settings, batched_inputs, jobs, backend=backend)
eval_list = list(itertools.chain(*eval_list))
assert len(eval_list) > 0, "parallel_lambda returned 0 rows"

if backend == 'ray':
# preserve alpha behavior since ray has its own JVM lifecycle
try:
from pyterrier_alpha.parallel import parallel_lambda # type: ignore
except ImportError as ie:
raise ImportError("pyterrier-alpha[parallel] must be installed for backend='ray'") from ie
eval_list = parallel_lambda(_evaluate_several_settings, batched_inputs, jobs, backend=backend)
eval_list = list(itertools.chain(*eval_list))
else:
# avoid process parallelism by using threads, which share the parent's JVM
# instead of calling fork() and leaving leave the JVM in a weird state
def _eval_chunk(chunk):
out = [_evaluate_one_setting(k, v) for k, v in chunk]
detach() # release JNI refs

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Flake says we should be decorating this method with @pt.java.required, but I think the actual point is that if the transformer doesnt involve Java (Terrier or Anserini) then jnius may not even be installed. So we need to detect jnuis (try import etc) and act accordingly.

return out
with ThreadPoolExecutor(max_workers=jobs) as ex:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Have you tested that you get a speed up with different Python threads calling Java? I think Terrier's data structures arent thread safe unless they are loaded with "concurrent:" prefix. At least with the forked JVM, if setup correctly, the results would be correct.

per_chunk = list(ex.map(_eval_chunk, batched_inputs))
eval_list = list(itertools.chain(*per_chunk))
assert len(eval_list) > 0, "GridScan produced 0 rows"

# resulting eval_list has the form [
# ( [(BR, 'wmodel', 'BM25'), (BR, 'c', 0.2)] , {"map" : 0.2654} )
Expand Down
Loading