【问题标题】:flink AggregateFunction accumulator slowflink AggregateFunction 累加器慢
【发布时间】: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

标签: aggregate apache-flink


【解决方案1】:

您可以尝试增加 flink 托管内存。第二个参数是内存大小。

config.setString("taskmanager.memory.managed.size", "8g")
val env: StreamExecutionEnvironment = StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(config)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2019-05-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多