Skip to content
Open
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
15 changes: 10 additions & 5 deletions src/spikeinterface/core/basesorting.py
Original file line number Diff line number Diff line change
Expand Up @@ -1076,11 +1076,16 @@ def _get_spike_vector_segment_slices(self):
if self._cached_spike_vector_segment_slices is None:
# compute the, this is needed when spikevector is loaded from format and not computed
num_seg = self.get_num_segments()
slices = np.searchsorted(self._cached_spike_vector["segment_index"], np.arange(num_seg + 1))
self._cached_spike_vector_segment_slices = np.zeros((num_seg, 2), dtype="int64")
for seg_index in range(num_seg):
self._cached_spike_vector_segment_slices[seg_index, 0] = slices[seg_index]
self._cached_spike_vector_segment_slices[seg_index, 1] = slices[seg_index + 1]
if num_seg == 1:
self._cached_spike_vector_segment_slices = np.array(
[[0, self._cached_spike_vector.size]], dtype="int64"
)
else:
slices = np.searchsorted(self._cached_spike_vector["segment_index"], np.arange(num_seg + 1))
self._cached_spike_vector_segment_slices = np.zeros((num_seg, 2), dtype="int64")
for seg_index in range(num_seg):
self._cached_spike_vector_segment_slices[seg_index, 0] = slices[seg_index]
self._cached_spike_vector_segment_slices[seg_index, 1] = slices[seg_index + 1]
return self._cached_spike_vector_segment_slices

def to_reordered_spike_vector(
Expand Down
114 changes: 87 additions & 27 deletions src/spikeinterface/core/node_pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,11 @@ def __init__(
parents = [parents]
self.parents = parents

self._kwargs = dict()
self._kwargs = dict(
time_series=time_series,
return_output=return_output,
parents=parents,
)

def get_margin(self):
# can optionally be overwritten
Expand Down Expand Up @@ -110,10 +114,15 @@ def __init__(self, recording, peaks):
self.peaks = peaks

# precompute segment slice
self.segment_slices = []
for segment_index in range(recording.get_num_segments()):
i0, i1 = np.searchsorted(peaks["segment_index"], [segment_index, segment_index + 1])
self.segment_slices.append(slice(i0, i1))
if recording.get_num_segments() > 1:
self.segment_slices = []
for segment_index in range(recording.get_num_segments()):
i0, i1 = np.searchsorted(peaks["segment_index"], [segment_index, segment_index + 1])
self.segment_slices.append(slice(i0, i1))
else:
self.segment_slices = None

self._kwargs.update(dict(peaks=peaks))

def get_margin(self):
return 0
Expand All @@ -122,15 +131,22 @@ def get_dtype(self):
return base_peak_dtype

def get_peak_slice(self, segment_index, start_frame, end_frame, max_margin):
sl = self.segment_slices[segment_index]
peaks_in_segment = self.peaks[sl]
if self.segment_slices is not None:
sl = self.segment_slices[segment_index]
peaks_in_segment = self.peaks[sl]
else:
peaks_in_segment = self.peaks
i0, i1 = np.searchsorted(peaks_in_segment["sample_index"], [start_frame, end_frame])
return i0, i1

def compute(self, traces, start_frame, end_frame, segment_index, max_margin, peak_slice):
# get local peaks
sl = self.segment_slices[segment_index]
peaks_in_segment = self.peaks[sl]
if self.segment_slices is not None:
sl = self.segment_slices[segment_index]
peaks_in_segment = self.peaks[sl]
else:
peaks_in_segment = self.peaks

# i0, i1 = np.searchsorted(peaks_in_segment["sample_index"], [start_frame, end_frame])
i0, i1 = peak_slice
local_peaks = peaks_in_segment[i0:i1]
Expand Down Expand Up @@ -195,7 +211,6 @@ def __init__(
category=FutureWarning,
stacklevel=2,
)

self._dtype = spike_peak_dtype

self.include_spikes_in_margin = include_spikes_in_margin
Expand All @@ -204,19 +219,40 @@ def __init__(

main_channel_ids = sorting.get_property("main_channel_id")
assert main_channel_ids is not None, "SpikeRetriever needs the sorting to have `main_channel_id`s."
main_channel_indices = recording.ids_to_indices(main_channel_ids)
self.peaks = sorting_to_peaks(sorting, main_channel_indices, self._dtype)
self.main_channel_indices = recording.ids_to_indices(main_channel_ids)
self.spike_vector, segment_slices = sorting.to_spike_vector(return_slices=True)
self.spike_sample_indices = np.asarray(self.spike_vector["sample_index"])

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.

Can you read self.spike_vector["sample_index"] directly in get_peak_slice instead of storing it. np.searchsorted works on the field view without a copy. Plus, when the node is pickled for spawn workers the view is pickled as its own array so every worker gets the sample indices twice

self.sorting = sorting
self._peaks = None

if not channel_from_template:
channel_distance = get_channel_distances(recording)
self.neighbours_mask = channel_distance <= radius_um
self.peak_sign = peak_sign

# precompute segment slice
self.segment_slices = []
for segment_index in range(recording.get_num_segments()):
i0, i1 = np.searchsorted(self.peaks["segment_index"], [segment_index, segment_index + 1])
self.segment_slices.append(slice(i0, i1))
# For mono-segment, we avoid an extra slice, otherwise make them a tuple for slicing
if sorting.get_num_segments() == 1:
self.segment_slices = None
else:
self.segment_slices = [slice(int(s0), int(s1)) for s0, s1 in segment_slices]

self._kwargs.update(
dict(
sorting=sorting,
channel_from_template=channel_from_template,
extremum_channel_inds=extremum_channel_inds,
radius_um=radius_um,
peak_sign=peak_sign,
include_spikes_in_margin=include_spikes_in_margin,
)
)

@property
def peaks(self):

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.

Who uses this? I could not find a caller in the repo. Accessing it brings back the full allocation that this PR removes. It also calls sorting_to_peaks without self._dtype so the in_margin field would be missing.

Should we remove it?

if self._peaks is not None:
return self._peaks
self._peaks = sorting_to_peaks(self.sorting, self.main_channel_indices)
return self._peaks

def get_margin(self):
return 0
Expand All @@ -225,26 +261,34 @@ def get_dtype(self):
return self._dtype

def get_peak_slice(self, segment_index, start_frame, end_frame, max_margin):
sl = self.segment_slices[segment_index]
peaks_in_segment = self.peaks[sl]
if self.segment_slices is not None:
sl = self.segment_slices[segment_index]
sample_indices_in_segment = self.spike_sample_indices[sl]
else:
sample_indices_in_segment = self.spike_sample_indices
if self.include_spikes_in_margin:
i0, i1 = np.searchsorted(
peaks_in_segment["sample_index"], [start_frame - max_margin, end_frame + max_margin]
)
i0, i1 = np.searchsorted(sample_indices_in_segment, [start_frame - max_margin, end_frame + max_margin])
else:
i0, i1 = np.searchsorted(peaks_in_segment["sample_index"], [start_frame, end_frame])
i0, i1 = np.searchsorted(sample_indices_in_segment, [start_frame, end_frame])
return i0, i1

def compute(self, traces, start_frame, end_frame, segment_index, max_margin, peak_slice):
# get local peaks
sl = self.segment_slices[segment_index]
peaks_in_segment = self.peaks[sl]
i0, i1 = peak_slice
if self.segment_slices is not None:
s0 = self.segment_slices[segment_index].start
spikes = self.spike_vector[s0 + i0 : s0 + i1]
else:
spikes = self.spike_vector[i0:i1]

local_peaks = peaks_in_segment[i0:i1]
local_peaks = np.zeros(spikes.size, dtype=self._dtype)
local_peaks["sample_index"] = spikes["sample_index"]
local_peaks["channel_index"] = self.main_channel_indices[spikes["unit_index"]]
local_peaks["amplitude"] = 0.0
local_peaks["segment_index"] = spikes["segment_index"]
local_peaks["unit_index"] = spikes["unit_index"]

# make sample index local to traces
local_peaks = local_peaks.copy()
local_peaks["sample_index"] -= start_frame - max_margin

# handle flag for margin
Expand Down Expand Up @@ -347,6 +391,15 @@ def __init__(
self.nafter = ms_to_samples(ms_after, sampling_frequency)
self.neighbours_mask = None

self._kwargs.update(
dict(
ms_before=ms_before,
ms_after=ms_after,
nbefore=nbefore,
nafter=nafter,
)
)


class ExtractDenseWaveforms(WaveformsNode):
def __init__(
Expand Down Expand Up @@ -470,6 +523,13 @@ def __init__(
self.neighbours_mask = self.channel_distance <= radius_um
self.max_num_chans = np.max(np.sum(self.neighbours_mask, axis=1))

self._kwargs.update(
dict(
radius_um=radius_um,
sparsity_mask=sparsity_mask,
)
)

def get_margin(self):
return max(self.nbefore, self.nafter)

Expand Down
5 changes: 2 additions & 3 deletions src/spikeinterface/postprocessing/amplitude_scalings.py
Original file line number Diff line number Diff line change
Expand Up @@ -178,13 +178,12 @@ def __init__(
PipelineNode.__init__(self, recording, parents=parents, return_output=return_output)
self.return_in_uV = return_in_uV
if return_in_uV and recording.has_scaleable_traces():
self._dtype = np.float32
self._gains = recording.get_channel_gains()
self._offsets = recording.get_channel_offsets()
else:
self._dtype = recording.get_dtype()
self._gains = None
self._offsets = None
self._dtype = np.float32
spike_retriever = find_parent_of_type(parents, SpikeRetriever)
assert isinstance(
spike_retriever, SpikeRetriever
Expand Down Expand Up @@ -268,7 +267,7 @@ def compute(self, traces, peaks):
collisions = {}

# compute the scaling for each spike
scalings = np.zeros(len(local_spikes), dtype=float)
scalings = np.zeros(len(local_spikes), dtype=self._dtype)
spike_collision_mask = np.zeros(len(local_spikes), dtype=bool)

for spike_index, spike in enumerate(local_spikes):
Expand Down
8 changes: 4 additions & 4 deletions src/spikeinterface/postprocessing/unit_locations.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,10 +7,10 @@

# this dict is for peak location
dtype_localize_by_method = {
"center_of_mass": [("x", "float64"), ("y", "float64")],
"grid_convolution": [("x", "float64"), ("y", "float64"), ("z", "float64")],
"peak_channel": [("x", "float64"), ("y", "float64")],
"monopolar_triangulation": [("x", "float64"), ("y", "float64"), ("z", "float64"), ("alpha", "float64")],
"center_of_mass": [("x", "float32"), ("y", "float32")],

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.

I am wary of the float64 to float32 change. dtype_localize_by_method is also used by localize_peaks in sortingcomponents (center_of_mass.py, monopolar.py and grid.py) so this changes the output for the sorters and motion estimation too, plus the dtype of saved extensions. I would like this to be its own PR so we can discuss it separately from the memory work.

"grid_convolution": [("x", "float32"), ("y", "float32"), ("z", "float32")],
"peak_channel": [("x", "float32"), ("y", "float32")],
"monopolar_triangulation": [("x", "float32"), ("y", "float32"), ("z", "float32"), ("alpha", "float32")],
}

possible_localization_methods = list(dtype_localize_by_method.keys())
Expand Down
9 changes: 9 additions & 0 deletions src/spikeinterface/preprocessing/detect_artifacts.py
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,15 @@ def __init__(
else:
self.diff_threshold_unscaled = None

self._kwargs.update(
dict(
saturation_threshold_uV=saturation_threshold_uV,
diff_threshold_uV=diff_threshold_uV,
proportion=proportion,
signed=signed,
)
)

def get_margin(self) -> int:
"""Return the number of margin samples required on each side of a chunk."""
return 0
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -91,8 +91,8 @@ def __init__(

self.all_channels = all_channels
self.peak_sign = peak_sign
self._kwargs.update(dict(all_channels=all_channels, peak_sign=peak_sign))
self._dtype = recording.get_dtype()
self._kwargs.update(dict(all_channels=all_channels, peak_sign=peak_sign))

def get_dtype(self):
return self._dtype
Expand Down Expand Up @@ -125,8 +125,8 @@ def __init__(
self.channel_distance = get_channel_distances(recording)
self.neighbours_mask = self.channel_distance <= radius_um
self.all_channels = all_channels
self._kwargs.update(dict(radius_um=radius_um, all_channels=all_channels))
self._dtype = recording.get_dtype()
self._kwargs.update(dict(radius_um=radius_um, all_channels=all_channels))

def get_dtype(self):
return self._dtype
Expand Down Expand Up @@ -168,6 +168,7 @@ def __init__(
self.radius_um = radius_um
self.sparse = sparse
self.noise_threshold = noise_threshold
self._dtype = recording.get_dtype()
self._kwargs.update(
dict(
projections=projections,
Expand All @@ -177,7 +178,6 @@ def __init__(
feature=feature,
)
)
self._dtype = recording.get_dtype()

def get_dtype(self):
return self._dtype
Expand Down
Loading