Streaming di dati di colonne con Apache Arrow

La traduzione dell'articolo è stata preparata appositamente per gli studenti del corso «Data Engineer».

Streaming di dati di colonne con Apache Arrow

Negli ultimi settimane, siamo stati con Nong Li aggiunto al Apache Arrow formato di streaming binario, completando il già esistente formato di file ad accesso casuale/IP. Abbiamo implementazioni in Java e C++ e binding per Python. In questo articolo spiegherò come funziona il formato e mostrerò come raggiungere un'elevata capacità di trasmissione dati per DataFrame pandas.

Streaming di dati colonnari

Una domanda comune che ricevo dagli utenti di Arrow è quella riguardante l'alto costo del trasferimento di grandi set di dati tabellari da un formato basato su righe o record a un formato colonnare. Per dataset di diversi gigabyte, la trasposizione in memoria o su disco può rivelarsi un compito difficile.

Per lo streaming di dati, che siano di origine basata su righe o colonnare, una delle opzioni resta l'invio di piccoli pacchetti di righe, ognuno dei quali all'interno contiene una disposizione per colonne.

In Apache Arrow, a collection of in-memory columnar arrays representing a chunk of a table is called a record batch. To represent a unified data structure for a logical table, multiple record batches can be assembled.

In an existing 'random access' file format, we write metadata containing the table schema and the block locations at the end of the file, allowing you to very cheaply select any record batch or any column from the dataset. In the streaming format, we send a series of messages: the schema, followed by one or more record batches.

Different formats look approximately like this:

Streaming di dati di colonne con Apache Arrow

Data streaming in PyArrow: application

To show you how this works, I will create an example dataset representing a single streaming chunk:

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)
    })	

Supponiamo ora di voler registrare 1 GB di dati, composti da chunk di dimensione 1 MB ciascuno, per un totale di 1024 chunk. Iniziamo creando il primo frame di dati delle dimensioni di 1 MB con 16 colonne:

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

df = generate_data(MEGABYTE, NCOLS)

Poi li convertirò in pyarrow.RecordBatch:

batch = pa.RecordBatch.from_pandas(df)

Ora creerò uno stream di output che scriverà in memoria e creerò StreamWriter:

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

Poi scriveremo 1024 chunk che costituiranno in totale 1 GB di dataset:

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

Poiché stiamo scrivendo in RAM, intero lo stream potrà essere ottenuto in un unico buffer:

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

In [14]: source
Out[14]: 

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

Poiché questi dati si trovano in memoria, la lettura dei pacchetti di record Arrow risulta un'operazione zero-copy. Apro StreamReader e leggo i dati in pyarrow.Table, e poi li converto in 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

Tutto ciò è sicuramente positivo, ma potrebbero sorgere domande. Quanto è veloce questo processo? In che modo la dimensione del chunk influisce sulle prestazioni di ottenimento del DataFrame pandas?

Prestazioni dello streaming dei dati

Man mano che la dimensione del chunk diminuisce, il costo di ricostruzione di un frame a colonna continuo del DataFrame in pandas aumenta a causa di schemi di accesso alla cache inefficienti. Ci sono anche alcuni sovraccarichi dovuti all'interazione con strutture dati C++ e i loro buffer di memoria.

Per 1 MB, come indicato sopra, sul mio laptop (Quad-core Xeon E3-1505M) otteniamo:

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

Pertanto, la larghezza di banda effettiva è di 7.75 Gb/s per il recupero di un DataFrame di 1 Gb da 1024 chunk da 1 MB. Cosa succede se utilizziamo chunk di dimensioni maggiori o minori? I risultati saranno i seguenti:

Streaming di dati di colonne con Apache Arrow

Le prestazioni diminuiscono notevolmente da 256K a 64K chunk. Sono rimasto sorpreso che i chunk di dimensioni 1 MB siano stati elaborati più velocemente rispetto a quelli di 16 MB. È necessario condurre un'indagine più approfondita per capire se si tratta di una distribuzione normale o se c'è qualche altro fattore in gioco.

Nell'attuale implementazione del formato, i dati non vengono compressi, quindi la dimensione in memoria e 'nelle linee' è sostanzialmente la stessa. In futuro, la compressione potrebbe diventare un'opzione aggiuntiva.

Risultato

Lo streaming di dati columnari può risultare un modo efficace per trasferire grandi set di dati verso strumenti analitici colonne, come pandas, utilizzando piccoli chunk. I servizi di dati che utilizzano uno storage orientato alle righe possono trasferire e trasporre piccoli chunk di dati più adatti alla cache L2 e L3 del tuo processore.

Codice completo

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)

Fonte: habr.com

Acquista hosting affidabile per siti web con protezione DDoS, server VPS VDS 🔥 Acquista hosting affidabile per siti web con protezione DDoS, server VPS VDS | ProHoster