【问题标题】:Why does having serializer in threadlocal result in memory leak in kafka producer?为什么在 threadlocal 中使用序列化程序会导致 kafka 生产者内存泄漏?
【发布时间】:2019-07-03 04:10:29
【问题描述】:

考虑以下设置

prop.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, ThriftSerializer.class.getName());


public class ThriftSerializer implements Serializer<TBase> {
    private final ThreadLocal<TSerializer> serializer = new ThreadLocalTSerializer();

    @Override
    public void configure(Map map, boolean b) {

    }

    @Override
    public byte[] serialize(String s, TBase event) {
        try {
            return serializer.get().serialize(event);
        } catch (TException e) {
            return new byte[0];
        }
    }

    @Override
    public void close() {

    }
}

以上代码导致内存泄漏

但我不明白为什么会这样。 kafka producer 会创建很多不死的线程吗?

如果上面的代码被替换为

@Override
public byte[] serialize(String s, TBase event) {
    TSerializer serializer = new TSerializer();

    try {
        return serializer.serialize(event);
    } catch (TException e) {
        return new byte[0];
    }
}

然后内存泄漏消失了,这是有道理的,但是对于每个事件,它都会创建需要进行垃圾收集的新对象,如果吞吐量很高,可能会导致 gc 压力

有人可以指出我理解这种行为的方向吗?

【问题讨论】:

    标签: java multithreading apache-kafka thread-safety


    【解决方案1】:

    据我所知,KafkaProducer 是线程安全的,跨线程共享单个生产者实例通常比拥有多个实例要快。

    但是 send 方法是异步的(除非你在 send 方法返回的 Future 对象上调用 .get() ,不建议这样做,这样你会等待每次发送,因此会以同步的方式处理它们)。

    根据文档,生产者由一个缓冲空间池组成,该池保存尚未传输到服务器的记录以及 一个后台 I/O 线程,负责将这些记录转换为请求并将它们传输到集群。使用后未能关闭生产者会泄漏这些资源。

    看来 send 方法实际上是使用后台线程来转换你的记录并将其发送到集群。

    最后你真的要关闭生产者吗?

        producer.flush();
        producer.close();
    

    当要关闭 Kafka 会话时,会调用序列化程序的 close 方法。 我的猜测是,您可以尝试在序列化程序的 close 方法中进行一些额外的清理,或者将其标记为符合垃圾回收条件。

    【讨论】:

      猜你喜欢
      • 2011-06-25
      • 1970-01-01
      • 1970-01-01
      • 2012-02-09
      • 2013-07-12
      • 1970-01-01
      • 2017-02-13
      • 1970-01-01
      相关资源
      最近更新 更多