Artikli tõlge on ette valmistatud spetsiaalselt kursuse üliõpilastele .

Viimase nädala jooksul oleme lisanud binaarse voogedastusformaadi, täiustades juba olemasolevat failiformaati random access/IPC. Meil on rakendused Java ja C++ keeltes ning Pythonile sidemed. Selles artiklis selgitan, kuidas formaat töötab ja näitan, kuidas saavutada väga kõrget andmesidet pandas DataFrame'i jaoks.
Veergude andmete voogedastus
Üks sagedane küsimus, mida saan Arrowi kasutajatelt, on suurte tabeliandmete kogumite ülekandmise kõrge hinna küsimus ridade või rekordipõhisest vormingust veergude formaati. Mitme gigabaidi suuruste andmestike korral võib mälus või kettal transponeerimine osutuda ülemäära keeruliseks.
Andmete voogedastuse jaoks, olenemata sellest, kas algandmed on ridade või veergude vormingus, on üks variant saata väikseid rida pakette, millest igaühes on veergudeks jaotatud sisu.
Apache Arrowis nimetatakse mälus olevat veergude massiivide kogu, mis esindab tabeli plokki, rekordipakiks (record batch). Ühe loogilise tabelistruktuuri kujundamiseks saab koguda mitu rekordipakki.
Olemasolevas random access failiformaatis salvestame metaandmed, mis sisaldavad tabeli skeemi ja plokkide asukohta faili lõpus, mis võimaldab teil väga odavalt valida mis tahes rekordipaki või mis tahes veeru andmestikust. Voogedastusformaadis edastame sõnumite seeria: skeemi, seejärel ühe või mitu rekordipakki.
Erinevad formaadid näevad välja umbes nii, nagu on kujutatud sellel joonisel:

Andmete voogedastus PyArrow's: rakendamine
Kuidas see töötab, näitamaks, loon näidisandmestiku, mis esindab ühte voogedastusplokki:
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 eeldame, et soovime salvestada 1 GB andmeid, mis koosneb igast 1 MB suurusest plokist, kokku 1024 plokki. Kõigepealt loome esimese 1 MB suuruse andmeraami, millel 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 väljundvoo, mis kirjutab RAM-i ja loon StreamWriter:
sink = pa.InMemoryOutputStream()
stream_writer = pa.StreamWriter(sink, batch.schema)Seejärel kirjutame 1024 plokki, mis kokku annavad 1 GB andmestikku:
for i in range(DATA_SIZE // MEGABYTE):
stream_writer.write_batch(batch)Kuna oleme kirjutanud RAM-i, saame kogu voog ühe puhvrina:
In [13]: source = sink.get_result()
In [14]: source
Out[14]:
In [15]: source.size
Out[15]: 1074750744 Kuna need andmed asuvad mälus, on rekordipakettide lugemine Arrowis zero-copy operatsioon. Avan StreamReaderi, 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 tore, kuid teil võivad tekkida küsimusi. Kui kiiresti see toimub? Kuidas ploki suurus mõjutab DataFrame'i saamise jõudlust?
Andmete voogedastuse jõudlus
Ploki suuruse vähenedes suureneb DataFrame'i pideva veergude kadreerimise rekonstrueerimise hind pandas'is mäluliidese ebaefektiivsete skeemide tõttu. Samuti on mõningaid kulusid C++ andmestruktuuride ja massiivide ning nende mälu puhvritega töötamisel.
1 MB puhul, nagu eespool nimetatud, olen ma oma sülearvutis (Quad-core Xeon E3-1505M) saanud:
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 DataFrame 1024 plokist, igaüks 1 MB. Mis juhtub, kui kasutame suuremaid või väiksemaid plokke? Sellised on tulemused:

Tõhusus langeb oluliselt 256K-st 64K plokkideni. Mind üllatas, et 1 MB plokke töödeldi kiiremini kui 16 MB. Tuleks teha põhjalikum uurimus ja välja selgitada, kas see on normaalne jaotumine või mõjutab midagi muud.
Praeguses formaadi rakenduses andmeid ei tihendata, seega on mälu ja 'juhtme' suurus peaaegu sama. Tulevikus võib pakkuda tihendamist alternatiivina.
Kokkuvõte
Veergraafide voogudefaktid võivad olla tõhus viis suurte andmekoguste edastamiseks veergude analüüsitööriistadesse, nagu pandas, väikeste tükkidena. Andmete teenuseid, mis kasutavad ridadele suunatud salvestust, saavad edastada ja transponeerida väikseid andmepartiisid, mis on L2 ja L3 vahemälu jaoks paremad.
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
