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:

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:

Datenmodell der Transformationsschicht

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:

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: dev

Dann 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: 2
sources:
- 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_location kann 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:

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_st
LEFT 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_rides
FROM {{ ref ("fact_services") }} AS srv
INNER JOIN {{ ref("dim_nl_train_stations") }} AS tr_st
ON srv.station_sk = tr_st.station_sk
INNER JOIN {{ ref("dim_nl_municipalities") }} AS m
ON tr_st.municipality_sk = m.municipality_sk
INNER JOIN {{ ref("dim_nl_provinces") }} AS p
ON m.province_sk = p.province_sk
WHERE service_year = {{ var('execution_year') }}
GROUP BY ALL

Die 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.parquet

PostgreSQL

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 srv
INNER JOIN {{ ref("rep_dim_nl_train_stations") }} AS tr_st
ON srv.station_sk = tr_st.station_sk
INNER JOIN {{ ref("rep_dim_nl_municipalities") }} AS mn
ON tr_st.municipality_sk = mn.municipality_sk
WHERE 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 ALL

Dank 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:

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: dev

Wir 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:

Terminal window
psql -U postgres
SELECT count(*), sum(number_of_rides)
FROM main_public.rep_fact_train_services_daily_agg;
count | sum
--------+----------
240826 | 17438151

Wichtig: Während PostgreSQL die Zieldatenbank ist und dbt merge als inkrementelle Strategie dafür bietet, geschieht die Ausführung der obigen Pipeline in DuckDB; der inkrementelle Load kann deshalb nur mit den Strategien append oder delete+insert erfolgen.

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.3
05:48:08 Registered adapter: duckdb=1.9.2
05:48:08 Found 10 models, 20 data tests, 4 sources, 565 macros
05:48:08
05: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:51
05:48:51 Completed successfully
05:48:51
05:48:51 Done. PASS=30 WARN=0 ERROR=0 SKIP=0 TOTAL=30

Fazit

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.