Kafka S3 Sink连接器同步22个主题OOM,拆分后正常的根因排查
Kafka S3 Sink连接器Java堆内存溢出问题排查
问题现象
使用Kafka S3 Sink连接器同步22个主题的数据到S3存储桶时,任务抛出不可恢复的内存溢出错误:
ERROR [prod-partnerbilling-sink-v3|task-2] WorkerSinkTask{id=prod-partnerbilling-sink-v3-2} Task threw an uncaught and unrecoverable exception. Task is being killed and will not recover until manually restarted (org.apache.kafka.connect.runtime.WorkerTask:193) org.apache.kafka.connect.errors.ConnectException: Exiting WorkerSinkTask due to unrecoverable exception. at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:609) at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:329) at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:232) at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:201) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:186) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:241) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829) Caused by: java.lang.OutOfMemoryError: Java heap space
将22个主题拆分为两组各11个主题,使用完全相同的连接器配置时,任务运行正常。
根本原因分析
- 单任务负载过载:单个Sink连接器任务同时处理22个主题的所有分区时,每个分区的消费缓冲区、待处理消息、中间序列化对象会持续占用堆内存。主题数量翻倍后,内存中积压的数据量超出JVM堆内存上限,触发溢出。拆分后每个任务的负载减半,内存消耗控制在阈值内。
- 任务分配不合理:如果连接器的
tasks.max配置值过小,22个主题的所有分区会集中分配到少量任务中(如报错的task-2承担了过多分区的处理),单个任务需要并行处理大量消息,内存占用急剧上升。 - Worker堆内存配置不足:Kafka Connect Worker默认的JVM堆内存(通过
KAFKA_HEAP_OPTS配置)通常较低,当同时处理多主题高吞吐量数据时,无法容纳消费缓存、S3批量上传缓存等数据。拆分后任务的内存需求降低,不会触发溢出。 - 批量上传配置过载:若连接器配置了较大的
flush.size或rotate.interval.ms,22个主题的消息会同时累积到批量上传阈值,内存中会持有大量待上传的消息数据,直接导致堆内存耗尽。拆分后每个连接器的批量数据量减少,内存占用可控。
内容的提问来源于stack exchange,提问作者scriptkidie
相关产品推荐
相关产品推荐

