最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

詳解大數(shù)據(jù)處理引擎Flink內(nèi)存管理

 更新時間:2021年05月20日 09:08:17   作者:華為云開發(fā)者社區(qū)  
Flink是jvm之上的大數(shù)據(jù)處理引擎,jvm存在java對象存儲密度低、full gc時消耗性能,gc存在stw的問題,同時omm時會影響穩(wěn)定性。針對頻繁序列化和反序列化問題flink使用堆內(nèi)堆外內(nèi)存可以直接在一些場景下操作二進制數(shù)據(jù),減少序列化反序列化消耗。本文帶你詳細理解其原理。

內(nèi)存模型

Flink可以使用堆內(nèi)和堆外內(nèi)存,內(nèi)存模型如圖所示:

flink使用內(nèi)存劃分為堆內(nèi)內(nèi)存和堆外內(nèi)存。按照用途可以劃分為task所用內(nèi)存,network memory、managed memory、以及framework所用內(nèi)存,其中task network managed所用內(nèi)存計入slot內(nèi)存。framework為taskmanager公用。

堆內(nèi)內(nèi)存包含用戶代碼所用內(nèi)存、heapstatebackend、框架執(zhí)行所用內(nèi)存。

堆外內(nèi)存是未經(jīng)jvm虛擬化的內(nèi)存,直接映射到操作系統(tǒng)的內(nèi)存地址,堆外內(nèi)存包含框架執(zhí)行所用內(nèi)存,jvm堆外內(nèi)存、Direct、native等。

Direct memory內(nèi)存可用于網(wǎng)絡(luò)傳輸緩沖。network memory屬于direct memory的范疇,flink可以借助于此進行zero copy,從而減少內(nèi)核態(tài)到用戶態(tài)copy次數(shù),從而進行更高效的io操作。

jvm metaspace存放jvm加載的類的元數(shù)據(jù),加載的類越多,需要的空間越大,overhead用于jvm的其他開銷,如native memory、code cache、thread stack等。

Managed Memory主要用于RocksDBStateBackend和批處理算子,也屬于native memory的范疇,其中rocksdbstatebackend對應(yīng)rocksdb,rocksdb基于lsm數(shù)據(jù)結(jié)構(gòu)實現(xiàn),每個state對應(yīng)一個列族,占有獨立的writebuffer,rocksdb占用native內(nèi)存大小為 blockCahe + writebufferNum * writeBuffer + index ,同時堆外內(nèi)存是進程之間共享的,jvm虛擬化大量heap內(nèi)存耗時較久,使用堆外內(nèi)存的話可以有效的避免該環(huán)節(jié)。但堆外內(nèi)存也有一定的弊端,即監(jiān)控調(diào)試使用相對復雜,對于生命周期較短的segment使用堆內(nèi)內(nèi)存開銷更低,flink在一些情況下,直接操作二進制數(shù)據(jù),避免一些反序列化帶來的開銷。如果需要處理的數(shù)據(jù)超出了內(nèi)存限制,則會將部分數(shù)據(jù)存儲到硬盤上。

內(nèi)存管理

類似于OS中的page機制,flink模擬了操作系統(tǒng)的機制,通過page來管理內(nèi)存,flink對應(yīng)page的數(shù)據(jù)結(jié)構(gòu)為dataview和MemorySegment,memorysegment是flink內(nèi)存分配的最小單位,默認32kb,其可以在堆上也可以在堆外,flink通過MemorySegment的數(shù)據(jù)結(jié)構(gòu)來訪問堆內(nèi)堆外內(nèi)存,借助于flink序列化機制(序列化機制會在下一小節(jié)講解),memorysegment提供了對二進制數(shù)據(jù)的讀取和寫入的方法,flink使用datainputview和dataoutputview進行memorysegment的二進制的讀取和寫入,flink可以通過HeapMemorySegment 管理堆內(nèi)內(nèi)存,通過HybridMemorySegment來管理堆內(nèi)和堆外內(nèi)存,MemorySegment管理jvm堆內(nèi)存時,其定義一個字節(jié)數(shù)組的引用指向內(nèi)存端,基于該內(nèi)部字節(jié)數(shù)組的引用進行操作的HeapMemorySegment。

public abstract class MemorySegment {
    /**
     * The heap byte array object relative to which we access the memory.
     *  如果為堆內(nèi)存,則指向訪問的內(nèi)存的引用,否則若內(nèi)存為非堆內(nèi)存,則為null
     * <p>Is non-<tt>null</tt> if the memory is on the heap, and is <tt>null</tt>, if the memory is
     * off the heap. If we have this buffer, we must never void this reference, or the memory
     * segment will point to undefined addresses outside the heap and may in out-of-order execution
     * cases cause segmentation faults.
     */
    protected final byte[] heapMemory;
    /**
     * The address to the data, relative to the heap memory byte array. If the heap memory byte
     * array is <tt>null</tt>, this becomes an absolute memory address outside the heap.
     * 字節(jié)數(shù)組對應(yīng)的相對地址
     */
    protected long address;  
}

HeapMemorySegment用來分配堆上內(nèi)存。

public final class HeapMemorySegment extends MemorySegment {
    /**
     * An extra reference to the heap memory, so we can let byte array checks fail by the built-in
     * checks automatically without extra checks.
     * 字節(jié)數(shù)組的引用指向該內(nèi)存段
     */
    private byte[] memory;
    public void free() {
        super.free();
        this.memory = null;
    }
 
    public final void get(DataOutput out, int offset, int length) throws IOException {
        out.write(this.memory, offset, length);
    }
}

HybridMemorySegment即支持onheap和offheap內(nèi)存,flink通過jvm的unsafe操作,如果對象o不為null,為onheap的場景,并且后面的地址或者位置是相對位置,那么會直接對當前對象(比如數(shù)組)的相對位置進行操作。如果對象o為null,操作的內(nèi)存塊不是JVM堆內(nèi)存,為off-heap的場景,并且后面的地址是某個內(nèi)存塊的絕對地址,那么這些方法的調(diào)用也相當于對該內(nèi)存塊進行操作。

public final class HybridMemorySegment extends MemorySegment {
  @Override
    public ByteBuffer wrap(int offset, int length) {
        if (address <= addressLimit) {
            if (heapMemory != null) {
                return ByteBuffer.wrap(heapMemory, offset, length);
            }
            else {
                try {
                    ByteBuffer wrapper = offHeapBuffer.duplicate();
                    wrapper.limit(offset + length);
                    wrapper.position(offset);
                    return wrapper;
                }
                catch (IllegalArgumentException e) {
                    throw new IndexOutOfBoundsException();
                }
            }
        }
        else {
            throw new IllegalStateException("segment has been freed");
        }
    }
}

flink通過MemorySegmentFactory來創(chuàng)建memorySegment,memorySegment是flink內(nèi)存分配的最小單位。對于跨memorysegment的數(shù)據(jù)方位,flink抽象出一個訪問視圖,數(shù)據(jù)讀取datainputView,數(shù)據(jù)寫入dataoutputview。

/**
 * This interface defines a view over some memory that can be used to sequentially read the contents of the memory.
 * The view is typically backed by one or more {@link org.apache.flink.core.memory.MemorySegment}.
 */
@Public
public interface DataInputView extends DataInput {
private MemorySegment[] memorySegments; // view持有的MemorySegment的引用, 該組memorysegment可以視為一個內(nèi)存頁,
flink可以順序讀取memorysegmet中的數(shù)據(jù)
/**
     * Reads up to {@code len} bytes of memory and stores it into {@code b} starting at offset {@code off}.
     * It returns the number of read bytes or -1 if there is no more data left.
     * @param b byte array to store the data to
     * @param off offset into byte array
     * @param len byte length to read
     * @return the number of actually read bytes of -1 if there is no more data left
     */
    int read(byte[] b, int off, int len) throws IOException;
}

dataoutputview是數(shù)據(jù)寫入的視圖,outputview持有多個memorysegment的引用,flink可以順序的寫入segment。

/**
 * This interface defines a view over some memory that can be used to sequentially write contents to the memory.
 * The view is typically backed by one or more {@link org.apache.flink.core.memory.MemorySegment}.
 */
@Public
public interface DataOutputView extends DataOutput {
private final List<MemorySegment> memory; // memorysegment的引用
/**
     * Copies {@code numBytes} bytes from the source to this view.
     * @param source The source to copy the bytes from.
     * @param numBytes The number of bytes to copy.
    void write(DataInputView source, int numBytes) throws IOException;
}

上一小節(jié)中講到的managedmemory內(nèi)存部分,flink使用memorymanager來管理該內(nèi)存,managedmemory只使用堆外內(nèi)存,主要用于批處理中的sorting、hashing、以及caching(社區(qū)消息,未來流處理也會使用到該部分),在流計算中作為rocksdbstatebackend的部分內(nèi)存。memeorymanager通過memorypool來管理memorysegment。

/**
 * The memory manager governs the memory that Flink uses for sorting, hashing, caching or off-heap state backends
 * (e.g. RocksDB). Memory is represented either in {@link MemorySegment}s of equal size or in reserved chunks of certain
 * size. Operators allocate the memory either by requesting a number of memory segments or by reserving chunks.
 * Any allocated memory has to be released to be reused later.
 * <p>The memory segments are represented as off-heap unsafe memory regions (both via {@link HybridMemorySegment}).
 * Releasing a memory segment will make it re-claimable by the garbage collector, but does not necessarily immediately
 * releases the underlying memory.
 */
public class MemoryManager {
 /**
     * Allocates a set of memory segments from this memory manager.
     * <p>The total allocated memory will not exceed its size limit, announced in the constructor.
     * @param owner The owner to associate with the memory segment, for the fallback release.
     * @param target The list into which to put the allocated memory pages.
     * @param numberOfPages The number of pages to allocate.
     * @throws MemoryAllocationException Thrown, if this memory manager does not have the requested amount
     *                                   of memory pages any more.
     */
    public void allocatePages(
            Object owner,
            Collection<MemorySegment> target,
            int numberOfPages) throws MemoryAllocationException {
}

private static void freeSegment(MemorySegment segment, @Nullable Collection<MemorySegment> segments) {
        segment.free();
        if (segments != null) {
            segments.remove(segment);
        }
    }
/**
     * Frees this memory segment.
     * <p>After this operation has been called, no further operations are possible on the memory
     * segment and will fail. The actual memory (heap or off-heap) will only be released after this
     * memory segment object has become garbage collected.
     */
    public void free() {
        // this ensures we can place no more data and trigger
        // the checks for the freed segment
        address = addressLimit + 1;
    }
}

對于上一小節(jié)中提到的NetWorkMemory的內(nèi)存,flink使用networkbuffer做了一層buffer封裝。buffer的底層也是memorysegment,flink通過bufferpool來管理buffer,每個taskmanager都有一個netwokbufferpool,該tm上的各個task共享該networkbufferpool,同時task對應(yīng)的localbufferpool所需的內(nèi)存需要從networkbufferpool申請而來,它們都是flink申請的堆外內(nèi)存。

上游算子向resultpartition寫入數(shù)據(jù)時,申請buffer資源,使用bufferbuilder將數(shù)據(jù)寫入memorysegment,下游算子從resultsubpartition消費數(shù)據(jù)時,利用bufferconsumer從memorysegment中讀取數(shù)據(jù),bufferbuilder與bufferconsumer一一對應(yīng)。同時這一流程也和flink的反壓機制相關(guān)。如圖

/**
 * A buffer pool used to manage a number of {@link Buffer} instances from the
 * {@link NetworkBufferPool}.
 * <p>Buffer requests are mediated to the network buffer pool to ensure dead-lock
 * free operation of the network stack by limiting the number of buffers per
 * local buffer pool. It also implements the default mechanism for buffer
 * recycling, which ensures that every buffer is ultimately returned to the
 * network buffer pool.
 * <p>The size of this pool can be dynamically changed at runtime ({@link #setNumBuffers(int)}. It
 * will then lazily return the required number of buffers to the {@link NetworkBufferPool} to
 * match its new size.
 */
class LocalBufferPool implements BufferPool {
@Nullable
    private MemorySegment requestMemorySegment(int targetChannel) throws IOException {
        MemorySegment segment = null;
        synchronized (availableMemorySegments) {
            returnExcessMemorySegments();

            if (availableMemorySegments.isEmpty()) {
                segment = requestMemorySegmentFromGlobal();
            }
            // segment may have been released by buffer pool owner
            if (segment == null) {
                segment = availableMemorySegments.poll();
            }
            if (segment == null) {
                availabilityHelper.resetUnavailable();
            }
            if (segment != null && targetChannel != UNKNOWN_CHANNEL) {
                if (subpartitionBuffersCount[targetChannel]++ == maxBuffersPerChannel) {
                    unavailableSubpartitionsCount++;
                    availabilityHelper.resetUnavailable();
                }
            }
        }
        return segment;
    }
    }
    /**
 * A result partition for data produced by a single task.
 *
 * <p>This class is the runtime part of a logical {@link IntermediateResultPartition}. Essentially,
 * a result partition is a collection of {@link Buffer} instances. The buffers are organized in one
 * or more {@link ResultSubpartition} instances, which further partition the data depending on the
 * number of consuming tasks and the data {@link DistributionPattern}.
 * <p>Tasks, which consume a result partition have to request one of its subpartitions. The request
 * happens either remotely (see {@link RemoteInputChannel}) or locally (see {@link LocalInputChannel})
  The life-cycle of each result partition has three (possibly overlapping) phases:
    Produce  Consume  Release  Buffer management State management
 */
public abstract class ResultPartition implements ResultPartitionWriter, BufferPoolOwner {
      @Override
    public BufferBuilder getBufferBuilder(int targetChannel) throws IOException, InterruptedException {
        checkInProduceState();
        return bufferPool.requestBufferBuilderBlocking(targetChannel);
    }
    }
}

自定義序列化框架

flink對自身支持的基本數(shù)據(jù)類型,實現(xiàn)了定制的序列化機制,flink數(shù)據(jù)集對象相對固定,可以只保存一份schema信息,從而節(jié)省存儲空間,數(shù)據(jù)序列化就是java對象和二進制數(shù)據(jù)之間的數(shù)據(jù)轉(zhuǎn)換,flink使用TypeInformation的createSerializer接口負責創(chuàng)建每種類型的序列化器,進行數(shù)據(jù)的序列化反序列化,類型信息在構(gòu)建streamtransformation時通過typeextractor根據(jù)方法簽名類信息等提取類型信息并存儲在streamconfig中。

/**
     * Creates a serializer for the type. The serializer may use the ExecutionConfig
     * for parameterization.
     * 創(chuàng)建出對應(yīng)類型的序列化器
     * @param config The config used to parameterize the serializer.
     * @return A serializer for this type.
     */
    @PublicEvolving
    public abstract TypeSerializer<T> createSerializer(ExecutionConfig config);
/**
 * A utility for reflection analysis on classes, to determine the return type of implementations of transformation
 * functions.
 */
@Public
public class TypeExtractor {
/**
 * Creates a {@link TypeInformation} from the given parameters.
     * If the given {@code instance} implements {@link ResultTypeQueryable}, its information
     * is used to determine the type information. Otherwise, the type information is derived
     * based on the given class information.
     * @param instance            instance to determine type information for
     * @param baseClass            base class of {@code instance}
     * @param clazz                class of {@code instance}
     * @param returnParamPos    index of the return type in the type arguments of {@code clazz}
     * @param <OUT>                output type
     * @return type information
     */
    @SuppressWarnings("unchecked")
    @PublicEvolving
    public static <OUT> TypeInformation<OUT> createTypeInfo(Object instance, Class<?> baseClass, Class<?> clazz,
 int returnParamPos) {
        if (instance instanceof ResultTypeQueryable) {
            return ((ResultTypeQueryable<OUT>) instance).getProducedType();
        } else {
            return createTypeInfo(baseClass, clazz, returnParamPos, null, null);
        }
    }
}

對于嵌套的數(shù)據(jù)類型,flink從最內(nèi)層的字段開始序列化,內(nèi)層序列化的結(jié)果將組成外層序列化結(jié)果,反序列時,從內(nèi)存中順序讀取二進制數(shù)據(jù),根據(jù)偏移量反序列化為java對象。flink自帶序列化機制存儲密度很高,序列化對應(yīng)的類型值即可。

flink中的table模塊在memorysegment的基礎(chǔ)上使用了BinaryRow的數(shù)據(jù)結(jié)構(gòu),可以更好地減少反序列化開銷,需要反序列化是可以只序列化相應(yīng)的字段,而無需序列化整個對象。

同時你也可以注冊子類型和自定義序列化器,對于flink無法序列化的類型,會交給kryo進行處理,如果kryo也無法處理,將強制使用avro來序列化,kryo序列化性能相對flink自帶序列化機制較低,開發(fā)時可以使用env.getConfig().disableGenericTypes()來禁用kryo,盡量使用flink框架自帶的序列化器對應(yīng)的數(shù)據(jù)類型。

緩存友好的數(shù)據(jù)結(jié)構(gòu)

cpu中L1、L2、L3的緩存讀取速度比從內(nèi)存中讀取數(shù)據(jù)快很多,高速緩存的訪問速度是主存的訪問速度的很多倍。另外一個重要的程序特性是局部性原理,程序常常使用它們最近使用的數(shù)據(jù)和指令,其中兩種局部性類型,時間局部性指最近訪問的內(nèi)容很可能短期內(nèi)被再次訪問,空間局部性是指地址相互臨近的項目很可能短時間內(nèi)被再次訪問。

結(jié)合這兩個特性設(shè)計緩存友好的數(shù)據(jù)結(jié)構(gòu)可以有效的提升緩存命中率和本地化特性,該特性主要用于排序操作中,常規(guī)情況下一個指針指向一個<key,v>對象,排序時需要根據(jù)指針pointer獲取到實際數(shù)據(jù),然后再進行比較,這個環(huán)節(jié)涉及到內(nèi)存的隨機訪問,緩存本地化會很低,使用序列化的定長key + pointer,這樣key就會連續(xù)存儲到內(nèi)存中,避免的內(nèi)存的隨機訪問,還可以提升cpu緩存命中率。對兩條記錄進行排序時首先比較key,如果大小不同直接返回結(jié)果,只需交換指針即可,不用交換實際數(shù)據(jù),如果相同,則比較指針實際指向的數(shù)據(jù)。

以上就是詳解大數(shù)據(jù)處理引擎Flink內(nèi)存管理的詳細內(nèi)容,更多關(guān)于大數(shù)據(jù)處理引擎Flink內(nèi)存管理的資料請關(guān)注腳本之家其它相關(guān)文章!

您可能感興趣的文章:

相關(guān)文章

  • SpringSecurity之SecurityContextHolder使用解讀

    SpringSecurity之SecurityContextHolder使用解讀

    這篇文章主要介紹了SpringSecurity之SecurityContextHolder使用解讀,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2023-03-03
  • 話說Spring Security權(quán)限管理(源碼詳解)

    話說Spring Security權(quán)限管理(源碼詳解)

    本篇文章主要介紹了話說Spring Security權(quán)限管理(源碼詳解) ,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2017-02-02
  • Java關(guān)于MyBatis緩存詳解

    Java關(guān)于MyBatis緩存詳解

    緩存的重要性是不言而喻的,使用緩存,我們可以避免頻繁的與數(shù)據(jù)庫進行交互,尤其是在查詢越多、緩存命中率越高的情況下,使用緩存對性能的提高更明顯。本文將給大家詳細的介紹,對大家的學習或工作具有一定的參考借鑒價值
    2021-09-09
  • 基于Spring AMQP實現(xiàn)消息隊列的示例代碼

    基于Spring AMQP實現(xiàn)消息隊列的示例代碼

    Spring AMQP作為Spring框架的一部分,是一套用于支持高級消息隊列協(xié)議(AMQP)的工具,AMQP是一種強大的消息協(xié)議,旨在支持可靠的消息傳遞,本文給大家介紹了如何基于Spring AMQP實現(xiàn)消息隊列,需要的朋友可以參考下
    2024-03-03
  • java8使用Stream API方法總結(jié)

    java8使用Stream API方法總結(jié)

    在本篇文章里小編給大家分享了關(guān)于java8使用Stream API方法相關(guān)知識點,需要的朋友們學習下。
    2019-04-04
  • idea SpringBoot+Gradle環(huán)境配置到項目打包

    idea SpringBoot+Gradle環(huán)境配置到項目打包

    Gradle是一個基于Java應(yīng)用的項目自動化構(gòu)建工具,本文介紹了在IDEA中創(chuàng)建Spring Boot Gradle項目,項目配置包括init.gradle和settings.gradle,感興趣的可以了解一下
    2024-11-11
  • Java 中的 Class.forName(類名) 使用及原理解析

    Java 中的 Class.forName(類名) 使用及原理解析

    Class.forName是Java中用于動態(tài)加載類的強大工具,廣泛應(yīng)用于數(shù)據(jù)庫驅(qū)動加載、反射機制和插件系統(tǒng)等場景,它通過ClassLoader加載類并執(zhí)行靜態(tài)初始化代碼,但在使用時需要注意類路徑、初始化副作用和類加載器的選擇等問題,感興趣的朋友一起看看吧
    2024-12-12
  • JavaEE Cookie的基本使用細節(jié)

    JavaEE Cookie的基本使用細節(jié)

    本章我們將學習會話跟蹤技術(shù)中的Cookie與Session,它在我們整個JavaEE的知識體系中是非常重要的,本節(jié)我們先介紹Cookie,廢話不多說,直接上正文
    2022-12-12
  • Java 程序員容易犯的10個SQL錯誤

    Java 程序員容易犯的10個SQL錯誤

    本文介紹了Java 程序員容易犯的10個SQL錯誤。具有很好的參考價值,下面跟著小編一起來看下吧
    2017-01-01
  • springboot + elasticsearch 實現(xiàn)聚合查詢的詳細代碼

    springboot + elasticsearch 實現(xiàn)聚合查詢的詳細代碼

    文章介紹了如何在Spring Boot 2.2.6中使用Elasticsearch進行聚合查詢,重點在于通過API創(chuàng)建索引和映射,而不是使用Spring Data Elasticsearch的自動創(chuàng)建功能,文章還提到在創(chuàng)建映射時,Elasticsearch會自動為keyword類型添加keyword屬性,感興趣的朋友一起看看吧
    2025-02-02

最新評論

高陵县| 南漳县| 长汀县| 盐城市| 清水河县| 五家渠市| 固始县| 乐都县| 舟山市| 广宁县| 普定县| 新闻| 尉氏县| 稻城县| 全南县| 蓝田县| 宜兰市| 绥江县| 志丹县| 稷山县| 来安县| 吴堡县| 新田县| 民县| 苍南县| 济宁市| 凤城市| 康乐县| 通渭县| 东阳市| 西贡区| 青河县| 保山市| 绥滨县| 闽清县| 晋城| 佛坪县| 涡阳县| 离岛区| 兖州市| 朔州市|