Skip to article frontmatterSkip to article content
Site not loading correctly?

This may be due to an incorrect BASE_URL configuration. See the MyST Documentation for reference.

Big Data & High-Performance Pipelines

Heinrich-Heine-Universität Düsseldorf

Datenverarbeitung jenseits des RAMs. Fokus: Zero-Copy, spaltenorientierte Formate und Lazy-Evaluation.

!uv venv --quiet
!uv sync --quiet
!uv pip --quiet install polars pyarrow pandas scikit-learn
  • Numpy 2.0 ABI: Achte auf Versionen: Pandas \ge 2.2.0, Matplotlib \ge 3.9.0, Scikit-Learn \ge 1.5.0.

  • Array-Alternativen: cupy für GPU-Beschleunigung oder jax.numpy für Auto-Differentiation.

Architektur-Prinzipien

  • Row- vs. Columnar: CSV/JSON (Row) speichert Datensätze zusammen – schlecht für analytische Abfragen. Parquet (Columnar) speichert Spalten zusammen – ideal für Aggregation und Kompression.

  • Zero-Copy (Apache Arrow): Standardisiertes Speicherformat. Verschiebt Daten zwischen Tools (Polars, PyTorch, Arrow) ohne Serialisierungs-Overhead.

  • Out-of-Core: Verarbeitung von Datensätzen, die den RAM übersteigen, durch Chunking (I/O-Iteratoren) oder Lazy-Evaluation.

Polars & Lazy Evaluation

Polars optimiert den gesamten Query-Plan (Predicate Pushdown, Projection Pushdown) vor der Ausführung. scan_parquet lädt nur Metadaten.

%%writefile polars_writer.py
import polars as pl
import numpy as np
import os

# Ensure data directory exists
os.makedirs("data", exist_ok=True)

# Generate a sample DataFrame
data = {
    "id": np.arange(1000),
    "value": np.random.randint(0, 200, 1000),
    "feature1": np.random.randn(1000),
    "feature2": np.random.randn(1000),
    "target": np.random.randint(0, 2, 1000)
}

df = pl.DataFrame(data)

# Save as Parquet for your scan_parquet tests
df.write_parquet("data/sample1.parquet")
df.write_parquet("data/sample2.parquet")

print("Dummy data created in ./data/ directory.")
Overwriting polars_writer.py
%%writefile polars_reader.py
import polars as pl
# Lazy: Kein I/O, nur Planerstellung
lazy_df = pl.scan_parquet("data/sample*.parquet") \
    .filter(pl.col("value") > 100) \
    .select(["id", "value"])

# Streaming: Verarbeitet in Chunks, RAM-sicher
df = lazy_df.collect(engine="streaming")
Writing polars_reader.py
!uv run polars_writer.py
!uv run polars_reader.py
Dummy data created in ./data/ directory.
import polars as pl
import matplotlib.pyplot as plt

hist_data = (
    pl.scan_parquet("data/sample*.parquet")
    .select(pl.col("target"))
    .collect(engine="streaming")
)

plt.figure(figsize=(10, 5))
plt.hist(hist_data["target"], bins=20, edgecolor="black")
plt.title("Distribution of 'target' from Parquet Files")
plt.xlabel("target")
plt.ylabel("Frequency")
plt.grid(axis="y", alpha=0.5)
plt.show()
<Figure size 1000x500 with 1 Axes>

ML-Pipeline-Integration

Scikit-Learn

Verwende partial_fit statt fit für Daten > RAM.

%%writefile sklearn_parquet.py
import polars as pl
from sklearn.linear_model import SGDClassifier

clf = SGDClassifier()
scanner = pl.scan_parquet("data/sample*.parquet")
for chunk in scanner.collect(engine="streaming").iter_slices():
    X = chunk.select(["feature1", "feature2"]).to_numpy()
    y = chunk.select(["target"]).to_numpy().ravel()
    clf.partial_fit(X, y, classes=[0, 1])
Writing sklearn_parquet.py
!uv run sklearn_parquet.py

PyTorch

Nutze torch.as_tensor() für Zero-Copy-Konvertierung von Arrow-Backends.

%%writefile torch_parquet.py
import torch
import polars as pl


class ParquetDataset(torch.utils.data.IterableDataset):
    def __init__(self, path, chunk_size=1000):
        self.path = path
        self.chunk_size = chunk_size

    def __iter__(self):
        df = pl.scan_parquet(self.path).collect(engine="streaming")
        for i in range(0, len(df), self.chunk_size):
            chunk = df.slice(i, self.chunk_size)
            features = chunk.select(["feature1", "feature2"])
            yield torch.as_tensor(features.to_numpy(), dtype=torch.float32)


dataset = ParquetDataset("data/sample*.parquet", chunk_size=1000)
loader = torch.utils.data.DataLoader(dataset, batch_size=None) # batch_size=None because iteration yields batches
model = torch.nn.Linear(2, 1)

# training loop
for batch in loader:
    output = model(batch.float())
    # ... proceed with loss calculation and optimization
Writing torch_parquet.py

Wir starten das Skript mit --with torch, damit torch nicht als projektweite Abhängigkeit deklariert werden muss (denn PyTorch is groß).

!uv run --with "torch" torch_parquet.py
Installed 29 packages in 231ms                                       

PyArrow

Man kann auch direkt mit PyArrow arbeiten (wobei wir hier sinnloserweise wieder in eine Parquet-Datei schreiben, wozu man den Umweg über Arrow nicht braucht).

%%writefile pyarrow_parquet.py
import polars as pl
import pyarrow as pa
import pyarrow.parquet as pq

lazy_df = pl.scan_parquet("data/sample*.parquet") \
    .filter(pl.col("value") > 100) \
    .select(["id", "value"])
df = lazy_df.collect(engine="streaming")
arrow_table = df.to_arrow()
pq.write_table(arrow_table, "data/compressed.parquet", compression="snappy")
Writing pyarrow_parquet.py
!uv run pyarrow_parquet.py
!stat -c %s data/sample*.parquet
!stat -c %s data/compressed.parquet
21495
21495
4149

Entscheidungsmatrix

SzenarioEmpfehlungGrund
< 2 GBPandasViel Code öffentlich
> 2 GB, analytischPolarsHohe Performance, Streaming-Support, Gut durchdachte API
> RAM, verteiltDask / RayHorizontale Skalierung auf Cluster-Ebene
!rm -rf data
!rm -f polars_reader.py polars_writer.py pyarrow_parquet.py sklearn_parquet.py torch_parquet.py