Streaming di dati da colonne utilizzando Apache Arrow

Traduzione dell'articolo preparata appositamente per gli studenti del corso «Data Engineer».

Streaming di dati da colonne utilizzando Apache Arrow

Negli ultimi settimane abbiamo Nong Li aggiunto al Apache Arrow formato di streaming binario, completando il già esistente formato di file per accesso casuale/IPC. Abbiamo implementazioni in Java e C++ e binding per Python. In questo articolo spiegherò come funziona il formato e mostrerò come è possibile raggiungere una capacità di trasmissione dati molto elevata per DataFrame pandas.

Streaming di dati colonnari

Una domanda comune che ricevo dagli utenti di Arrow è riguardo al costo elevato del trasferimento di grandi set di dati tabulari da un formato a righe o record a un formato colonnare. Per dataset di diverse gigabyte, la trasposizione in memoria o su disco può diventare un compito gravoso.

Per lo streaming di dati, indipendentemente dal fatto che i dati originali siano in formato righe o colonne, una delle opzioni è inviare piccoli pacchetti di righe, ognuno dei quali contiene una disposizione per colonne.

In Apache Arrow, una collezione di array colonnari in memoria, rappresentante un chunk di tabella, è chiamata pacchetto di registrazioni (record batch). Per rappresentare una singola struttura dati di una tabella logica è possibile assemblare più pacchetti di registrazioni.

Nel formato di file esistente per accesso casuale, scriviamo metadati contenenti lo schema della tabella e la posizione dei blocchi alla fine del file, consentendoti di selezionare estremamente economicamente qualsiasi pacchetto di registrazioni o qualsiasi colonna dal set di dati. Nel formato di streaming inviamo una serie di messaggi: lo schema, seguito da uno o più pacchetti di registrazioni.

I vari formati appaiono più o meno come mostrato in questa immagine:

Streaming di dati da colonne utilizzando Apache Arrow

Streaming di dati in PyArrow: applicazione

Per mostrarti come funziona, creerò un esempio di dataset rappresentante un chunk di streaming:

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

Ora, supponiamo di voler scrivere 1 GB di dati, costituiti da chunk delle dimensioni di 1 MB ciascuno, per un totale di 1024 chunk. Per iniziare, creiamo il primo DataFrame della dimensione 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à nella memoria e creerò StreamWriter:

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

Successivamente scriveremo 1024 chunk, che alla fine costituiranno 1 GB di dati:

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

Poiché abbiamo scritto in RAM, possiamo ottenere l'intero stream 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 si traduce in un'operazione zero-copy. Apro lo 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ò è ovviamente interessante, ma potreste porvi delle domande. Quanto velocemente avviene tutto questo? Come influisce la dimensione del chunk sulle prestazioni dell'ottenimento di DataFrame pandas?

Prestazioni del data streaming

Man mano che la dimensione del chunk diminuisce, il costo di ricostruzione di un frame a colonne continuo del DataFrame in pandas aumenta a causa di schemi di accesso alla memoria cache inefficaci. Ci sono anche alcune spese generali nella gestione delle strutture dati C++ e dei 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

Quindi, la larghezza di banda effettiva è di 7.75 GB/s per recuperare il DataFrame da 1 GB da 1024 chunk da 1 MB. Cosa succede se utilizziamo chunk di dimensioni maggiori o minori? Ecco quali risultati otteniamo:

Streaming di dati da colonne utilizzando Apache Arrow

Le prestazioni diminuiscono significativamente passando da chunk da 256K a 64K. Sono rimasto sorpreso che chunk da 1 MB venissero elaborati più velocemente rispetto a quelli da 16 MB. Vale la pena fare un'analisi più approfondita e capire se si tratta di una distribuzione normale o se c'è qualcos'altro che influisce.

Nell'attuale implementazione del formato, i dati non sono compressi, quindi la dimensione in memoria e "in transito" è praticamente la stessa. In futuro, la compressione potrebbe diventare un'opzione aggiuntiva.

Risultato

La trasmissione di dati a colonne può rivelarsi un metodo efficace per trasferire grandi set di dati negli strumenti di analisi a colonne, come pandas, utilizzando piccoli chunk. I servizi dati che utilizzano un'archiviazione orientata alle righe possono trasmettere e trasporre piccoli chunk di dati, che sono più adatti per la 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, VPS VDS server 🔥 Acquista hosting affidabile per siti web con protezione DDoS, VPS VDS server | ProHoster