Группировать данные в flow в kotlin по времени
Возникла у меня тут задача:
- Есть поток данных. Довольно плотный: тысячи объектов в секунду.
- Его нужно писать в БД, но обрабатывать по одному долго.
- Соответственно, нужно собирать некоторое количество записей за определённый промежуток времени и писать одной транзакцией.
Среди стандартных методов ничего нужного не нашёл.
Сам, пока что, додумался только до такого:
suspend fun <T> Flow<T>.groupByTime(timeout : Duration): Flow<List<T>> = flow<List<T>> {
while (true) {
val list = mutableListOf<T>()
coroutineScope {
val job = launch(currentCoroutineContext()) {
collect { list.add(it) }
}
job.start()
delay(timeout)
job.cancelAndJoin()
}
if (list.size > 0) emit(list)
}
}
Если у кого есть решения лучше, предлагайте.
Ответы (1 шт):
Автор решения: IR42
→ Ссылка
У вас как минимум 3 ошибки:
- Постоянно вызывается this.collect, это приведёт к тому, что вы будете получать одни и те же значения.
- list используется разными корутинами, можно получить состояние гонки https://kotlinlang.org/docs/shared-mutable-state-and-concurrency.html#the-problem
- Ваша функция никогда не закончится из-за
while (true)без выхода. Спасёт только cancel на корутину, в которой запускается этот while.
Как мне кажется, правильным решением будет использовать select.
Пример (Использовать на свой страх и риск! Не протестировано, может неправильно работать в особых случаях):
@ExperimentalCoroutinesApi
fun <T> Flow<T>.groupByTime(timeout: Long): Flow<List<T>> = flow {
coroutineScope {
val list = mutableListOf<T>()
val values = produce { collect { send(it) } }
val ticker = produce {
while (true) {
delay(timeout)
send(Unit)
}
}
suspend fun emitList() {
if (list.isNotEmpty()) {
val value = list.toList()
emit(value)
list.clear()
}
}
whileSelect {
values.onReceiveCatching { result ->
result
.onSuccess { list += it }
.onFailure {
it?.let { throw it }
ticker.cancel()
emitList()
}
.isSuccess
}
ticker.onReceive {
emitList()
true
}
}
}
}