Группировать данные в flow в kotlin по времени

Возникла у меня тут задача:

  1. Есть поток данных. Довольно плотный: тысячи объектов в секунду.
  2. Его нужно писать в БД, но обрабатывать по одному долго.
  3. Соответственно, нужно собирать некоторое количество записей за определённый промежуток времени и писать одной транзакцией.

Среди стандартных методов ничего нужного не нашёл.

Сам, пока что, додумался только до такого:

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 ошибки:

  1. Постоянно вызывается this.collect, это приведёт к тому, что вы будете получать одни и те же значения.
  2. list используется разными корутинами, можно получить состояние гонки https://kotlinlang.org/docs/shared-mutable-state-and-concurrency.html#the-problem
  3. Ваша функция никогда не закончится из-за 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
            }
        }
    }
}
→ Ссылка