Veergude andmete voogedastus Apache Arrow abil

Artikli tĂ”lge on ette valmistatud spetsiaalselt kursuse ĂŒliĂ”pilastele „Andmeinsener“.

Veergude andmete voogedastus Apache Arrow abil

Viimase kuu jooksul oleme me Nong Li lisanud Apache Arrow binaarse voogesitusvormingu, tÀiustades juba olemasolevat juhusliku juurdepÀÀsu/IP formaati. Meil on rakendused Java ja C++ keeles ning Python jaoks sidemed. Selles artiklis selgitan, kuidas see formaat töötab ja nÀitan, kuidas saavutada vÀga kÔrget andme edastuskiiruset DataFrame pandas'e jaoks.

Veergude andmete voogedastus

KĂŒsimus, mida Arrowi kasutajatelt tihti kuulen, on seotud suurte tabelandmete komplektide ĂŒleviimise kĂ”rgete kuludega ridakujulisest vĂ”i kirjekujulisest formaadist veergude formaati. Mitme gigabaidi suuruste andmehulkade korral vĂ”ib mĂ€lu vĂ”i ketta peal transponimine osutuda ĂŒletamatuks ĂŒlesandeks.

Andmete voogedastamiseks, sĂ”ltumata sellest, kas algandmed on ridakujulised vĂ”i veergude formaadis, on ĂŒks valik saata vĂ€ikeseid ridapakette, kus igaĂŒhes on veergude struktuur.

Apache Arrowis on mĂ€lus veergude massiivide kogum, mis esindab tabeli chunk'i, nimetatakse kirje pakkumiks (record batch). Üksikute andmestruktuuride esindamiseks saab kokku koguda mitmeid kirje pakkumeid.

Olemasolevas juhusliku juurdepÀÀsu failiformaadis kirjutame metainfot, mis sisaldab tabeli skeemi ja plokkide asukohta faili lĂ”ppu, vĂ”imaldades teil vĂ€ga odavalt valida mis tahes kirje paku vĂ”i mis tahes veeru andmestikust. Voogesitusformaadis saadame sĂ”numite seeria: skeemi ja seejĂ€rel ĂŒhe vĂ”i mitu kirje pakkumist.

Erinevad formaadid nÀevad vÀlja umbes nii, nagu see joonis nÀitab:

Veergude andmete voogedastus Apache Arrow abil

Andmete voogedastus PyArrow'is: rakendus

Kuidas nĂ€idata, kuidas see töötab, loodan luua nĂ€idisandmestiku, mis esindab ĂŒhte voogu:

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

NĂŒĂŒd oletame, et tahame kirjutada 1 GB andmeid, mis koosnevad 1 MB suurustest chunk'idest, kokku 1024 chunk'i. Esiteks loome 1 MB suuruse andmeraami, kus on 16 veergu:

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

df = generate_data(MEGABYTE, NCOLS)

SeejÀrel konverteerin need pyarrow.RecordBatch:

batch = pa.RecordBatch.from_pandas(df)

NĂŒĂŒd loon ma vĂ€ljundi voo, mis kirjutab mĂ€llu ja loon StreamWriter:

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

SeejÀrel kirjutame 1024 lÔiku, mis kokku moodustavad 1 GB andmekogumist:

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

Kuna kirjutasime RAM-i, saame kogu voog ĂŒhes puhvrisse kĂ€tte:

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

In [14]: source
Out[14]: 

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

Kuna need andmed on mÀlus, on Arrow'i kirje pakettide lugemine zero-copy operatsioon. Avan StreamReader'i, loen andmed pyarrow.Table, ja seejÀrel konverteerin need 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

KĂ”ik see on muidugi kena, kuid teil vĂ”ivad tekkida kĂŒsimused. Kui kiiresti see toimub? Kuidas lĂ”ike suurus mĂ”jutab pandas DataFrame'i saamise jĂ”udlust?

Andmevoo jÔudlus

LÔike suuruse vÀhenedes suureneb andmevoo pideva veergude struktuuri rekonstrueerimise maksumus pandas'is, mistÔttu on vahemÀlu juurdepÀÀsu skeemid ebaefektiivsed. Samuti on mÔningaid kulusid C++ andmestruktuuride ja massiivide ning nende mÀlupuhvriga töötamisel.

1 MB puhul, nagu eespool mainitud, on minu sĂŒlearvuti (Quad-core Xeon E3-1505M) puhul tulemuseks:

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

Seega on efektiivne lÀbilaskevÔime 7,75 GB/s, et taastada 1 GB suurune DataFrame 1024 1 MB lÔigust. Mis juhtub, kui kasutame suurema vÔi vÀiksema suurusega lÔike? Tulemused on:

Veergude andmete voogedastus Apache Arrow abil

JĂ”udlus vĂ€heneb mĂ€rkimisvÀÀrselt 256K kuni 64K lĂ”ikude puhul. Mind ĂŒllatas, et 1 MB suurused lĂ”igud töötati kiiremini kui 16 MB. Tasu on teha pĂ”hjalikum uurimistöö, et mĂ”ista, kas see on normaaljaotuste mĂ”ju vĂ”i on siin veel mĂ”ni tegur.

Praeguses formaadis ei pakita andmeid pĂ”himĂ”tteliselt, seega on mĂ€lus suurus ja ĂŒlekanal, kus andmed liiguvad, ligikaudu sama. Tulevikus vĂ”ib pakkimine saada lisavalikuks.

KokkuvÔte

Veergu pĂ”hinev andmete voogedastus vĂ”ib olla tĂ”hus viis suurte andmepakettide edastamiseks veergude analĂŒĂŒsi tööriistadesse, nĂ€iteks pandas, vĂ€ikeste blokkgeneraatorite kaudu. Andmete teenused, mis kasutavad ridadele suunatud salvestust, saavad edastada ja transponida vĂ€ikeseid andmeblokke, mis sobivad paremini teie protsessori L2 ja L3 vahemĂ€llu.

TĂ€ielik kood

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)

Allikas: habr.com

Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid đŸ”„ Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid | ProHoster