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
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
package com.kakao.actionbase.core.edge.payload

import com.kakao.actionbase.core.metadata.common.AggregationType

import com.fasterxml.jackson.annotation.JsonSubTypes
import com.fasterxml.jackson.annotation.JsonTypeInfo

data class AggregationSweepRequest(
val items: List<SweepItem>,
)

data class SweepItem(
val type: AggregationType,
@field:JsonTypeInfo(
use = JsonTypeInfo.Id.NAME,
include = JsonTypeInfo.As.EXTERNAL_PROPERTY,
property = "type",
)
@field:JsonSubTypes(
JsonSubTypes.Type(value = TopkSweepItem::class, name = "TOPK"),
)
val item: SweepItemPayload,
)

sealed interface SweepItemPayload

data class TopkSweepItem(
val database: String,
val table: String,
val topk: String,
val source: String,
val target: String,
val direction: String,
val ranges: String = "",
val entity: String,
val topkDimensionValue: String,
val dimensionValues: String = "",
val properties: Map<String, String> = emptyMap(),
val refreshAt: Long = -1,
) : SweepItemPayload
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
package com.kakao.actionbase.core.edge.payload

data class AggregationSweepResult(
val database: String,
val table: String,
val topk: String,
val entity: String,
val status: String,
val error: String?,
)
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
package com.kakao.actionbase.core.edge.payload

data class AggregationsSweepResponse(
val items: List<Item>,
) {
data class Item(
val database: String,
val table: String,
val topk: String,
val entity: String,
val status: String,
val error: String?,
)

companion object {
fun from(sweepResults: List<AggregationSweepResult>): AggregationsSweepResponse =
AggregationsSweepResponse(
items =
sweepResults.map { result ->
Item(
database = result.database,
table = result.table,
topk = result.topk,
entity = result.entity,
status = result.status,
error = result.error,
)
},
)
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@ package com.kakao.actionbase.engine.service

import com.kakao.actionbase.core.edge.payload.AggregationItemPayload
import com.kakao.actionbase.core.edge.payload.AggregationResult
import com.kakao.actionbase.core.edge.payload.AggregationSweepResult
import com.kakao.actionbase.core.edge.payload.SweepItem
import com.kakao.actionbase.core.metadata.QualifiedAggregations
import com.kakao.actionbase.core.metadata.common.AggregationType
import com.kakao.actionbase.engine.AggregationEngine
Expand All @@ -23,4 +25,12 @@ class AggregationService(
.fromIterable(items)
.flatMap { item -> Flux.merge(handlersByType.values.map { it.aggregate(item) }) }
.collectList()

fun sweep(items: List<SweepItem>): Mono<List<AggregationSweepResult>> =
Flux
.fromIterable(items)
.flatMap { item -> handler(item.type).sweep(item.item) }
.collectList()

private fun handler(type: AggregationType): AggregationHandler = handlersByType[type] ?: error("No aggregation handler for type $type")
}
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,12 @@ package com.kakao.actionbase.engine.service.aggregation

import com.kakao.actionbase.core.edge.payload.AggregationItemPayload
import com.kakao.actionbase.core.edge.payload.AggregationResult
import com.kakao.actionbase.core.edge.payload.AggregationSweepResult
import com.kakao.actionbase.core.edge.payload.SweepItemPayload
import com.kakao.actionbase.core.metadata.common.AggregationType

import reactor.core.publisher.Flux
import reactor.core.publisher.Mono

/**
* One aggregation kind's write path. [AggregationService] dispatches to the handler
Expand All @@ -16,4 +19,7 @@ interface AggregationHandler {

/** Aggregates a single edge event and writes its result rows. */
fun aggregate(item: AggregationItemPayload): Flux<AggregationResult>

/** Recomputes a single refreshed event and re-writes its result row. */
fun sweep(item: SweepItemPayload): Mono<AggregationSweepResult>
}
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,11 @@ import com.kakao.actionbase.v2.core.metadata.Direction as V2Direction
import com.kakao.actionbase.core.edge.Edge
import com.kakao.actionbase.core.edge.payload.AggregationItemPayload
import com.kakao.actionbase.core.edge.payload.AggregationResult
import com.kakao.actionbase.core.edge.payload.AggregationSweepResult
import com.kakao.actionbase.core.edge.payload.EdgeBulkMutationRequest.MutationItem
import com.kakao.actionbase.core.edge.payload.MutationResult
import com.kakao.actionbase.core.edge.payload.SweepItemPayload
import com.kakao.actionbase.core.edge.payload.TopkSweepItem
import com.kakao.actionbase.core.metadata.common.AggregationConstants
import com.kakao.actionbase.core.metadata.common.AggregationType
import com.kakao.actionbase.core.metadata.common.Aggregations
Expand Down Expand Up @@ -54,6 +57,58 @@ class TopkAggregationHandler(
},
).flatMap { (event, direction, topk) -> processTopk(event, direction, topk) }

override fun sweep(item: SweepItemPayload): Mono<AggregationSweepResult> {
require(item is TopkSweepItem) { "TopkAggregationHandler handles TopkSweepItem, got ${item::class.simpleName}" }

fun result(
status: String,
error: String? = null,
) = AggregationSweepResult(
database = item.database,
table = item.table,
topk = item.topk,
entity = item.entity,
status = status,
error = error,
)

val tb = engine.getTableBinding(database = item.database, alias = item.table)
val group =
tb.schema.groups
.firstOrNull { g -> g.aggregations.topk.any { it.topk == item.topk } }
?: return Mono.just(result(SKIPPED))

val topk = group.aggregations.topk.first { it.topk == item.topk }
val direction = Direction.valueOf(item.direction)
val directedSource = if (direction == Direction.IN) item.target else item.source
val dimensionValues = AggregationConstants.Topk.splitValues(item.dimensionValues)
val (rankDatabase, rankTable) = parseFqn(topk.rank)

return writeRank(
sourceDatabase = item.database,
sourceTable = tb.table,
group = group.group,
start = directedSource,
direction = direction,
ranges = item.ranges.takeIf { it.isNotEmpty() },
ranking =
Ranking(
database = rankDatabase,
table = rankTable,
topk = item.topk,
entity = item.entity,
topkDimensionValue = item.topkDimensionValue,
dimensionValues = dimensionValues,
properties = item.properties,
),
).map { results ->
result(status = if (results.any { it.status == ERROR }) ERROR else SUCCESS)
}.onErrorResume { err ->
logger.error("topk sweep failed for {}.{} topk={}", item.database, item.table, item.topk, err)
Mono.just(result(status = ERROR, error = err.message))
}
}

private fun createTopkEvent(item: AggregationItemPayload): List<EdgeAggregationEvent> {
val tb = engine.getTableBinding(database = item.database, alias = item.table)

Expand Down Expand Up @@ -122,7 +177,7 @@ class TopkAggregationHandler(

/**
* Runs the aggregation query for one ranking and writes its rank row.
* Used by the aggregate flow, which then enqueues a refresh message
* Shared by the aggregate flow (which then enqueues a refresh message) and the sweep flow
* (which recomputes from an already-resolved [Ranking]).
*/
private fun writeRank(
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
package com.kakao.actionbase.engine.service

import com.kakao.actionbase.core.edge.payload.SweepItem
import com.kakao.actionbase.core.edge.payload.TopkSweepItem
import com.kakao.actionbase.core.metadata.QualifiedAggregations
import com.kakao.actionbase.core.metadata.common.AggregationType
import com.kakao.actionbase.engine.AggregationEngine
Expand All @@ -10,10 +12,11 @@ import org.junit.jupiter.api.Test
import io.kotest.matchers.collections.shouldContainExactlyInAnyOrder
import io.mockk.every
import io.mockk.mockk
import reactor.test.StepVerifier

/**
* `AggregationService` is a thin dispatcher: it forwards metadata lookups to the engine and routes
* each item to the handler registered for its type. Per-type behavior (top-K ranking) lives
* each item to the handler registered for its type. Per-type behavior (top-K ranking, refresh) lives
* in the handlers, and dispatch through a real handler is exercised end-to-end by
* `MetadataAggQueryControllerE2ETest` — so here we only pin what the dispatcher itself owns.
*/
Expand All @@ -39,4 +42,32 @@ class AggregationServiceTest {
service.getAggregations(AggregationType.TOPK) shouldContainExactlyInAnyOrder listOf(entry)
}
}

@Nested
inner class Sweep {
@Test
fun `errors when no handler is registered for the type`() {
val service = AggregationService(engine, emptyList())

StepVerifier
.create(service.sweep(listOf(sweepItem())))
.verifyError(IllegalStateException::class.java)
}
}
}

private fun sweepItem(): SweepItem =
SweepItem(
type = AggregationType.TOPK,
item =
TopkSweepItem(
database = "commerce",
table = "orders",
topk = "top_purchased",
source = "user1",
target = "item1",
direction = "OUT",
entity = "user1",
topkDimensionValue = "item1",
),
)
Loading
Loading