Przesyłanie danych kolumnowych za pomocą Apache Arrow

Tłumaczenie artykułu przygotowano specjalnie dla studentów kursu Data Engineer.

Przesyłanie danych kolumnowych za pomocą Apache Arrow

W ciągu ostatnich kilku tygodni dodaliśmy Nong Li do Apache Arrow binarne strumieniowe format, uzupełniając już istniejący format plików random access/IP. Mamy realizacje w Java i C++ oraz powiązania do Pythona. W tym artykule wyjaśnię, jak działa ten format i pokażę, jak można osiągnąć bardzo wysoką przepustowość danych dla DataFrame pandas.

Strumieniowe przesyłanie danych kolumnowych

Powszechnym pytaniem, które otrzymuję od użytkowników Arrow, jest kwestia wysokiego kosztu przenoszenia dużych zbiorów danych tabelarycznych z formatu opartego na wierszach lub rekordach do formatu kolumnowego. Dla multigigabajtowych zbiorów danych transponowanie w pamięci lub na dysku może okazać się niewykonalnym zadaniem.

Aby przesyłać dane strumieniowo, niezależnie od tego, czy dane źródłowe są w formacie wierszowym, czy kolumnowym, jedną z opcji pozostaje wysyłanie małych pakietów wierszy, z których każdy zawiera wewnętrznie układ kolumnowy.

W Apache Arrow kolekcja kolumnowych tablic w pamięci, reprezentująca kawałek tabeli, nazywana jest pakietem rekordów (record batch). Aby przedstawić jednolitą strukturę danych logicznej tabeli, można połączyć kilka pakietów rekordów.

W istniejącym formacie plików 'random access' zapisujemy metadane, zawierające schemat tabeli i lokalizację bloków na końcu pliku, co umożliwia bardzo tanie wybieranie dowolnego pakietu rekordów lub dowolnej kolumny z zestawu danych. W formacie strumieniowym wysyłamy serię wiadomości: schemat, a następnie jeden lub kilka pakietów rekordów.

Różne formaty wyglądają mniej więcej tak, jak przedstawione na tym rysunku:

Przesyłanie danych kolumnowych za pomocą Apache Arrow

Strumieniowe przesyłanie danych w PyArrow: zastosowanie

Aby pokazać, jak to działa, stworzę przykład zbioru danych, reprezentującego jeden strumieniowy kawałek:

import time
import numpy as np
import pandas as pd
import pyarrow as pa

def generate_data(total_size, ncols):
    nrows = int(total_size / ncols / np.dtype('float64').itemsize)
    return pd.DataFrame({
        'c' + str(i): np.random.randn(nrows)
        for i in range(ncols)
    })	

Teraz, załóżmy, że chcemy zapisać 1 Gb danych, składających się z kawałków o wielkości 1 Mb każdy, w sumie 1024 kawałki. Najpierw utwórzmy pierwszy DataFrame o rozmiarze 1 Mb z 16 kolumnami:

KILOBYTE = 1 << 10
MEGABYTE = KILOBYTE * KILOBYTE
DATA_SIZE = 1024 * MEGABYTE
NCOLS = 16

df = generate_data(MEGABYTE, NCOLS)

Następnie konwertuję je na pyarrow.RecordBatch:

batch = pa.RecordBatch.from_pandas(df)

Teraz stworzę strumień wyjściowy, który będzie pisał do pamięci operacyjnej i utworzy StreamWriter:

sink = pa.InMemoryOutputStream()
stream_writer = pa.StreamWriter(sink, batch.schema)

Następnie zapiszemy 1024 fragmenty, które w sumie będą miały 1GB zbioru danych:

for i in range(DATA_SIZE // MEGABYTE):
    stream_writer.write_batch(batch)

Ponieważ pisaliśmy w RAM, cały strumień będziemy mogli uzyskać w jednym buforze:

In [13]: source = sink.get_result()

In [14]: source
Out[14]: 

In [15]: source.size
Out[15]: 1074750744

Ponieważ te dane znajdują się w pamięci, odczytywanie pakietów z zapisami Arrow jest operacją zero-copy. Otwieram StreamReader, odczytuję dane do pyarrow.Table, a następnie konwertuję je na DataFrame pandas:

In [16]: reader = pa.StreamReader(source)

In [17]: table = reader.read_all()

In [18]: table
Out[18]: 

In [19]: df = table.to_pandas()

In [20]: df.memory_usage().sum()
Out[20]: 1073741904

To wszystko jest oczywiście dobre, ale mogą się pojawić pytania. Jak szybko to się dzieje? Jak rozmiar fragmentu wpływa na wydajność uzyskiwania DataFrame pandas?

Wydajność przesyłania danych

W miarę zmniejszania rozmiaru fragmentu koszt rekonstrukcji ciągłego ramki kolumnowej DataFrame w pandas rośnie z powodu nieefektywnych schematów dostępu do pamięci podręcznej. Istnieją również pewne koszty związane z pracą ze strukturami danych C++ oraz tablicami i ich buforami pamięci.

Dla 1 MB, jak podano powyżej, na moim laptopie (Quad-core Xeon E3-1505M) uzyskuję:

In [20]: %timeit pa.StreamReader(source).read_all().to_pandas()
10 loops, best of 3: 129 ms per loop

Efektywna przepustowość wynosi 7.75 GB/s przy odtwarzaniu DataFrame o objętości 1GB z 1024 fragmentów po 1MB. Co się stanie, jeśli użyjemy większych lub mniejszych fragmentów? Oto takie wyniki:

Przesyłanie danych kolumnowych za pomocą Apache Arrow

Wydajność znacznie spada przy przejściu z fragmentów 256K do 64K. Zaskoczyło mnie, że fragmenty o rozmiarze 1 MB były przetwarzane szybciej niż 16 MB. Warto przeprowadzić dokładniejsze badania i zrozumieć, czy to normalny rozkład, czy wpływa na to coś innego.

W obecnej implementacji formatu dane w zasadzie nie są kompresowane, więc rozmiar w pamięci i 'w przewodach' jest mniej więcej taki sam. W przyszłości kompresja może stać się dodatkową opcją.

Podsumowanie

Transmisja danych kolumnowych może być skutecznym sposobem przesyłania dużych zestawów danych do narzędzi analitycznych kolumnowych, takich jak pandas, w małych kawałkach. Usługi danych korzystające z magazynu opartego na wierszach mogą przesyłać i transponować małe kawałki danych, które są bardziej przyjazne dla pamięci podręcznej L2 i L3 twojego procesora.

Pełny kod

import time
import numpy as np
import pandas as pd
import pyarrow as pa

def generate_data(total_size, ncols):
    nrows = total_size / ncols / np.dtype('float64').itemsize
    return pd.DataFrame({
        'c' + str(i): np.random.randn(nrows)
        for i in range(ncols)
    })

KILOBYTE = 1 << 10
MEGABYTE = KILOBYTE * KILOBYTE
DATA_SIZE = 1024 * MEGABYTE
NCOLS = 16

def get_timing(f, niter):
    start = time.clock_gettime(time.CLOCK_REALTIME)
    for i in range(niter):
        f()
    return (time.clock_gettime(time.CLOCK_REALTIME) - start) / NITER

def read_as_dataframe(klass, source):
    reader = klass(source)
    table = reader.read_all()
    return table.to_pandas()
NITER = 5
results = []

CHUNKSIZES = [16 * KILOBYTE, 64 * KILOBYTE, 256 * KILOBYTE, MEGABYTE, 16 * MEGABYTE]

for chunksize in CHUNKSIZES:
    nchunks = DATA_SIZE // chunksize
    batch = pa.RecordBatch.from_pandas(generate_data(chunksize, NCOLS))

    sink = pa.InMemoryOutputStream()
    stream_writer = pa.StreamWriter(sink, batch.schema)

    for i in range(nchunks):
        stream_writer.write_batch(batch)

    source = sink.get_result()

    elapsed = get_timing(lambda: read_as_dataframe(pa.StreamReader, source), NITER)

    result = (chunksize, elapsed)
    print(result)
    results.append(result)

Źródło: habr.com

Kup solidny hosting stron z ochroną przed DDoS, serwery VPS VDS 🔥 Kup solidny hosting stron z ochroną przed DDoS, serwery VPS VDS | ProHoster