Skip to content
Merged
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
3 changes: 1 addition & 2 deletions machine-learning/immich_ml/sessions/rknn/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -64,8 +64,7 @@ def run(
run_options: Any = None,
) -> list[NDArray[np.float32]]:
input_data: list[NDArray[np.float32]] = [np.ascontiguousarray(v) for v in input_feed.values()]
self.rknnpool.put(input_data)
res = self.rknnpool.get()
res = self.rknnpool.run(input_data)
if res is None:
raise RuntimeError("RKNN inference failed!")
return res
Expand Down
18 changes: 8 additions & 10 deletions machine-learning/immich_ml/sessions/rknn/rknnpool.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,9 @@
# Following Apache License 2.0

import logging
from concurrent.futures import Future, ThreadPoolExecutor
import threading
from concurrent.futures import ThreadPoolExecutor
from pathlib import Path
from queue import Queue
from typing import Callable

import numpy as np
Expand Down Expand Up @@ -66,20 +66,18 @@ def __init__(
func: Callable[["RKNNLite", list[NDArray[np.float32]]], list[NDArray[np.float32]]],
) -> None:
self.tpes = tpes
self.queue: Queue[Future[list[NDArray[np.float32]]]] = Queue()
self.rknn_pool = [init_rknn(model_path) for _ in range(tpes)]
self.pool = ThreadPoolExecutor(max_workers=tpes)
self.func = func
self.num = 0
self.lock = threading.Lock()

def put(self, inputs: list[NDArray[np.float32]]) -> None:
self.queue.put(self.pool.submit(self.func, self.rknn_pool[self.num % self.tpes], inputs))
self.num += 1
def run(self, inputs: list[NDArray[np.float32]]) -> list[NDArray[np.float32]]:
with self.lock:
idx = self.num % self.tpes
self.num += 1

def get(self) -> list[NDArray[np.float32]] | None:
if self.queue.empty():
return None
fut = self.queue.get()
fut = self.pool.submit(self.func, self.rknn_pool[idx], inputs)
return fut.result()

def release(self) -> None:
Expand Down
2 changes: 1 addition & 1 deletion machine-learning/test_main.py
Original file line number Diff line number Diff line change
Expand Up @@ -564,7 +564,7 @@ def test_run_rknn(self, rknn_session: mock.Mock, mocker: MockerFixture) -> None:

session.run(None, input_feed)

rknn_session.return_value.put.assert_called_once_with([input1, input2])
rknn_session.return_value.run.assert_called_once_with([input1, input2])
assert np_spy.call_count == 2
np_spy.assert_has_calls([mock.call(input1), mock.call(input2)])

Expand Down
Loading