Transmetimi i të dhënave kolonë me Apache Arrow

Përkthimi i artikullit është përgatitur posaçërisht për studentët e kursit «Inxhiner i Të Dhënave».

Transmetimi i të dhënave kolonë me Apache Arrow

Në javët e fundit ne kemi Nong Li shtuar në Apache Arrow formatin binar të transmetimit, duke plotësuar formatin ekzistues të skedareve random access/IPC. Ne kemi zbatime në Java dhe C++ dhe lidhje Python. Në këtë artikull do t'ju tregoj se si funksionon formati dhe do të tregoj se si mund të arrijmë një kapacitet shumë të lartë transmetimi të të dhënave për DataFrame pandas.

Transmetimi i të dhënave kolonore

Një pyetje e zakonshme që marr nga përdoruesit e Arrow është rreth kostos së lartë për të transferuar grupe të mëdha të dhënash tabelare nga një format të orientuar në rreshta ose regjistrime në një format kolonor. Për datasetet shumëgigabajt, transponimi në memorie ose në disk mund të jetë një detyrë e papërballueshme.

Për transmetimin e të dhënave, pavarësisht nëse të dhënat fillestare janë në format rreshtash ose kolonash, një nga opsionet mbetet dërgimi i paketeve të vogla rreshtash, secila e cila brenda përmban një kompozim sipas kolonave.

Në Apache Arrow, koleksioni i array-ev kolonor në memorie, që përfaqëson një çank të tabelës, quhet paketë regjistrash (record batch). Për të paraqitur një strukturë të vetme të të dhënave të tabelës logjike, mund të mbledhim disa paketa regjistrash.

Në formatin ekzistues të skedareve «random access», ne shkruajmë metadata që përmban skemën e tabelës dhe vendndodhjen e bllokëve në fund të skedarit, që ju lejon të zgjidhni shumë lirë çdo paketë regjistrash ose çdo kolonë nga një grup të dhënash. Në formatin e transmetimit, ne dërgojmë një seri mesazhesh: skemën dhe pastaj një ose më shumë paketa regjistrash.

Formatet e ndryshme duken përafërsisht siç paraqitet në këtë figurë:

Transmetimi i të dhënave kolonë me Apache Arrow

Transmetimi i të dhënave në PyArrow: aplikimi

Për t'ju treguar se si funksionon, do të krijoj një shembull të një dataset-i që përfaqëson një çank transmitimi:

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

Tani, supozoni se duam të shkruajmë 1 GB të dhënash, që përbëhet nga çanka me përmasë 1 MB çdo, duke përfituar 1024 çanka gjithsej. Së pari, le të krijojmë frekën e parë të të dhënave me përmasë 1 MB me 16 kolona:

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

df = generate_data(MEGABYTE, NCOLS)

Pastaj unë do t'i konvertoj në pyarrow.RecordBatch:

batch = pa.RecordBatch.from_pandas(df)

Tani do të krijoj një rrjedhë dalëse që do të shkruaj në memorjen e përkohshme dhe do të krijoj StreamWriter:

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

Pastaj do të shkruajmë 1024 çanka, që në fund do të përbëjnë një grup të dhënash 1GB:

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

Duke qenë se po shkruanim në RAM, ne do të mund ta marrim tërë rrjedhën në një buffer:

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

In [14]: source
Out[14]: 

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

Duke qenë se këto të dhëna ndodhen në memorie, leximi i paketave të regjistrave Arrow rezulton një operacion zero-copy. Unë hap një StreamReader, lexoj të dhënat në pyarrow.Table, pastaj i konvertoj në 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

Të gjitha këto, sigurisht, janë të mira, por mund të keni pyetje. Sa shpejt ndodh kjo? Si ndikon përmasa e çankut në performancën e marrjes së DataFrame pandas?

Performanca e transmetimit të të dhënave

Ndërsa përmasa e çankut të transmetimit zvogëlohet, kostoja e rikonstruksionit të një kornizë të vazhdueshme kolonore të DataFrame në pandas rritet për shkak të skemave të papërshkueshme të aksesit në memorie. Ka gjithashtu disa kosto të mbulimeve nga puna me struktura të dhënash C++ dhe array të tyre dhe bufferat e memorie.

Për 1 MB, siç u përmend më herët, në kompjuterin tim (Quad-core Xeon E3-1505M) rezulton:

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

Rezultati Ă«shtĂ« se kapaciteti efektiv Ă«shtĂ« 7.75 GB/s pĂ«r rikthimin e njĂ« DataFrame 1GB nga 1024 çanka 1MB. ÇfarĂ« ndodh nĂ«se pĂ«rdorim çanka mĂ« tĂ« mĂ«dha ose mĂ« tĂ« vogla? KĂ«to janĂ« rezultatet qĂ« do tĂ« marrim:

Transmetimi i të dhënave kolonë me Apache Arrow

Performanca e zvogëlohet ndjeshëm nga çankarët 256K në 64K. Më befasoi se çankat me përmasë 1 MB përpunoheshin më shpejt se ato me 16 MB. Duhet bërë një hulumtim më të detajuar për të kuptuar nëse kjo është një shpërndarje normale apo ka diçka tjetër që ndikon.

Në realizimin aktual të formatit, të dhënat nuk kompresohen ashtu siç është, kështu që përmasat në memorie dhe «në tela» janë përafërsisht të njëjta. Në të ardhmen, kompresimi mund të bëhet një opsion shtesë.

Përfundimi

Transmetimi i të dhënave kolonë mund të jetë një mënyrë efektive për të transferuar grupe të mëdha të dhënash në mjetet analitike kolonë, për shembull në pandas, duke përdorur copa të vogla. Shërbimet e të dhënave që përdorin një magazinë të orientuar nga rreshti, mund të transmetojnë dhe transponojnë copa të vogla të dhënash, të cilat janë më të përshtatshme për caches L2 dhe L3 të procesorit tuaj.

Kodi i plotë

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)

Burimi: habr.com

Bleni hostim tĂ« besueshĂ«m pĂ«r faqe me mbrojtje nga DDoS, serverĂ« VPS VDS đŸ”„ Bleni hostim tĂ« besueshĂ«m pĂ«r faqe me mbrojtje nga DDoS, serverĂ« VPS VDS | ProHoster