2021-06-25

Parquet präzise mit DuckDB abfragen

Hannes Mühleisen, Mark Raasveldt

Apache Parquet ist das gängigste „Big Data“-Speicherformat für Analysen. In Parquet-Dateien liegen Daten in einem spaltenkomprimierten Binärformat. Jede Parquet-Datei speichert eine einzelne Tabelle. Die Tabelle ist in Row Groups unterteilt, die jeweils eine Teilmenge der Zeilen enthalten. Innerhalb einer Row Group sind die Tabellendaten spaltenweise gespeichert.

Example Parquet file shown visually. The Parquet file (taxi.parquet) is divided into row groups that each have two columns (pickup_at and dropoff_at)

Das Parquet-Format hat mehrere Eigenschaften, die es für analytische Einsatzfälle geeignet machen:

  1. Die spaltenweise Darstellung bedeutet, dass einzelne Spalten (effizient) gelesen werden können. Die gesamte Datei muss nicht immer gelesen werden!
  2. Die Datei enthält pro Spalte Statistiken in jeder Row Group (Min-/Max-Wert und die Anzahl der NULL-Werte). Diese Statistiken erlauben es dem Reader, Row Groups zu überspringen, wenn sie nicht benötigt werden.
  3. Die spaltenweise Kompression reduziert die Dateigröße deutlich und damit den Speicherbedarf von Datensätzen. Aus Big Data wird oft Medium Data.

DuckDB und Parquet

DuckDBs abhängigkeitsfreier Parquet-Reader kann SQL-Abfragen direkt auf Parquet-Dateien ausführen, ohne Import- oder Analyseschritt. Dank des natürlichen Spaltenformats von Parquet ist das sehr schnell!

DuckDB liest die Parquet-Dateien streaming, Sie können also Abfragen auf großen Parquet-Dateien ausführen, die nicht in den Hauptspeicher passen.

DuckDB erkennt automatisch, welche Spalten und Zeilen für eine gegebene Abfrage nötig sind. So lassen sich deutlich größere und komplexere Parquet-Dateien analysieren, ohne manuelle Optimierungen oder zusätzliche Hardware.

Als Bonus kann DuckDB das alles parallel und über mehrere Parquet-Dateien gleichzeitig mit der Glob-Syntax tun.

Als kurzer Vorgeschmack: Dieses Snippet führt eine SQL-Abfrage direkt auf einer Parquet-Datei aus.

Das DuckDB-Paket installieren:

Terminal window
pip install duckdb

Die Parquet-Datei herunterladen:

Terminal window
wget https://blobs.duckdb.org/data/taxi_2019_04.parquet

Dann folgendes Python-Skript ausführen:

import duckdb
print(duckdb.query('''
SELECT count(*)
FROM 'taxi_2019_04.parquet'
WHERE pickup_at BETWEEN '2019-04-15' AND '2019-04-20'
''').fetchall())

Automatischer Filter- und Projection-Pushdown

Schauen wir uns die vorherige Abfrage genauer an, um die Stärke des Parquet-Formats in Kombination mit DuckDBs Query-Optimizer zu verstehen.

SELECT count(*)
FROM 'taxi_2019_04.parquet'
WHERE pickup_at BETWEEN '2019-04-15' AND '2019-04-20';

In dieser Abfrage lesen wir eine einzelne Spalte aus unserer Parquet-Datei (pickup_at). Alle anderen in der Datei gespeicherten Spalten können vollständig übersprungen werden, weil wir sie für die Antwort nicht brauchen.

Projection & filter pushdown into Parquet file example.

Zusätzlich beeinflussen nur Zeilen mit einem pickup_at zwischen dem 15. und 20. April 2019 das Ergebnis. Alle Zeilen, die dieses Prädikat nicht erfüllen, können übersprungen werden.

Die Statistiken in der Parquet-Datei sind hier sehr nützlich. Jede Row Group mit einem Max-Wert von pickup_at unter 2019-04-15 oder einem Min-Wert über 2019-04-20 kann übersprungen werden. In manchen Fällen können wir so ganze Dateien überspringen.

DuckDB versus Pandas

Um zu zeigen, wie wirksam diese automatischen Optimierungen sind, führen wir einige Abfragen auf Parquet-Dateien mit Pandas und DuckDB aus.

Wir nutzen einen Teil des bekannten New-York-Taxi-Datensatzes als Parquet-Dateien, konkret Daten von April, Mai und Juni 2019. Die Dateien sind zusammen etwa 360 MB groß und enthalten rund 21 Millionen Zeilen mit je 18 Spalten. Die drei Dateien liegen im Ordner taxi/.

Die Beispiele sind als interaktives Notebook auf Google Colab verfügbar. Die hier genannten Zeiten stammen aus dieser Umgebung, damit sie nachvollziehbar sind.

Mehrere Parquet-Dateien lesen

Zuerst schauen wir uns einige Zeilen im Datensatz an. Im Ordner taxi/ liegen drei Parquet-Dateien. DuckDB unterstützt die Glob-Syntax, mit der sich alle drei Dateien gleichzeitig abfragen lassen.

con.execute("""
SELECT *
FROM 'taxi/*.parquet'
LIMIT 5""").df()
pickup_at dropoff_at passenger_count trip_distance rate_code_id
2019-04-01 00:04:09 2019-04-01 00:06:35 1 0.5 1
2019-04-01 00:22:45 2019-04-01 00:25:43 1 0.7 1
2019-04-01 00:39:48 2019-04-01 01:19:39 1 10.9 1
2019-04-01 00:35:32 2019-04-01 00:37:11 1 0.2 1
2019-04-01 00:44:05 2019-04-01 00:57:58 1 4.8 1

Obwohl die Abfrage alle Spalten aus drei (recht großen) Parquet-Dateien wählt, ist sie sofort fertig. DuckDB verarbeitet die Parquet-Datei streaming und hört nach den ersten Zeilen auf zu lesen, weil das für die Abfrage reicht.

Versuchen wir dasselbe in Pandas, merken wir, dass es nicht so einfach ist: Pandas kann nicht mehrere Parquet-Dateien in einem Aufruf lesen. Zuerst müssen wir pandas.concat nutzen, um die drei Dateien zusammenzufügen:

import pandas
import glob
df = pandas.concat(
[pandas.read_parquet(file)
for file
in glob.glob('taxi/*.parquet')])
print(df.head(5))

Unten die Zeiten für beide Abfragen.

System Time (s)
DuckDB 0.015
Pandas 12.300

Pandas braucht deutlich länger. Es muss nicht nur jede der drei Parquet-Dateien vollständig lesen, sondern die drei getrennten Pandas-DataFrames auch noch zusammenfügen.

Zu einer einzelnen Datei zusammenfügen

Das Zusammenfügen können wir umgehen, indem wir aus den drei kleineren Teilen eine große Parquet-Datei erzeugen. Dafür nutzen wir die Bibliothek pyarrow, die mehrere Parquet-Dateien lesen und in eine große Datei streamen kann. Der pyarrow-Parquet-Reader ist derselbe, den Pandas intern verwendet.

import pyarrow.parquet as pq
# concatenate all three Parquet files
pq.write_table(pq.ParquetDataset('taxi/').read(), 'alltaxi.parquet', row_group_size=100000)

Hinweis: DuckDB kann Parquet-Dateien auch schreiben, über das COPY-Statement.

Die große Datei abfragen

Wiederholen wir das vorherige Experiment, diesmal mit der einzelnen Datei.

# DuckDB
con.execute("""
SELECT *
FROM 'alltaxi.parquet'
LIMIT 5""").df()
# Pandas
pandas.read_parquet('alltaxi.parquet')
.head(5)
System Time (s)
DuckDB 0.02
Pandas 7.50

Pandas ist besser als zuvor, weil das Zusammenfügen entfällt. Die gesamte Datei muss aber weiterhin in den Speicher gelesen werden, was Zeit und Speicher kostet.

Für DuckDB spielt es kaum eine Rolle, wie viele Parquet-Dateien eine Abfrage liest.

Zeilen zählen

Angenommen, wir wollen wissen, wie viele Zeilen unser Datensatz hat. Das geht so:

# DuckDB
con.execute("""
SELECT count(*)
FROM 'alltaxi.parquet'
""").df()
# Pandas
len(pandas.read_parquet('alltaxi.parquet'))
System Time (s)
DuckDB 0.015
Pandas 7.500

DuckDB ist sehr schnell fertig, weil es automatisch erkennt, was aus der Parquet-Datei gelesen werden muss, und die nötigen Reads minimiert. Pandas muss die gesamte Datei erneut lesen und braucht deshalb dieselbe Zeit wie zuvor.

Für diese Abfrage können wir Pandas durch manuelle Optimierung verbessern. Für einen Count brauchen wir nur eine Spalte aus der Datei. Geben wir im read_parquet-Befehl manuell eine einzelne Spalte an, erhalten wir dasselbe Ergebnis deutlich schneller.

len(pandas.read_parquet('alltaxi.parquet', columns=['vendor_id']))
System Time (s)
DuckDB 0.015
Pandas 7.500
Pandas (optimized) 1.200

Das ist viel schneller, braucht aber immer noch mehr als eine Sekunde, weil die gesamte Spalte vendor_id als Pandas-Spalte in den Speicher muss, nur um die Zeilen zu zählen.

Zeilen filtern

Häufig filtert man auf die interessanten Teile eines Datensatzes. Angenommen, wir wollen wissen, wie viele Taxifahrten nach dem 30. Juni 2019 stattfinden. In DuckDB geht das so:

con.execute("""
SELECT count(*)
FROM 'alltaxi.parquet'
WHERE pickup_at > '2019-06-30'
""").df()

Die Abfrage ist in 45ms fertig und liefert:

count
167022

In Pandas können wir dieselbe Operation naiv ausführen.

# pandas naive
len(pandas.read_parquet('alltaxi.parquet')
.query("pickup_at > '2019-06-30'"))

Das liest wieder die gesamte Datei in den Speicher, die Abfrage braucht 7,5 s. Mit manuellem Projection Pushdown kommen wir auf 0,9 s. Immer noch deutlich langsamer als DuckDB.

# pandas projection pushdown
len(pandas.read_parquet('alltaxi.parquet', columns=['pickup_at'])
.query("pickup_at > '2019-06-30'"))

Der pyarrow-Parquet-Reader erlaubt jedoch auch Filter-Pushdown in den Scan. Damit landen wir bei deutlich konkurrenzfähigeren 70ms.

len(pandas.read_parquet('alltaxi.parquet', columns=['pickup_at'], filters=[('pickup_at', '>', '2019-06-30')]))
System Time (s)
DuckDB 0.05
Pandas 7.50
Pandas (projection pushdown) 0.90
Pandas (projection & filter pushdown) 0.07

Das zeigt: Die Ergebnisse liegen nicht daran, dass DuckDBs Parquet-Reader schneller wäre als der von pyarrow. DuckDB ist bei diesen Abfragen besser, weil seine Optimizer automatisch alle benötigten Spalten und Filter aus der SQL-Abfrage ziehen und ohne manuellen Aufwand im Parquet-Reader nutzen.

Interessanterweise sind sowohl der pyarrow-Parquet-Reader als auch DuckDB deutlich schneller als dieselbe Operation nativ in Pandas auf einem materialisierten DataFrame.

# read the entire Parquet file into Pandas
df = pandas.read_parquet('alltaxi.parquet')
# run the query natively in Pandas
# note: we only time this part
print(len(df[['pickup_at']].query("pickup_at > '2019-06-30'")))
System Time (s)
DuckDB 0.05
Pandas 7.50
Pandas (projection pushdown) 0.90
Pandas (projection & filter pushdown) 0.07
Pandas (native) 0.26

Aggregationen

Schließlich eine komplexere Aggregation. Angenommen, wir wollen die Anzahl der Fahrten pro Fahrgastzahl berechnen. Mit DuckDB und SQL sieht das so aus:

con.execute("""
SELECT passenger_count, count(*)
FROM 'alltaxi.parquet'
GROUP BY passenger_count""").df()

Die Abfrage ist in 220ms fertig und liefert:

passenger_count count
0 408742
1 15356631
2 3332927
3 944833
4 439066
5 910516
6 546467
7 106
8 72
9 64

Für SQL-Skeptiker und als Vorgeschmack auf einen späteren Beitrag: DuckDB hat auch eine „Relational API“, mit der sich Abfragen pythonischer deklarieren lassen. Das Äquivalent zur SQL-Abfrage oben liefert dasselbe Ergebnis und dieselbe Leistung:

con.from_parquet('alltaxi.parquet')
.aggregate('passenger_count, count(*)')
.df()

Zum Vergleich dieselbe Abfrage in Pandas wie zuvor.

# naive
pandas.read_parquet('alltaxi.parquet')
.groupby('passenger_count')
.agg({'passenger_count' : 'count'})
# projection pushdown
pandas.read_parquet('alltaxi.parquet', columns=['passenger_count'])
.groupby('passenger_count')
.agg({'passenger_count' : 'count'})
# native (parquet file pre-loaded into memory)
df.groupby('passenger_count')
.agg({'passenger_count' : 'count'})
System Time (s)
DuckDB 0.22
Pandas 7.50
Pandas (projection pushdown) 0.58
Pandas (native) 0.51

DuckDB ist in allen drei Szenarien schneller als Pandas, ohne manuelle Optimierungen und ohne die Parquet-Datei vollständig in den Speicher zu laden.

Fazit

DuckDB kann Abfragen effizient direkt auf Parquet-Dateien ausführen, ohne eine initiale Ladephase. Das System nutzt automatisch alle fortgeschrittenen Parquet-Eigenschaften, um die Ausführung zu beschleunigen.

DuckDB ist ein freies und quelloffenes Datenbankmanagementsystem (MIT-Lizenz). Es will das SQLite für Analysen sein und bietet ein schnelles, effizientes Datenbanksystem ohne externe Abhängigkeiten. Es ist nicht nur für Python verfügbar, sondern auch für C/C++, R, Java und mehr.