-
Notifications
You must be signed in to change notification settings - Fork 191
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Implement engine="thread" in ChunkRecordingExecutor #3526
base: main
Are you sure you want to change the base?
Implement engine="thread" in ChunkRecordingExecutor #3526
Conversation
After a long fight this is now working. |
@zm711 : if you have time could you make some test using windows ? |
Yep. Will do once tests pass :) |
@alejoe91 @zm711 @chrishalcrow @h-mayorquin : ready for review. Also we should discuss changing |
failing tests are really weird. The arrays that are not equal look exactly the same from the values |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
few tweaks as well.
I haven't look deeply but since Windows and Macs have a different mp context at baseline do we think that could be contributing or that this is just floating point issues?
script in case people want to know import spikeinterface.full as si
sorting, rec = si.generate_ground_truth_recording()
analyzer = si.create_sorting_analyzer(rec, sorting)
analyzer.compute(['random_spikes', 'waveforms'])
%timeit analyzer.compute(['principal_components'])
# or
%timeit analyzer.compute(['principal_components'], n_jobs=2) Okay with default settings on this branch using the
compared to going back to main
I'm thinking this might be deeper than our code. Let me know if you want a different type of test Sam. For real data with n_jobs we often have PCA take 1-2 hours. So I'm not currently testing with real data. Do we want me to play around with n_job number and threads? |
Hi Zach and Alessio. Thanks fo testing. pca is not the best parralelisation we have. Yes, this failing test are super strange because it is a floating point issue which is not the same with multiprocessing or thread. This hard to figure. As a side note : I modified the waveforms_tools that estimate templates with both thread and process and this do not give the same results. I need to set decimals=4 for almost equal which is very easy to pass. I did not except this. Maybe we have somthing deep with parralisation that can change the result. |
I can test |
test are passing now |
I forgot.... easy testing of |
okay detect_peaks on main
with this pr
So an improvement but still worse than not using multiprocessing. |
Now trying a longer simulated recording to see if the more realistic length if we see a big difference. object is GroundTruthRecording: 4 channels - 25.0kHz - 1 segments - 250,000,000 samples
10,000.00s (2.78 hours) - float32 dtype - 3.73 GiB
|
merci Zach. How to document this is a big topic. The parralelism was quite empirical until now. For big linux station, it faster with many channel and many core but do not scale we reach a plateau quite soon in speed. |
@samuelgarcia . i'm testing it and i have a bug
File "/media/cure/Secondary/pierre/softwares/spikeinterface/src/spikeinterface/sorters/internal/spyking_circus2.py", line 190, in _run_from_folder |
It works only if you modify L537 in job_tools.py by recording_slices2 = [(thread_local_data, ) + tuple(args) for args in recording_slices] (need to add the tuple() ) |
Hello, just some quick benchmarking on a M4 Mac, doing THIS PR: CURRENT MAIN: So the "thread" option doesn't seem to be working as expected for this system. |
Did some benchmarks on linux, and while it is slower than process, this is working and we do get a speedup. However, @samuelgarcia, I think you should also extend the possibilities to the get_poolexecutor() in job_tools.py. Some functions in components are relying on this, and maybe we should propagate the thread option there also |
I did more tests on windows on a quite old machine i5-4460 3.2Ghz (4 cores) Running this simple example from spikeinterface.generation import generate_drifting_recording
from spikeinterface.sortingcomponents.peak_detection import detect_peaks
from spikeinterface import get_noise_levels
import time
all_job_kwargs = [
dict(pool_engine="process", n_jobs=2, mp_context="spawn", max_threads_per_worker=2),
dict(pool_engine="process", n_jobs=4, mp_context="spawn", max_threads_per_worker=1),
dict(pool_engine="thread", n_jobs=4, mp_context=None, max_threads_per_worker=1),
dict(pool_engine="thread", n_jobs=2, mp_context=None, max_threads_per_worker=2),
dict(n_jobs=1),
]
rec, _, sorting = generate_drifting_recording(
num_units=50,
duration=120.0,
sampling_frequency=30000.0,
probe_name="Neuropixel-128",
)
# print(rec)
noise_levels = get_noise_levels(rec, return_scaled=False)
for job_kwargs in all_job_kwargs:
print()
print(job_kwargs)
t0 = time.perf_counter()
peaks = detect_peaks(rec, method="locally_exclusive", noise_levels=noise_levels, **job_kwargs)
t1 = time.perf_counter()
print("time included the spawn:", t1-t0) Give this
So for me 4 conclusions:
Note: n_jobs=1 do not limit max_threads_per_worker underlying libs are threading under the wood. @chrishalcrow @zm711 could you rerun the same code on your machine because we clearly do not have the results. If someone have a stronger windows machine I would be happy to change n_jobs=4 to something bigger. |
I'm processing something now, but I can run your exact script on mine with n=8 to see what happens in your condition. For me I found an improvement with increasing channels. So I think you and I agree for Windows. I found it was worse with low channel counts, but better with high channel counts. See this comment where I tested higher channel count :) So that means we need to run the test on mac. I could also do that from mine (which is M1pro) but Chris has a newer he could try with. |
On a bigger machine with 40cores with on Linux. with the same script but n = 10 or 20 or 40 with conter balance the max_threads_per_worker. I have this.
Here the trend is:
Note : the spawn time will be negligable for long recording (here in this bench duration are very very short). |
EDIT: hadn't pulled. So re-ran after pulling. Results from running your script:
Did it with the duration x10'd too:
|
When I do the same process on real data,
Any ideas why this would be different than the generated recording? |
Thanks Chris and Zach. Indeed my benchmark is not realistic finally. |
src/spikeinterface/core/job_tools.py
Outdated
else: # windows and mac | ||
# on windows and macos the fork is forbidden and process+spwan is super slow at startup | ||
# so let's go to threads | ||
pool_engine = "thread" |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Change mac default to "process"
…parralel_with_thread # Conflicts: # src/spikeinterface/core/waveform_tools.py
for more information, see https://pre-commit.ci
More informative progresse bar. darwin default job_kwargs
…nterface into parralel_with_thread
for more information, see https://pre-commit.ci
…parralel_with_thread
…nterface into parralel_with_thread
@alejoe91 : ready to merge I think |
Was there a tested/profiled case where this was useful at the end? |
It is commented in the test. in test_job_tools.py. |
The main use case is to avoid the spawn on windows and also for to use thread in a HPC situation because the fork consumes too much memory and slurm is killing my task... |
@h-mayorquin @zm711 @alejoe91
I wanted to do this since a long time : implement thread for ChunkRecordingExecutor.
Now we can do :
Maybe this will help a lot windows user and will avoid use "spawn" for multiprocessing (which is too slow at startup).
My bet is that thread will have IO (read data from disk) and computation withou the GIL without computing lib, so it should good enought for parralell computing. Lets see.
Here a very first as a proof of concept. (Without any real test).
TODO:
max_threads_per_process
tomax_threads_per_worker
. See note.n_worker
ornum_worker
ormax_worker
(Alessio will be unhappy)get_best_job_kwargs()
that will be platform dependant.When using engine thread we will need to disambigu the concept of worker that can be thread and
max_threads_per_worker
which related to the hidden thread use by computing lib like numpy/scipy/blas/lapcak/sklearn.EDIT:
Important notes:
many
ProcessPoolExecutor
are hard coded in several places (principal_components, split, merge...)we should also do this but this will be a separated PR. because theses parallelisation are not on the recording axis.
Note that tricky stuff between loops, thread and process is initialzing variable to a private dict per worker. This is done in
ChunkRecordingExecutor
but this should be propagated for other use cases that do not need recording slices.