在Camel路由中使用Kafka消费者时Apache Karaf CPU负载过高问题
问题:Apache Karaf+Camel消费Kafka Topic时CPU负载过高
环境与配置
- 运行环境:Apache Karaf 4.4.0,Apache Camel 3.11.7,3节点Apache Kafka 2.13-3.3.1
- 虚拟机配置:2核CPU、32GiB内存
- 消费规模:约10000个Topic,每个Topic对应一个独立消费者
- 集成方式:Java编写Camel路由,以Bundle形式部署到Karaf
Camel路由代码
from ("kafka://crs.topic?brokers=PLAINTEXT://kafka1:9092,kafka2:9093,kafka3:9094&keyDeserializer=org.apache.kafka.common.serialization.StringDeserializer&valueDeserializer=com.rwe.remit.ejb.backend.kafkaclient.api.JacksonReadingSerializer&groupId=testing&heartbeatIntervalMs=120000&maxPollIntervalMs=86400000&sessionTimeoutMs=86400000&maxPollIntervalMs=86400000&deliveryTimeoutMs=86400000&requestTimeoutMs=86400000")
Kafka主节点配置
process.roles=broker,controller node.id=1 controller.quorum.voters=1@kafka1:19092,2@kafka2:19093,3@kafka3:19094 listeners=PLAINTEXT://kafka1:9092,CONTROLLER://kafka1:19092 inter.broker.listener.name=PLAINTEXT advertised.listeners=PLAINTEXT://kafka1:9092 controller.listener.names=CONTROLLER listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,SSL:SSL,SASL_PLAINTEXT:SASL_PLAINTEXT,SASL_SSL:SASL_SSL num.network.threads=3 num.io.threads=8 socket.send.buffer.bytes=102400 socket.receive.buffer.bytes=102400 socket.request.max.bytes=104857600 log.dirs=/data/apache-karaf/node1/log num.partitions=3 num.recovery.threads.per.data.dir=1 offsets.topic.replication.factor=3 transaction.state.log.replication.factor=3 transaction.state.log.min.isr=3 log.retention.hours=168 log.segment.bytes=1073741824 log.retention.check.interval.ms=864000000 group.max.session.timeout.ms=86400000
问题现象
启动全部900个Bundle后,消费者开始监听Topic,但未发送任何测试数据时,Karaf所在Linux虚拟机CPU负载已达100%并持续保持。此前使用相同Karaf+Camel架构消费ActiveMQ队列时无性能问题。
线程转储信息
以下是两个Kafka消费者的线程转储,其余近万个Topic的消费者状态一致,均持续处于RUNNABLE状态且占用CPU:
"Camel (crs-rwe-tenants-test) thread #71 - KafkaConsumer[crs.topic1]" #394 daemon prio=5 os_prio=0 tid=0x00007f17b8374000 nid=0x1acf51 runnable [0x00007f1796253000] java.lang.Thread.State: RUNNABLE at java.lang.Integer.valueOf(Integer.java:832) at sun.nio.ch.EPollSelectorImpl.updateSelectedKeys(EPollSelectorImpl.java:120) at sun.nio.ch.EPollSelectorImpl.doSelect(EPollSelectorImpl.java:98) at sun.nio.ch.SelectorImpl.lockAndDoSelect(SelectorImpl.java:86) - locked <0x00000006d72f07c8> (a sun.nio.ch.Util$3) - locked <0x00000006d72f07b8> (a java.util.Collections$UnmodifiableSet) - locked <0x00000006d72cea40> (a sun.nio.ch.EPollSelectorImpl) at sun.nio.ch.SelectorImpl.select(SelectorImpl.java:97) at org.apache.kafka.common.network.Selector.select(Selector.java:873) at org.apache.kafka.common.network.Selector.poll(Selector.java:465) at org.apache.kafka.clients.NetworkClient.poll(NetworkClient.java:560) at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:280) at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:251) at org.apache.kafka.clients.consumer.KafkaConsumer.pollForFetches(KafkaConsumer.java:1306) at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1242) at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1168) at org.apache.camel.component.kafka.KafkaConsumer$KafkaFetchRecords.doPollRun(KafkaConsumer.java:351) at org.apache.camel.component.kafka.KafkaConsumer$KafkaFetchRecords.doRun(KafkaConsumer.java:279) at org.apache.camel.component.kafka.KafkaConsumer$KafkaFetchRecords.run(KafkaConsumer.java:246) at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) at java.util.concurrent.FutureTask.run(FutureTask.java:266) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) at java.lang.Thread.run(Thread.java:750) "Camel (crs-rwe-tenants-allgaeukraft) thread #79 - KafkaConsumer[crs.topic2]" #413 daemon prio=5 os_prio=0 tid=0x00007f17b895b000 nid=0x1acf64 runnable [0x00007f1795142000] java.lang.Thread.State: RUNNABLE at sun.nio.ch.EPollArrayWrapper.epollWait(Native Method) at sun.nio.ch.EPollArrayWrapper.poll(EPollArrayWrapper.java:269) at sun.nio.ch.EPollSelectorImpl.doSelect(EPollSelectorImpl.java:93) at sun.nio.ch.SelectorImpl.lockAndDoSelect(SelectorImpl.java:86) - locked <0x00000006d78c5610> (a sun.nio.ch.Util$3) - locked <0x00000006d78c5600> (a java.util.Collections$UnmodifiableSet) - locked <0x00000006d78a3710> (a sun.nio.ch.EPollSelectorImpl) at sun.nio.ch.SelectorImpl.select(SelectorImpl.java:97) at org.apache.kafka.common.network.Selector.select(Selector.java:873) at org.apache.kafka.common.network.Selector.poll(Selector.java:465) at org.apache.kafka.clients.NetworkClient.poll(NetworkClient.java:560) at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:280) at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:251) at org.apache.kafka.clients.consumer.KafkaConsumer.pollForFetches(KafkaConsumer.java:1306) at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1242) at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1168) at org.apache.camel.component.kafka.KafkaConsumer$KafkaFetchRecords.doPollRun(KafkaConsumer.java:351) at org.apache.camel.component.kafka.KafkaConsumer$KafkaFetchRecords.doRun(KafkaConsumer.java:279) at org.apache.camel.component.kafka.KafkaConsumer$KafkaFetchRecords.run(KafkaConsumer.java:246) at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) at java.util.concurrent.FutureTask.run(FutureTask.java:266) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) at java.lang.Thread.run(Thread.java:750)
疑问
- 为何Linux虚拟机负载过高?
- 是否有方法查看Java或Apache Karaf的内部进程与流量?
- 线程持续处于RUNNABLE状态且占满CPU的原因是什么?
- 能否配置消费者仅在Topic有数据时运行,无数据时处于待机状态?
内容的提问来源于stack exchange,提问作者Amjad Farajallah
相关产品推荐
相关产品推荐

