2025-04-04
Vollständig lokale Datentransformation mit dbt und DuckDB
Petrica Leuca
Einführung
Das Data Build Tool, dbt, ist ein Open-Source-Transformations-Framework, das Datenteams erlaubt, Software-Engineering-Best-Practices im gelieferten Code zu übernehmen, etwa Git-Workflow und Unit Testing. Weitere bemerkenswerte Features von dbt sind Data Lineage, Dokumentation und Datentests als Teil der Ausführungspipeline.
In diesem Artikel zeigen wir, wie man vollständig lokale Datentransformationen mit dbt und DuckDB durchführt.
Dazu nutzen wir den Adapter dbt-duckdb, der den dbt-Standard mit DuckDBs Verarbeitungspower integriert.
Datenmodell
Wir nutzen in diesem Beitrag zwei offene Datensätze: Eisenbahnservices, bereitgestellt vom Team hinter der Anwendung Rijden de Treinen (Fahren die Züge?), und Kartografie-Informationen über die Niederlande, bereitgestellt von cartomap. Die Datensätze sind organisiert in:
- eine persistierte DuckDB-Datenbank, in der wir die Eisenbahnservicedaten von 2024 speichern, gehostet auf Cloudflare;
- eine Provinzen-GeoJSON-Datei mit den geografischen Informationen über die niederländischen Provinzen, gehostet auf GitHub;
- eine Gemeinden-GeoJSON-Datei mit den geografischen Informationen über niederländische Gemeinden, zusammen mit dem Code gespeichert.
Nach einer ersten Exploration der obigen Daten können wir beobachten, dass eine Provinz eine oder viele Gemeinden haben kann, eine Gemeinde keine oder viele Bahnhöfe haben kann und ein Zugservice-Datensatz mit genau einem Bahnhof verbunden ist. Deshalb entscheiden wir uns für das folgende Datenmodell:
- eine Dimensionstabelle
dim_nl_provincesfür Informationen über die niederländischen Provinzen; - eine Dimensionstabelle
dim_nl_municipalitiesfür Informationen über die niederländischen Gemeinden, verknüpft mitdim_nl_provinces; - eine Dimensionstabelle
dim_nl_train_stationsfür Informationen über die Bahnhöfe in den Niederlanden, verknüpft mitdim_nl_municipalities; - eine Faktentabelle
fact_servicesfür Informationen über die Zugservices, verknüpft mitdim_nl_train_stations.
Der Zweck der Verarbeitung der Zugservices und der niederländischen Kartografiedaten ist, die Daten in einer Struktur zu organisieren, die in künftigen Eisenbahn-Datenanalyse-Use-Cases leicht genutzt werden kann. Solche kuratierten Datenstrukturen werden oft Data Marts genannt.
Tipp: Der folgende Code ist auf GitHub verfügbar.
Daten mit DuckDB und dbt verarbeiten
Nachdem wir unser Projekt initialisiert haben, konfigurieren wir die Verbindungsdetails für DuckDB in der Datei profiles.yml. Neben der Angabe, ob die Datenbank im Speicher oder auf Disk persistiert sein soll, spezifizieren wir auch:
- welche Erweiterungen für die Datenverarbeitung nötig sind, z. B. spatial;
- externe Datenbanken, angehängt von der lokalen Disk oder anderen Storage-Lösungen.
dutch_railway_network:
outputs: dev: type: duckdb path: data/dutch_railway_network.duckdb extensions: - spatial - httpfs threads: 5 attach: - path: 'https://blobs.duckdb.org/nl-railway/train_stations_and_services.duckdb' type: duckdb alias: external_db target: devDann konfigurieren wir die Datei sources.yml unter dem Verzeichnis models, indem wir
externe Quellen (etwa Dateien) und Tabellendefinitionen
aus den angehängten Datenbank(en) angeben:
version: 2sources: - name: geojson_external tables: - name: nl_provinces config: external_location: "https://cartomap.github.io/nl/wgs84/provincie_2025.geojson" - name: nl_municipalities config: external_location: "seeds/gemeente_2025.geojson" - name: external_db database: external_db schema: main tables: - name: stations - name: services
external_locationkann auf CSV-, Parquet- oder JSON-Dateien zeigen. Sowohl das lokale Dateisystem als auch Remote-Endpunkte (z. B. HTTP oder S3) werden unterstützt.
Mit definiertem Profil und definierter Quelle können wir jetzt die Daten laden.
Daten nach DuckDB laden
In dbt heißt die Art, Daten im Zielsystem zu speichern, materialization. Der Adapter dbt-duckdb bietet die folgenden Materialisierungsoptionen:
table, ersetzt die Zieltabelle bei jedem Lauf;incremental, mit den Optionenappendunddelete+insert, ändert die Daten in der Tabelle, aber nicht die Tabelle selbst (falls sie existiert);snapshot, implementiert eine Tabelle Slowly Changing Dimension Type 2 mit Zeitgültigkeitsintervallen;view, ersetzt die Zielsicht bei jedem Lauf.
Ein weiteres Feature dieses Adapters: Der in den Datenverarbeitungsskripten genutzte SQL-Dialekt hat alle freundlichen SQL-Erweiterungen, die DuckDB bietet.
Um die Daten in der Tabelle dim_nl_provinces zu refreshen, nutzen wir die Spatial-Funktion st_read, die die in sources.yml definierte GeoJSON-Datei nl_provinces automatisch liest und parst.
{{ config(materialized='table') }}
SELECT {{ dbt_utils.generate_surrogate_key(['id']) }} AS province_sk, id AS province_id, statnaam AS province_name, geom AS province_geometry, {{ common_columns() }}FROM st_read({{ source("geojson_external", "nl_provinces") }}) AS src;Ähnlich refreshen wir die Daten für dim_nl_municipalities und fact_services vollständig.
Um die Beziehung zwischen einem Bahnhofsstandort und einer Gemeinde aufzubauen,
nutzen wir die Spatial-Funktion st_contains, die true liefert, wenn eine Geometrie eine andere Geometrie enthält:
{{ config(materialized='table') }}
SELECT {{ dbt_utils.generate_surrogate_key(['tr_st.code']) }} AS station_sk, tr_st.id AS station_id, tr_st.code AS station_code, tr_st.name_long AS station_name, tr_st.type AS station_type, st_point(tr_st.geo_lng, tr_st.geo_lat) AS station_geo_location, coalesce(dim_mun.municipality_sk, 'unknown') AS municipality_sk, {{ common_columns() }}FROM {{ source("external_db", "stations") }} AS tr_stLEFT JOIN {{ ref ("dim_nl_municipalities") }} AS dim_mun ON st_contains( dim_mun.municipality_geometry, st_point(tr_st.geo_lng, tr_st.geo_lat) )WHERE tr_st.country = 'NL';Um aus den externen Quellen zu lesen, referenzieren wir die Quelle, indem wir Quellen- und Tabellennamen angeben.
Daten aus DuckDB exportieren
Ein großer Vorteil von DuckDB für die Datenverarbeitung ist die Fähigkeit, Daten in Dateien zu exportieren (etwa CSV, JSON und Parquet) und Daten direkt in PostgreSQL- oder MySQL-Datenbanken zu refreshen.
Externe Dateien
Das Feature, Daten in Dateien zu exportieren, wird vom Adapter dbt-duckdb mit der Materialisierung external ermöglicht. Mit der Materialisierung external können wir Daten in die Dateitypen CSV, JSON und Parquet an einen angegebenen Storage-Ort (lokal oder extern) exportieren. Der Load-Typ ist full refresh, bestehende Dateien werden also überschrieben.
Im folgenden Verarbeitungsschritt exportieren wir aggregierte Zugservice-Daten auf Monatsebene in eine Parquet-Datei, partitioniert nach Jahr und Monat:
{{ config( materialized='external', location="data/exports/nl_train_services_aggregate", options={ "partition_by": "service_year, service_month", "overwrite": True } )}}
SELECT year(service_date) AS service_year, month(service_date) AS service_month, service_type, service_company, tr_st.station_sk, tr_st.station_name, m.municipality_sk, m.municipality_name, p.province_sk, p.province_name, count(*) AS number_of_ridesFROM {{ ref ("fact_services") }} AS srvINNER JOIN {{ ref("dim_nl_train_stations") }} AS tr_st ON srv.station_sk = tr_st.station_skINNER JOIN {{ ref("dim_nl_municipalities") }} AS m ON tr_st.municipality_sk = m.municipality_skINNER JOIN {{ ref("dim_nl_provinces") }} AS p ON m.province_sk = p.province_skWHERE service_year = {{ var('execution_year') }}GROUP BY ALLDie exportierten Dateien liegen in einer Hive-partitionierten Verzeichnisstruktur.
./service_year=2024/service_month=1:49255 Apr 2 14:54 data_0.parquet
...
./service_year=2024/service_month=12:48031 Apr 2 14:54 data_0.parquetPostgreSQL
Nachdem wir die Daten in unserem Eisenbahnservices-Data-Mart verarbeitet haben, können wir daraus eine Tagesaggregation auf Bahnhofsebene erzeugen, organisiert in einem Star-Schema-Modell, indem die Dimensionsschlüssel Teil der Daten sind:
{{ config( materialized='incremental', incremental_strategy='delete+insert', unique_key=""" service_date, service_type, service_company, station_sk """ )}}
SELECT service_date, service_type, service_company, srv.station_sk, mn.municipality_sk, province_sk, count(*) AS number_of_rides, {{ common_columns() }}FROM {{ ref ("fact_services") }} AS srvINNER JOIN {{ ref("rep_dim_nl_train_stations") }} AS tr_st ON srv.station_sk = tr_st.station_skINNER JOIN {{ ref("rep_dim_nl_municipalities") }} AS mn ON tr_st.municipality_sk = mn.municipality_skWHERE NOT service_arrival_cancelled
{% if is_incremental() %} AND srv.invocation_id = ( SELECT invocation_id FROM {{ ref("fact_services") }} ORDER BY last_updated_dt DESC LIMIT 1 ) {% endif %}GROUP BY ALLDank DuckDBs Fähigkeit, sich mit einer PostgreSQL-Datenbank zu verbinden und dahin zu schreiben, können wir den obigen Verarbeitungsschritt zu unserem dbt-Projekt unter dem Verzeichnis models/reverse_etl hinzufügen.
Um uns mit einer PostgreSQL-Datenbank zu verbinden, müssen wir in profiles.yml angeben:
- die
postgres-Erweiterung; - den PostgreSQL-Connection-String im Abschnitt
attach.
dutch_railway_network:
outputs: dev: type: duckdb path: data/dutch_railway_network.duckdb extensions: - ... - postgres threads: 5 attach: - ... - path: "postgresql://postgres:{{ env_var('DBT_DUCKDB_PG_PWD') }}@localhost:5466/postgres" type: postgres alias: postgres_db target: devWir müssen außerdem die Datenbankdetails des Modells in dbt_project.yml konfigurieren:
models: dutch_railway_network: transformation: schema: main +docs: node_color: 'silver' reverse_etl: database: postgres_db schema: public +docs: node_color: '#d5b85a'Mit dieser Konfiguration werden alle Modelle aus dem Verzeichnis transformation auf dem Schema main_main ausgeführt, während die Modelle aus reverse_etl auf dem Schema main_public ausgeführt werden.
Nach dem Ausführen der Modelle mit dbt run --model +reverse_etl sind die Daten aus PostgreSQL abfragbar:
psql -U postgresSELECT count(*), sum(number_of_rides)FROM main_public.rep_fact_train_services_daily_agg; count | sum--------+---------- 240826 | 17438151Wichtig: Während PostgreSQL die Zieldatenbank ist und
dbtmergeals inkrementelle Strategie dafür bietet, geschieht die Ausführung der obigen Pipeline in DuckDB; der inkrementelle Load kann deshalb nur mit den Strategienappendoderdelete+inserterfolgen.
Ausführungsdetails
Die obige Implementierung besteht aus 10 Modellen und 20 Datentests und verarbeitet 400 MB Daten aus der angehängten DuckDB-Datenbank zusammen mit kleinen Daten in GeoJSON-Dateien. Die Gesamtlaufzeit auf einem einzelnen Thread und einem MacBook Pro mit 12 GB liegt zwischen 40 und 45 Sekunden. Von der Gesamtlaufzeit werden etwa 30 Sekunden für die Verarbeitung der Zugservice-Daten und 4 Sekunden für das Schreiben der aggregierten Daten nach PostgreSQL aufgewendet:
05:48:07 Running with dbt=1.9.305:48:08 Registered adapter: duckdb=1.9.205:48:08 Found 10 models, 20 data tests, 4 sources, 565 macros05:48:0805:48:08 Concurrency: 1 threads (target='dev')...05:48:45 19 of 30 OK created sql table model main_main.fact_services .................... [OK in 32.60s]...05:48:50 26 of 30 OK created sql incremental model postgres_db.main_public.rep_fact_train_services_daily_agg [OK in 3.74s]...05:48:51 Finished running 2 external models, 1 incremental model, 7 table models, 20 data tests in 0 hours 0 minutes and 42.63 seconds (42.63s).05:48:5105:48:51 Completed successfully05:48:5105:48:51 Done. PASS=30 WARN=0 ERROR=0 SKIP=0 TOTAL=30Fazit
In diesem Beitrag haben wir gezeigt, wie DuckDB mit dbt integriert und Teil des Datenverarbeitungs-Ökosystems ist – anhand von Data-Mart-Erzeugung, Dateiexporten und Reverse ETL in eine PostgreSQL-Datenbank.