MemorySegment

MemorySegment 是 Paimon 项目中用于内存管理的核心类。可以把它理解为一个功能强大且高效的内存块的“包装器”。它屏蔽了底层内存是在 JVM 堆内(On-Heap)还是堆外(Off-Heap)的差异,为上层应用提供了统一、高效的内存访问接口。

下面我们从几个方面来详细解析这个类:

核心设计思想

MemorySegment 的主要目标是提供高性能的、可预测的内存操作。

  • 统一内存模型:无论是堆内内存(由 byte[] 数组表示)还是堆外内存(由 java.nio.ByteBuffer.allocateDirect 分配),MemorySegment 都提供了相同的 API 进行读写。这大大简化了上层数据结构(如哈希表、排序器等)的设计,它们无需关心内存的具体来源。
  • 极致的性能:该类大量使用了 sun.misc.Unsafe API。Unsafe 允许 Java 代码像 C/C++ 一样直接操作内存地址,绕过了 JVM 的一些安全检查(如数组边界检查),从而获得了极高的访问速度。这对于大数据处理中频繁的序列化/反序列化和数据拷贝操作至关重要。
  • 精细的控制:它提供了对多字节数据类型(如 int, long)的字节序(Endianness)的精细控制,这在需要跨平台或与网络协议交互时非常重要。

MemorySegment 内部通过几个关键字段来管理内存:

// ... existing code ...
    public static final sun.misc.Unsafe UNSAFE = MemoryUtils.UNSAFE;

    public static final long BYTE_ARRAY_BASE_OFFSET = UNSAFE.arrayBaseOffset(byte[].class);

    public static final boolean LITTLE_ENDIAN =
            (ByteOrder.nativeOrder() == ByteOrder.LITTLE_ENDIAN);

    @Nullable private byte[] heapMemory;

    @Nullable private ByteBuffer offHeapBuffer;

    private long address;

    private final int size;
// ... existing code ...
  • UNSAFE: sun.misc.Unsafe 的实例,是实现高性能操作的核心工具。
  • BYTE_ARRAY_BASE_OFFSET: 一个静态常量,表示 Java byte[] 数组中第一个元素的起始地址相对于数组对象头的偏移量。Unsafe 操作堆内存时必须使用这个偏移量。
  • LITTLE_ENDIAN: 判断当前操作系统的原生字节序。true 表示小端序(Little Endian),false 表示大端序(Big Endian)。
  • heapMemory: 如果 MemorySegment 管理的是堆内内存,这个字段会引用一个 byte[] 数组。否则为 null。
  • offHeapBuffer: 如果管理的是堆外内存,这个字段会引用一个 Direct ByteBuffer。否则为 null。这个引用主要是为了防止堆外内存在 MemorySegment 仍在使用时被 GC 回收。
  • address: 核心字段。它代表了内存块的“基地址”。
    • 对于堆内内存,address 的值就是 BYTE_ARRAY_BASE_OFFSET。
    • 对于堆外内存,address 是通过 getByteBufferAddress(buffer) 获取的实际内存地址。
  • size: 表示这个内存块的总字节数。

通过反射绕开语言限制,获取Unsafe

private static sun.misc.Unsafe getUnsafe() {
        try {
            Field unsafeField = sun.misc.Unsafe.class.getDeclaredField("theUnsafe");
            unsafeField.setAccessible(true);
            return (sun.misc.Unsafe) unsafeField.get(null);

对象的创建

MemorySegment 的构造函数是私有的,强制用户通过静态工厂方法来创建实例,这使得创建逻辑更清晰。


// ... existing code ...
    public static MemorySegment wrap(byte[] buffer) {
        return new MemorySegment(buffer, null, BYTE_ARRAY_BASE_OFFSET, buffer.length);
    }

    public static MemorySegment wrapOffHeapMemory(ByteBuffer buffer) {
        return new MemorySegment(null, buffer, getByteBufferAddress(buffer), buffer.capacity());
    }

    public static MemorySegment allocateHeapMemory(int size) {
        return wrap(new byte[size]);
    }

    public static MemorySegment allocateOffHeapMemory(int size) {
        return wrapOffHeapMemory(ByteBuffer.allocateDirect(size));
    }
// ... existing code ...
  • allocateHeapMemory(int size): 在 JVM 堆上分配指定大小的 byte[] 数组,并用 MemorySegment 包装。
  • allocateOffHeapMemory(int size): 分配指定大小的堆外内存(Direct Buffer),并用 MemorySegment 包装。
  • wrap(byte[] buffer): 包装一个已存在的 byte[] 数组。
  • wrapOffHeapMemory(ByteBuffer buffer): 包装一个已存在的堆外 ByteBuffer。

原始类型读写 (Getters/Setters)

MemorySegment 为所有 Java 基本类型都提供了 get 和 put 方法。

// ... existing code ...
    public byte get(int index) {
        return UNSAFE.getByte(heapMemory, address + index);
    }

    public void put(int index, byte b) {
        UNSAFE.putByte(heapMemory, address + index, b);
    }
// ... existing code ...
    public int getInt(int index) {
        return UNSAFE.getInt(heapMemory, address + index);
    }
// ... existing code ...
    public void putInt(int index, int value) {
        UNSAFE.putInt(heapMemory, address + index, value);
    }
// ... existing code ...

以 getInt 为例,UNSAFE.getInt(object, offset) 方法:

  • 第一个参数 object:如果是堆内内存,它就是 heapMemory 这个 byte[] 对象;如果是堆外内存,它必须是 null。MemorySegment 的设计巧妙地利用了这一点,heapMemory 字段在堆外模式下本身就是 null。
  • 第二个参数 offset:这是内存的绝对地址或相对偏移量。
    • 对于堆内内存,它是 address (即 BYTE_ARRAY_BASE_OFFSET) + index。
    • 对于堆外内存,它是 address (即堆外内存的起始地址) + index。

通过这种方式,getInt(index) 这一行代码就能够同时兼容堆内和堆外两种情况。

Unsafe操作模式

模式一:操作“某个对象内部”的内存(相对地址模式)

  • 触发条件:当你给它的第一个参数 base 传入一个具体的 Java 对象时(比如 byte[] 数组)。
  • 工作方式:它会把第二个参数 offset 当作是相对于这个对象起始位置的偏移量。
  • 对应场景:MemorySegment 的堆内内存(On-Heap)。
    • heapMemory 是一个 byte[] 实例(不是 null)。
    • address 是 BYTE_ARRAY_BASE_OFFSET(一个固定的偏移量)。
    • UNSAFE.getByte(heapMemory, address + index) 的含义是:“请在 heapMemory 这个对象内部,从它的基础偏移量 address 开始,再前进 index 个字节,然后读取那里的数据。”

模式二:操作“任意绝对地址”的内存(绝对地址模式)

  • 触发条件:当你给它的第一个参数 base 传入 null 时。
  • 工作方式:null 在这里是一个特殊信号。Unsafe 看到 null 后,就不会去尝试操作一个对象了。它会完全忽略第一个参数,并把第二个参数 offset 当作一个**操作系统内存中的绝对地址(一个具体的内存门牌号)**来直接操作。
  • 对应场景:MemorySegment 的堆外内存(Off-Heap)。
    • heapMemory 被特意设置为了 null。
    • address 是通过 getByteBufferAddress() 获取到的一个巨大的数字,它就是堆外内存的绝对起始地址。
    • UNSAFE.getByte(heapMemory, address + index) 的含义是:“我收到的第一个参数是 null,切换到绝对地址模式!现在,请直接把 address + index 这个数字当作一个内存地址,去那个地址上读取一个字节的数据。”

此外,它还提供了处理字节序的方法:

// ... existing code ...
    public int getIntLittleEndian(int index) {
        if (LITTLE_ENDIAN) {
            return getInt(index);
        } else {
            return Integer.reverseBytes(getInt(index));
        }
    }

    public int getIntBigEndian(int index) {
        if (LITTLE_ENDIAN) {
            return Integer.reverseBytes(getInt(index));
        } else {
            return getInt(index);
        }
    }
// ... existing code ...

这些方法确保了在不同架构的机器上也能正确地读写数据。

批量数据读写 (Bulk Read/Write)

对于大数据块的复制,逐字节操作效率很低。MemorySegment 提供了多种批量操作方法,它们都基于 UNSAFE.copyMemory,性能极高。

// ... existing code ...
    public void get(int index, byte[] dst, int offset, int length) {
        // check the byte array offset and length and the status
        if ((offset | length | (offset + length) | (dst.length - (offset + length))) < 0) {
            throw new IndexOutOfBoundsException();
        }

        UNSAFE.copyMemory(
                heapMemory, address + index, dst, BYTE_ARRAY_BASE_OFFSET + offset, length);
    }

    public void put(int index, byte[] src, int offset, int length) {
// ... existing code ...
        UNSAFE.copyMemory(
                src, BYTE_ARRAY_BASE_OFFSET + offset, heapMemory, address + index, length);
    }
// ... existing code ...
    public void copyTo(int offset, MemorySegment target, int targetOffset, int numBytes) {
        final byte[] thisHeapRef = this.heapMemory;
        final byte[] otherHeapRef = target.heapMemory;
        final long thisPointer = this.address + offset;
        final long otherPointer = target.address + targetOffset;

        UNSAFE.copyMemory(thisHeapRef, thisPointer, otherHeapRef, otherPointer, numBytes);
    }
// ... existing code ...
  • get(int, byte[], ...) / put(int, byte[], ...): 在 MemorySegment 和 byte[] 数组之间批量复制数据。
  • copyTo(int, MemorySegment, ...): 在两个 MemorySegment 之间直接复制内存,这是最高效的方式,因为它可能完全在堆外内存之间或者堆内内存之间进行,无需用户空间和内核空间的数据拷贝。

比较与交换

MemorySegment 还提供了高效的内存区域比较和交换方法。

// ... existing code ...
    public int compare(MemorySegment seg2, int offset1, int offset2, int len) {
        while (len >= 8) {
            long l1 = this.getLongBigEndian(offset1);
            long l2 = seg2.getLongBigEndian(offset2);

            if (l1 != l2) {
                return (l1 < l2) ^ (l1 < 0) ^ (l2 < 0) ? -1 : 1;
            }

            offset1 += 8;
            offset2 += 8;
            len -= 8;
        }
        while (len > 0) {
            int b1 = this.get(offset1) & 0xff;
            int b2 = seg2.get(offset2) & 0xff;
            int cmp = b1 - b2;
            if (cmp != 0) {
                return cmp;
            }
            offset1++;
            offset2++;
            len--;
        }
// ... existing code ...
        return 0;
    }

    public boolean equalTo(MemorySegment seg2, int offset1, int offset2, int length) {
        int i = 0;

        // we assume unaligned accesses are supported.
        // Compare 8 bytes at a time.
        while (i <= length - 8) {
            if (getLong(offset1 + i) != seg2.getLong(offset2 + i)) {
                return false;
            }
            i += 8;
        }
// ... existing code ...
        return true;
    }
// ... existing code ...

compare 和 equalTo 方法都进行了优化:它们尽可能一次性比较8个字节(一个 long),而不是逐字节比较,这大大提升了比较的效率。

总结

MemorySegment 是一个精心设计的底层内存操作组件。它通过 Unsafe 实现了与 C 语言级别相近的内存访问性能,同时通过优雅的封装,为上层应用提供了一个安全、易用、统一的接口来处理堆内和堆外内存。它是 Paimon 高性能数据处理能力的重要基石。

MemorySize

MemorySize 是 Paimon 中一个非常重要的基础工具类。它的核心目标是为系统提供一个统一、类型安全的方式来表示和处理内存大小。在大型数据处理系统中,内存相关的配置项(如缓冲区大小、缓存大小等)非常普遍,用户通常喜欢使用 "128MB"、"2G" 这样的可读字符串来配置,而程序内部需要将它们精确地转换为字节(bytes)进行计算。MemorySize 正是连接用户配置和程序内部计算的桥梁。

MemorySize 的设计遵循了几个关键原则:

  • 不可变性 (Immutability): 类的核心数据 private final long bytes; 是 final 的,一旦 MemorySize 对象被创建,其代表的字节数就不能被改变。这使得它在多线程环境下是安全的,并且行为可预测。
  • 内部统一表示: 无论用户输入的是 "1GB" 还是 "1024MB",在 MemorySize 内部都统一存储为 long 类型的字节数。这是所有计算和比较的基准。
  • 丰富的 API: 提供了从字符串解析、不同单位的转换、算术运算到人性化字符串输出的全套 API。

核心属性与构造方法

// ... existing code ...
public class MemorySize implements java.io.Serializable, Comparable<MemorySize> {

    // ... existing code ...
    public static final MemorySize VALUE_128_MB = MemorySize.ofMebiBytes(128);

    public static final MemorySize VALUE_256_MB = MemorySize.ofMebiBytes(256);

// ... existing code ...
    /** The memory size, in bytes. */
    private final long bytes;

// ... existing code ...
    /**
     * Constructs a new MemorySize.
     *
     * @param bytes The size, in bytes. Must be zero or larger.
     */
    public MemorySize(long bytes) {
        checkArgument(bytes >= 0, "bytes must be >= 0");
        this.bytes = bytes;
    }

    public static MemorySize ofMebiBytes(long mebiBytes) {
        return new MemorySize(mebiBytes << 20);
    }

    public static MemorySize ofKibiBytes(long kibiBytes) {
        return new MemorySize(kibiBytes << 10);
    }
// ... existing code ...
  • bytes 字段: 这是类的基石,所有内存大小的最终表示。
  • 构造器 MemorySize(long bytes): 这是唯一的构造器,接受一个 long 型的字节数。它通过 checkArgument 确保了传入的字节数必须大于等于0。
  • 静态工厂方法 ofMebiBytes, ofKibiBytes 等: 提供了便捷的、类型安全的方式来创建 MemorySize 对象。这里使用了位移运算(<< 20 代表乘以 2^20,即 1MB),这比乘法效率更高,也是处理2的幂次单位的常用技巧。
  • 预定义常量 ZERO, VALUE_128_MB 等: 为常用值提供了静态常量,避免了重复创建对象的开销,也提高了代码的可读性。

字符串解析 (Parsing)

这是 MemorySize 最核心和最常用的功能之一,它负责将字符串配置(如 "100m", "2g")转换为 MemorySize 对象。

// ... existing code ...
    public static MemorySize parse(String text) throws IllegalArgumentException {
        return new MemorySize(parseBytes(text));
    }
// ... existing code ...
    public static long parseBytes(String text) throws IllegalArgumentException {
        checkNotNull(text, "text");

        final String trimmed = text.trim();
        checkArgument(!trimmed.isEmpty(), "argument is an empty- or whitespace-only string");

        final int len = trimmed.length();
        int pos = 0;

        char current;
        while (pos < len && (current = trimmed.charAt(pos)) >= '0' && current <= '9') {
            pos++;
        }

        final String number = trimmed.substring(0, pos);
        final String unit = trimmed.substring(pos).trim().toLowerCase(Locale.US);

        if (number.isEmpty()) {
            throw new NumberFormatException("text does not start with a number");
        }

        final long value = Long.parseLong(number); // ...

        final long multiplier = parseUnit(unit).map(MemoryUnit::getMultiplier).orElse(1L);
        final long result = value * multiplier;

        // check for overflow
        if (result / multiplier != value) {
            // ...
        }

        return result;
    }
// ... existing code ...

解析过程 parseBytes 非常严谨:

  1. 预处理: 去除首尾空格。
  2. 分离数值和单位: 从头开始遍历字符串,找到第一个非数字字符的位置,从而将字符串分割成 "数值部分" 和 "单位部分"。
  3. 解析数值: 将 "数值部分" 解析为 long。
  4. 解析单位: 将 "单位部分"(转为小写)传入 parseUnit 方法,该方法会匹配 MemoryUnit 枚举中定义的单位别名,找到对应的乘数(multiplier)。如果单位部分为空,则默认乘数为1(即单位是 bytes)。
  5. 计算结果: 将数值和乘数相乘。
  6. 溢出检查: if (result / multiplier != value) 这是一个非常巧妙的溢出检查。如果 value * multiplier 的结果超出了 long 的最大值,那么 result / multiplier 的结果将不再等于原始的 value。

MemoryUnit 枚举

// ... existing code ...
    public enum MemoryUnit {
        BYTES(new String[] {"b", "bytes"}, 1L),
        KILO_BYTES(new String[] {"k", "kb", "kibibytes"}, 1024L),
        MEGA_BYTES(new String[] {"m", "mb", "mebibytes"}, 1024L * 1024L),
        GIGA_BYTES(new String[] {"g", "gb", "gibibytes"}, 1024L * 1024L * 1024L),
        TERA_BYTES(new String[] {"t", "tb", "tebibytes"}, 1024L * 1024L * 1024L * 1024L);

        private final String[] units;

        private final long multiplier;
// ... existing code ...

这个内部枚举是解析逻辑的核心。它定义了每种内存单位的多种文本别名(如 "g", "gb", "gibibytes" 都代表 GIGA_BYTES)以及对应的字节乘数。值得注意的是,这里的单位都是基于1024的二进制单位(KiB, MiB, GiB),这是计算机科学中内存计算的标准。

单位获取与格式化输出

  • 单位获取:

    // ... existing code ...
        /** Gets the memory size in Mebibytes (= 1024 Kibibytes). */
        public int getMebiBytes() {
            return (int) (bytes >> 20);
        }
    // ... existing code ...
    

    getKibiBytes(), getMebiBytes() 等方法使用位移运算(>>)来快速获取不同单位下的值。例如 >> 20 就是除以 2^20 (1 MiB)。getMebiBytes() 返回 int 是一个值得注意的细节,因为在很多框架(如 Flink)的API中,MB 级别的内存大小通常用 int 表示。

  • 格式化输出:

    • toString(): 生成一个规范化的字符串。它会寻找能整除总字节数的最大单位进行显示。例如,new MemorySize(2048) 会输出 "2 kb"。
    • toHumanReadableString(): 生成一个更易读的字符串,它会使用最接近的单位,并以浮点数显示,同时附上原始字节数。例如 new MemorySize(1500000) 可能会输出 "1.431mb (1500000 bytes)"。

算术运算

// ... existing code ...
    public MemorySize add(MemorySize that) {
        return new MemorySize(Math.addExact(this.bytes, that.bytes));
    }

    public MemorySize subtract(MemorySize that) {
        return new MemorySize(Math.subtractExact(this.bytes, that.bytes));
    }
// ... existing code ...

add, subtract 等方法提供了类型安全的算术运算。它们返回一个新的 MemorySize 对象(符合不可变性原则),并使用 Math.addExact 等方法,这些方法在发生溢出时会直接抛出 ArithmeticException,增强了代码的健壮性。

总结

MemorySize 是一个设计精良、健壮且功能完备的工具类。它通过将内存大小这一概念封装成一个独立类型,极大地简化了在系统中处理内存配置的复杂性,有效地防止了因单位混淆、手动计算错误或数据溢出等问题导致的 bug。对于任何需要处理用户自定义内存配置的系统来说,它都是一个典范级的实现。

Buffer 类:内存中的数据容器

Buffer 是 Paimon 内存管理中的一个基础且核心的组件。它并非直接管理内存,而是作为 MemorySegment 的一个轻量级包装,为其赋予了“已用大小”的概念。

@Public
public class Buffer {

    private final MemorySegment segment;

    private int size;

    public Buffer(MemorySegment segment, int size) {
        this.segment = segment;
        this.size = size;
    }

    public static Buffer create(MemorySegment segment) {
        return create(segment, 0);
    }

    public static Buffer create(MemorySegment segment, int size) {
        return new Buffer(segment, size);
    }

    public MemorySegment getMemorySegment() {
        return segment;
    }

    public int getSize() {
        return size;
    }

    public int getMaxCapacity() {
        return segment.size();
    }

    public ByteBuffer getNioBuffer(int index, int length) {
        return segment.wrap(index, length).slice();
    }

    public void setSize(int size) {
        this.size = size;
    }
}

核心设计与属性

  • private final MemorySegment segment;: 这是 Buffer 的基石。MemorySegment 代表了一块连续的、原始的内存(可以是堆内或堆外内存)。Buffer 本身不分配内存,而是引用一个已经分配好的 MemorySegment。这个引用是 final 的,意味着一个 Buffer 对象生命周期内都与同一个 MemorySegment 绑定。
  • private int size;: 这是 Buffer 设计的精髓所在。它是一个可变字段,用于追踪 segment 中当前被认为有效数据的字节数。这与 segment 本身的大小(即最大容量)是两个不同的概念。

关键方法分析

  • getMaxCapacity() vs getSize():
    • getMaxCapacity() 返回底层 MemorySegment 的总大小,代表这个 Buffer 最多能容纳多少数据。
    • getSize() 返回当前 Buffer 中实际存放的数据大小。size 总是小于或等于 maxCapacity。
  • setSize(int size): 提供了修改 size 的能力。这使得 Buffer 非常灵活,可以被复用。例如,一个组件可以从网络或文件中读取数据填充到 Buffer 中,然后调用 setSize() 来更新有效数据的大小。
  • getNioBuffer(...): 提供了与 Java NIO ByteBuffer 的互操作性。它允许你从底层 MemorySegment 的任意位置(index)开始,获取一个指定长度(length)的 ByteBuffer 视图,方便使用 NIO 的 API 进行读写。
  • 可变性: 与我们之前分析的 MemorySize 不同,Buffer 是一个可变对象(因为 size 可变)。这符合它的定位——一个可被填充、消费和重置的数据容器。

应用场景

在 Paimon 中,Buffer 无处不在。例如,当从磁盘文件中读取一个数据块时,会使用 Buffer 作为承载数据的目的地:

// ... existing code ...
    public boolean readBufferFromFileChannel(Buffer buffer) throws IOException {
// ... existing code ...
        int size = header.getInt();
        if (size > buffer.getMaxCapacity()) {
            throw new IllegalStateException(
                    "Buffer is too small for data: "
// ... existing code ...
        }
// ... existing code ...
        buffer.setSize(size);
        return false;
    }
// ... existing code ...

在这个例子里,readBufferFromFileChannel 方法接收一个 Buffer,从文件读取数据填充它,最后调用 buffer.setSize(size) 来标记这个 Buffer 现在包含了多少有效数据。

Logo

助力广东及东莞地区开发者,代码托管、在线学习与竞赛、技术交流与分享、资源共享、职业发展,成为松山湖开发者首选的工作与学习平台

更多推荐