JVM Profiler Reporter介绍

简介: 开篇 JVM Profiler采集完数据后可以通过多种途径上报数据,对接Console,File,redis,kafka等,这篇文章会把源码罗列一下毕竟都很简单。

开篇

 JVM Profiler采集完数据后可以通过多种途径上报数据,对接Console,File,redis,kafka等,这篇文章会把源码罗列一下毕竟都很简单。
 JVM Profiler提供灵活的框架可以集成更多的Reporter,只要实现Reporter接口即可,看你个人意愿了,反正github上有源码,直接集成编译打包即可。


img_32f35386c75b58da37ff4cd050da4b01.png


ConsoleOutputReporter

  • 简单明了的通过Sytem.out.println来上报监控数据。
public class ConsoleOutputReporter implements Reporter {
    @Override
    public void report(String profilerName, Map<String, Object> metrics) {
        System.out.println(String.format("ConsoleOutputReporter - %s: %s", profilerName, JsonUtils.serialize(metrics)));
    }

    @Override
    public void close() {
    }
}


FileOutputReporter

  • 在指定的目录创建采集数据记录文件。
  • 通过FileWriter来往文件写入数据。
public class FileOutputReporter implements Reporter {
    private static final AgentLogger logger = AgentLogger.getLogger(FileOutputReporter.class.getName());
    
    private String directory;
    private ConcurrentHashMap<String, FileWriter> fileWriters = new ConcurrentHashMap<>();
    private volatile boolean closed = false;
    
    public FileOutputReporter() {
    }

    public String getDirectory() {
        return directory;
    }

    public void setDirectory(String directory) {
        synchronized (this) {
            if (this.directory == null || this.directory.isEmpty()) {
                this.directory = directory;
            } else {
                throw new RuntimeException(String.format("Cannot set directory to %s because it is already has value %s", directory, this.directory));
            }
        }
    }

    @Override
    public synchronized void report(String profilerName, Map<String, Object> metrics) {
        if (closed) {
            logger.info("Report already closed, do not report metrics");
            return;
        }
        
        FileWriter writer = ensureFile(profilerName);
        try {
            writer.write(JsonUtils.serialize(metrics));
            writer.write(System.lineSeparator());
            writer.flush();
        } catch (IOException e) {
            e.printStackTrace();
        }
    }

    @Override
    public synchronized void close() {
        closed = true;
        
        List<FileWriter> copy = new ArrayList<>(fileWriters.values());
        for (FileWriter entry : copy) {
            try {
                entry.flush();
                entry.close();
            } catch (IOException e) {
                e.printStackTrace();
            }
        }
    }
    
    private FileWriter ensureFile(String profilerName) {
        synchronized (this) {
            if (directory == null || directory.isEmpty()) {
                try {
                    directory = Files.createTempDirectory("jvm_profiler_").toString();
                } catch (IOException e) {
                    e.printStackTrace();
                }
            }
        }

        return fileWriters.computeIfAbsent(profilerName, t -> createFileWriter(t));
    }
    
    private FileWriter createFileWriter(String profilerName) {
        String path = Paths.get(directory, profilerName + ".json").toString();
        try {
            return new FileWriter(path, true);
        } catch (IOException e) {
            throw new RuntimeException("Failed to create file writer: " + path, e);
        }
    }
}


KafkaOutputReporter

  • 依赖kafka-client的jar包来构建KafkaProducer。
  • 通过producer.send来发送采集数据。
public class KafkaOutputReporter implements Reporter {
    private String brokerList = "localhost:9092";
    private boolean syncMode = false;
    
    private String topicPrefix;
    
    private ConcurrentHashMap<String, String> profilerTopics = new ConcurrentHashMap<>();

    private Producer<String, byte[]> producer;

    public KafkaOutputReporter() {
    }
    
    public KafkaOutputReporter(String brokerList, boolean syncMode, String topicPrefix) {
        this.brokerList = brokerList;
        this.syncMode = syncMode;
        this.topicPrefix = topicPrefix;
    }

    @Override
    public void report(String profilerName, Map<String, Object> metrics) {
        ensureProducer();

        String topicName = getTopic(profilerName);
        
        String str = JsonUtils.serialize(metrics);
        byte[] message = str.getBytes(StandardCharsets.UTF_8);

        Future<RecordMetadata> future = producer.send(
                new ProducerRecord<String, byte[]>(topicName, message));

        if (syncMode) {
            producer.flush();
            try {
                future.get();
            } catch (InterruptedException | ExecutionException e) {
                throw new RuntimeException(e);
            }
        }
    }

   // 省略一些非核心的代码
    private void ensureProducer() {
        synchronized (this) {
            if (producer != null) {
                return;
            }

            Properties props = new Properties();
            props.put("bootstrap.servers", brokerList);
            props.put("retries", 10);
            props.put("batch.size", 16384);
            props.put("linger.ms", 0);
            props.put("buffer.memory", 16384000);
            props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
            props.put("value.serializer", org.apache.kafka.common.serialization.ByteArraySerializer.class.getName());

            if (syncMode) {
                props.put("acks", "all");
            }

            producer = new KafkaProducer<>(props);
        }
    }
}


RedisOutputReporter

  • 依赖jedis包来实现redis的读写。
  • redis当中存储的采集数据的key是机器ip和时间戳的组合,value是采集的数据。
public class RedisOutputReporter implements Reporter {

    private static final AgentLogger logger = AgentLogger.getLogger(RedisOutputReporter.class.getName());
    private JedisPool redisConn = null;

    //JedisPool should always be used as it is thread safe.
    public void report(String profilerName, Map<String, Object> metrics) {
        ensureJedisConn();
        try {
            Jedis jedisClient = redisConn.getResource();
            jedisClient.set(createOriginStamp(profilerName), JsonUtils.serialize(metrics));
            redisConn.returnResource(jedisClient);
        } catch (Exception err) {
            logger.warn(err.toString());
        }
    }

    public String createOriginStamp(String profilerName) {
        try {
            return (profilerName + "-" + InetAddress.getLocalHost().getHostAddress() + "-" + System.currentTimeMillis());
        } catch (UnknownHostException err) {
            logger.warn("Address could not be determined and will be omitted!");
            return (profilerName + "-" + System.currentTimeMillis());
        }
    }

    public void close() {
        synchronized (this) {
            redisConn.close();
            redisConn = null;
        }
    }

    private void ensureJedisConn() {
        synchronized (this) {
            if (redisConn == null || redisConn.isClosed()) {
                redisConn = new JedisPool(System.getenv("JEDIS_PROFILER_CONNECTION"));
                return;
            }
        }
    }
}
目录
相关文章
|
10月前
|
存储 缓存 Java
我们来说一说 JVM 的内存模型
我是小假 期待与你的下一次相遇 ~
612 5
|
10月前
|
存储 缓存 算法
深入理解JVM《JVM内存区域详解 - 世界的基石》
Java代码从编译到执行需经javac编译为.class字节码,再由JVM加载运行。JVM内存分为线程私有(程序计数器、虚拟机栈、本地方法栈)和线程共享(堆、方法区)区域,其中堆是GC主战场,方法区在JDK 8+演变为使用本地内存的元空间,直接内存则用于提升NIO性能,但可能引发OOM。
|
Arthas 存储 算法
深入理解JVM,包含字节码文件,内存结构,垃圾回收,类的声明周期,类加载器
JVM全称是Java Virtual Machine-Java虚拟机JVM作用:本质上是一个运行在计算机上的程序,职责是运行Java字节码文件,编译为机器码交由计算机运行类的生命周期概述:类的生命周期描述了一个类加载,使用,卸载的整个过类的生命周期阶段:类的声明周期主要分为五个阶段:加载->连接->初始化->使用->卸载,其中连接中分为三个小阶段验证->准备->解析类加载器的定义:JVM提供类加载器给Java程序去获取类和接口字节码数据类加载器的作用:类加载器接受字节码文件。
1093 55
|
Arthas 监控 Java
Arthas memory(查看 JVM 内存信息)
Arthas memory(查看 JVM 内存信息)
1085 6
|
缓存 监控 算法
JVM简介—2.垃圾回收器和内存分配策略
本文介绍了Java垃圾回收机制的多个方面,包括垃圾回收概述、对象存活判断、引用类型介绍、垃圾收集算法、垃圾收集器设计、具体垃圾回收器详情、Stop The World现象、内存分配与回收策略、新生代配置演示、内存泄漏和溢出问题以及JDK提供的相关工具。
JVM简介—2.垃圾回收器和内存分配策略
|
存储 缓存 算法
JVM简介—1.Java内存区域
本文详细介绍了Java虚拟机运行时数据区的各个方面,包括其定义、类型(如程序计数器、Java虚拟机栈、本地方法栈、Java堆、方法区和直接内存)及其作用。文中还探讨了各版本内存区域的变化、直接内存的使用、从线程角度分析Java内存区域、堆与栈的区别、对象创建步骤、对象内存布局及访问定位,并通过实例说明了常见内存溢出问题的原因和表现形式。这些内容帮助开发者深入理解Java内存管理机制,优化应用程序性能并解决潜在的内存问题。
819 29
JVM简介—1.Java内存区域
|
存储 设计模式 监控
如何快速定位并优化CPU 与 JVM 内存性能瓶颈?
如何快速定位并优化CPU 与 JVM 内存性能瓶颈?
429 0
如何快速定位并优化CPU 与 JVM 内存性能瓶颈?
|
存储 算法 Java
JVM: 内存、类与垃圾
分代收集算法将内存分为新生代和老年代,分别使用不同的垃圾回收算法。新生代对象使用复制算法,老年代对象使用标记-清除或标记-整理算法。
275 6
|
存储 设计模式 监控
快速定位并优化CPU 与 JVM 内存性能瓶颈
本文介绍了 Java 应用常见的 CPU & JVM 内存热点原因及优化思路。
1423 166
|
存储 Java 程序员
【JVM】——JVM运行机制、类加载机制、内存划分
JVM运行机制,堆栈,程序计数器,元数据区,JVM加载机制,双亲委派模型
499 10