Flink High CPU usage for specific TaskManager

Viewed 417

With Flink v1.13, I have 3 taskmanagers and 1 jobmanager with un-bounded stream using rocksdb as a state backend (all of them are running separate servers). When I start my flink application, one of the taskmanager's cpu load is higher than other. I am just calculating taskmanager jvm cpu load with promethues using: (flink_taskmanager_Status_JVM_CPU_Load{job="Flink"}*100)

When I start the application, CPU load:

  • TaskManager 1 => 46.0 %,
  • TaskManager 2 => 54.0 %,
  • TaskManager 3 => 76.0 %,

After 4-6 six hours, CPU load:

  • TaskManager 1 => 40-45.0 %,
  • TaskManager 2 => 50-57.0 %,
  • TaskManager 3 => 88-95.0 %, (almost 100% !!)

How can i distribute load equally ?

Rocksdb can also be reason for the high cpu? (Or how can i measure the rocksdb cpu usage ?)

1 Answers

This skew might be caused by

  • having more active task slots in TM3
  • data skew (e.g., a key (or set of keys) with more than its fair share of event traffic has been assigned to TM3)
  • if the keyspace is very small, perhaps TM3 has been assigned more than its fair share of keys

You might start by looking at the network metrics to see if the event traffic is well balanced, and at the checkpointing metrics to see the checkpoint sizes are roughly comparable (RocksDB checkpoint sizes are a very poor indicator of active state size, but if the checkpoints from the various TMs are wildly different in size that would be interesting information).

To get more insight, you can look at various RocksDB metrics.

Once you know the cause, you can work on fixing it. Possible solutions include:

  • pre-aggregate (e.g., some sort of local-global aggregation)
  • increase the number of keys (finer-grained keys)
  • increase the parallelism
Related