Improve chunk algorithm
This commit is contained in:
parent
0355c0bda9
commit
502fd472ed
2 changed files with 96 additions and 7 deletions
|
|
@ -28,16 +28,47 @@ class WorkerDataConverter(
|
||||||
// First try to serialize all requests at once. In the vast majority of cases this will work.
|
// First try to serialize all requests at once. In the vast majority of cases this will work.
|
||||||
return serializeRequests(notificationEventRequests)
|
return serializeRequests(notificationEventRequests)
|
||||||
.map { listOf(it) }
|
.map { listOf(it) }
|
||||||
.recoverCatching {
|
.recoverCatching { t ->
|
||||||
if (it is DataForWorkManagerIsTooBig) {
|
if (t is DataForWorkManagerIsTooBig) {
|
||||||
// Perform serialization on sublists, workDataOf have failed because of size limit
|
// Perform serialization on sublists, workDataOf have failed because of size limit
|
||||||
Timber.w(it, "Failed to serialize ${notificationEventRequests.size} notification requests, trying with chunks of $CHUNK_SIZE.")
|
Timber.w(t, "Failed to serialize ${notificationEventRequests.size} notification requests, split the requests per room.")
|
||||||
// TODO Do not split rooms
|
// Group the requests per rooms
|
||||||
notificationEventRequests.chunked(CHUNK_SIZE).mapNotNull { chunk ->
|
val requestsSortedPerRoom = notificationEventRequests.groupBy { it.roomId }.values
|
||||||
serializeRequests(chunk).getOrNull()
|
// Build a list of sublist with size at most CHUNK_SIZE, and with all rooms kept together
|
||||||
|
buildList {
|
||||||
|
val currentChunk = mutableListOf<NotificationEventRequest>()
|
||||||
|
for (requests in requestsSortedPerRoom) {
|
||||||
|
if (currentChunk.size + requests.size <= CHUNK_SIZE) {
|
||||||
|
// Can add the whole room requests to the current chunk
|
||||||
|
currentChunk.addAll(requests)
|
||||||
|
} else {
|
||||||
|
// Add the current chunk
|
||||||
|
add(currentChunk.toList())
|
||||||
|
// Start a new chunk with the current room requests
|
||||||
|
currentChunk.clear()
|
||||||
|
// If a room has more requests than CHUNK_SIZE, we need to split them
|
||||||
|
requests.chunked(CHUNK_SIZE) { chunk ->
|
||||||
|
if (chunk.size == CHUNK_SIZE) {
|
||||||
|
add(chunk.toList())
|
||||||
|
} else {
|
||||||
|
currentChunk.addAll(chunk)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
// Add any remaining requests
|
||||||
|
add(currentChunk.toList())
|
||||||
}
|
}
|
||||||
|
.filter { it.isNotEmpty() }
|
||||||
|
.also {
|
||||||
|
Timber.d("Split notification requests into ${it.size} chunks for WorkManager serialization")
|
||||||
|
it.forEach { requests ->
|
||||||
|
Timber.d(" - Chunk with ${requests.size} requests")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
.mapNotNull { serializeRequests(it).getOrNull() }
|
||||||
} else {
|
} else {
|
||||||
throw it
|
throw t
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -15,6 +15,7 @@ import io.element.android.libraries.matrix.test.AN_EVENT_ID
|
||||||
import io.element.android.libraries.matrix.test.AN_EVENT_ID_2
|
import io.element.android.libraries.matrix.test.AN_EVENT_ID_2
|
||||||
import io.element.android.libraries.matrix.test.A_ROOM_ID
|
import io.element.android.libraries.matrix.test.A_ROOM_ID
|
||||||
import io.element.android.libraries.matrix.test.A_ROOM_ID_2
|
import io.element.android.libraries.matrix.test.A_ROOM_ID_2
|
||||||
|
import io.element.android.libraries.matrix.test.A_ROOM_ID_3
|
||||||
import io.element.android.libraries.matrix.test.A_SESSION_ID
|
import io.element.android.libraries.matrix.test.A_SESSION_ID
|
||||||
import io.element.android.libraries.matrix.test.A_SESSION_ID_2
|
import io.element.android.libraries.matrix.test.A_SESSION_ID_2
|
||||||
import io.element.android.libraries.push.api.push.NotificationEventRequest
|
import io.element.android.libraries.push.api.push.NotificationEventRequest
|
||||||
|
|
@ -60,6 +61,63 @@ class WorkerDataConverterTest {
|
||||||
val serialized = sut.serialize(data)
|
val serialized = sut.serialize(data)
|
||||||
assertThat(serialized.getOrNull()?.size).isGreaterThan(1)
|
assertThat(serialized.getOrNull()?.size).isGreaterThan(1)
|
||||||
assertThat(serialized.getOrNull()?.size).isEqualTo(100 / WorkerDataConverter.CHUNK_SIZE)
|
assertThat(serialized.getOrNull()?.size).isEqualTo(100 / WorkerDataConverter.CHUNK_SIZE)
|
||||||
|
// All the items are present
|
||||||
|
val deserialized = serialized.getOrNull()?.flatMap { sut.deserialize(it)!! }
|
||||||
|
assertThat(deserialized).containsExactlyElementsIn(data)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `serializing lots of data leads to several work data generated case 2`() {
|
||||||
|
val data = List(101) {
|
||||||
|
NotificationEventRequest(
|
||||||
|
sessionId = A_SESSION_ID,
|
||||||
|
roomId = A_ROOM_ID,
|
||||||
|
eventId = EventId(AN_EVENT_ID.value + it),
|
||||||
|
providerInfo = "info$it",
|
||||||
|
)
|
||||||
|
}
|
||||||
|
val sut = WorkerDataConverter(DefaultJsonProvider())
|
||||||
|
val serialized = sut.serialize(data)
|
||||||
|
assertThat(serialized.getOrNull()?.size).isGreaterThan(1)
|
||||||
|
assertThat(serialized.getOrNull()?.size).isEqualTo(100 / WorkerDataConverter.CHUNK_SIZE + 1)
|
||||||
|
// All the items are present
|
||||||
|
val deserialized = serialized.getOrNull()?.flatMap { sut.deserialize(it)!! }
|
||||||
|
assertThat(deserialized).containsExactlyElementsIn(data)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `serializing lots of data leads to several work data generated case 3`() {
|
||||||
|
val data1 = List(15) {
|
||||||
|
NotificationEventRequest(
|
||||||
|
sessionId = A_SESSION_ID,
|
||||||
|
roomId = A_ROOM_ID,
|
||||||
|
eventId = EventId(AN_EVENT_ID.value + it),
|
||||||
|
providerInfo = "info".repeat(100) + it,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
val data2 = List(3) {
|
||||||
|
NotificationEventRequest(
|
||||||
|
sessionId = A_SESSION_ID,
|
||||||
|
roomId = A_ROOM_ID_2,
|
||||||
|
eventId = EventId(AN_EVENT_ID.value + it),
|
||||||
|
providerInfo = "info".repeat(100) + it,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
val data3 = List(7) {
|
||||||
|
NotificationEventRequest(
|
||||||
|
sessionId = A_SESSION_ID,
|
||||||
|
roomId = A_ROOM_ID_3,
|
||||||
|
eventId = EventId(AN_EVENT_ID.value + it),
|
||||||
|
providerInfo = "info".repeat(100) + it,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
val data = (data1 + data2 + data3).shuffled()
|
||||||
|
val sut = WorkerDataConverter(DefaultJsonProvider())
|
||||||
|
val serialized = sut.serialize(data)
|
||||||
|
assertThat(serialized.getOrNull()?.size).isEqualTo(2)
|
||||||
|
// All the items are present
|
||||||
|
val deserialized = serialized.getOrNull()?.flatMap { sut.deserialize(it)!! }
|
||||||
|
assertThat(deserialized).containsExactlyElementsIn(data)
|
||||||
}
|
}
|
||||||
|
|
||||||
private fun testIdentity(data: List<NotificationEventRequest>) {
|
private fun testIdentity(data: List<NotificationEventRequest>) {
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue