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

Viimase kuu jooksul oleme me lisanud 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:

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]: 1073741904KĂ”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 loopSeega 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:

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
