【发布时间】:2020-02-29 07:00:33
【问题描述】:
我正在使用 flink 进行数据聚合 类似的代码
.keyby()
.timewidow()
.aggregate(new AggCount(), new AggWindow())
class AggCountNew() extends AggregateFunction[((String, String, Int), Message), OutMessage, OutMessage] {
private val logger = LoggerFactory.getLogger(getClass)
override def createAccumulator(): OutMessage = OutMessage(null, null, 0, new mutable.OpenHashMap[String, Long](128), new mutable.OpenHashMap[String, Long](32768), Set(), Set())
override def add(value: ((String, String, Int), OutMessage), accumulator: OutMessage): OutMessage = {
val dataMap= accumulator.dataMap
dataMap(value._2.device) = dataMap.getOrElse[Long](value._2.device, 0) + 1
accumulator
}
override def getResult(accumulator: HbaseOneDayMessage): HbaseOneDayMessage = {
return accumulator
}
override def merge(a: Message, b: Message): HbaseOneDayMessage = {
a
}
}
当dataMap有超过1000个key时,吞吐量很低
【问题讨论】:
-
我已经初始化了HashMap能力
-
您使用的是哪个州的后端?
-
你的问题很模糊。你到底看到/期待什么?实际上要计算什么?
-
RocksDBStateBackend