Erfahre, wie Du die Datenextraktion mit Snowflake User-Defined Table Functions in Deine dbt Data Pipelines integrierst. Der Beitrag zeigt am Beispiel von Salesforce, wie Du Daten direkt in Snowflake lädst und ohne separates Ingestion-Tool in dbt-Modellen materialisierst.
Das fehlende E in dbt: Datenextraktion mit dbt Core und Snowflake
Mathias Heinze
Mathias Heinze
Principal Consultant
Mathias ist Berater mit langjähriger Erfahrung in Data Vault Modellierung, DWH-Architektur und DWH-Automatisierung und hat Versicherungen zu diesen Themen beraten. Seine Projekterfahrung umfasst dabei den kompletten Wertschöpfungsprozess von der Analyse, Planung und Konzeption bis hin zur Umsetzung mit Fokus auf Snowflake und Oracle Technologien. Er verfügt zudem über fundiertes Wissen in Domain-Specific-Languages (DSL) zur Codegenerierung, Big Data und Cloud Technologien sowie agilen Projektvorgehen.
Erfahre, wie Du die Datenextraktion mit Snowflake User-Defined Table Functions in Deine dbt Data Pipelines integrierst. Der Beitrag zeigt am Beispiel von Salesforce, wie Du Daten direkt in Snowflake lädst und ohne separates Ingestion-Tool in dbt-Modellen materialisierst.
Inhaltsverzeichnis
dbt übernimmt Transformationen – die Extraktion erfolgt meist extern
dbt (data build tool) hat sich zu einem weit verbreiteten Standard für Datentransformationen im modernen Data Engineering entwickelt. Es ermöglicht Entwicklern, modulare, versionskontrollierte und getestete SQL-Pipelines zu erstellen und damit bewährte Prinzipien der Softwareentwicklung auf Datenprozesse anzuwenden.
dbt ist jedoch für das T in klassischen ETL- oder ELT-Pipelines konzipiert. Die Extraktion und das Laden von Daten deckt das Tool nicht eigenständig ab.
In herkömmlichen ELT-Architekturen übernehmen spezialisierte Tools wie Fivetran, Airbyte oder individuelle Ingestion-Skripte die Übernahme von Daten aus Datenbanken, SaaS-Anwendungen und APIs. Das kann den Betriebsaufwand erhöhen: Es kommen weitere Tools, Infrastruktur, Kontextwechsel und gegebenenfalls Kosten hinzu.
Was wäre, wenn sich die Extraktion direkt in dbt abbilden ließe? Statt eine separate Orchestrierungsschicht zu betreiben, eine VM zu verwalten oder zusätzliche Connector-Lizenzen zu nutzen, könnten Teams SQL-native Pipelines vom Quellsystem bis zum Data Mart erstellen.
Mit Snowflake lässt sich dieser Ansatz über ein schlankes Integrationsmuster umsetzen.
Snowflake User-Defined Table Functions für die Datenextraktion
Die entscheidende Funktionalität: Snowflake unterstützt Python-basierte User-Defined Table Functions (UDTFs) über Snowpark. Diese Funktionen führen Python-Code innerhalb von Snowflake aus und geben die Ergebnisse als gewöhnliche SQL-Tabelle zurück. Da dbt Core nativ in Snowflake verfügbar ist, muss dbt zudem nicht separat gehostet oder deployed werden.
Dadurch lassen sich folgende Schritte umsetzen:
Eine Python-Funktion erstellen, die sich mit einem externen Quellsystem verbindet
Die Funktion als Snowflake-Tabellenfunktion registrieren
Die Funktion in einem in Snowflake ausgeführten dbt-Modell mit Standard-SQL aufrufen
Aus Sicht von dbt handelt es sich um eine normale Tabellenabfrage. Im Hintergrund verbindet sich Snowflake mit dem Quellsystem, ruft die angeforderten Daten ab und gibt sie innerhalb einer SQL-Anweisung als Zeilen und Spalten zurück.
So entsteht eine durchgängige ELT-Data-Pipeline in dbt, ohne Technologiebruch.
Beispiel-Implementierung
Die folgenden Skripte zeigen eine Beispiel-Implementierung für Salesforce.
Erstellung einer Snowpark-Python-Tabellenfunktion
USE DATABASE DB_MAHT_DBT;
USE SCHEMA UTIL;
-- create secrets
CREATE OR REPLACE SECRET DB_MAHT_DBT.UTIL.SALESFORCE_CREDS
TYPE = GENERIC_STRING
SECRET_STRING = '{"username": "<username>",
"password": "<password>",
"security_token": "<security-token>",
"instance_url": "<instance-url>"}';
-- create network rule
CREATE OR REPLACE NETWORK RULE DB_MAHT_DBT.UTIL.SALESFORCE_NETWORK_RULE
MODE = EGRESS
TYPE = HOST_PORT
VALUE_LIST = (
'login.salesforce.com',
'<instance-url>' );
-- create external access integration
CREATE OR REPLACE EXTERNAL ACCESS INTEGRATION SALESFORCE_ACCESS_INTEGRATION
ALLOWED_NETWORK_RULES = (DB_MAHT_DBT.UTIL.SALESFORCE_NETWORK_RULE)
ALLOWED_AUTHENTICATION_SECRETS = (DB_MAHT_DBT.UTIL.SALESFORCE_CREDS)
ENABLED = TRUE;
-- create UDTF
CREATE OR REPLACE FUNCTION DB_MAHT_DBT.UTIL.FCT_LOAD_SALESFORCE_TABLE(TARGET_TABLE STRING, FIELD_LIST STRING DEFAULT '', WHERE_CLAUSE STRING DEFAULT '', SOQL_QUERY STRING DEFAULT '')
RETURNS TABLE ("DATA" OBJECT)
LANGUAGE PYTHON
RUNTIME_VERSION = '3.11'ARTIFACT_REPOSITORY = snowflake.snowpark.pypi_shared_repository
PACKAGES = ('simple-salesforce', 'snowflake-snowpark-python')
EXTERNAL_ACCESS_INTEGRATIONS = (SALESFORCE_ACCESS_INTEGRATION)
SECRETS = ('sf_creds' = DB_MAHT_DBT.UTIL.SALESFORCE_CREDS)
HANDLER = 'pythonDataReader'AS
$$
import _snowflake
import json
from simple_salesforce import Salesforce
classpythonDataReader:defprocess(self, target_table, field_list, where_clause, soql_query): creds = json.loads(_snowflake.get_generic_secret_string('sf_creds'))
sf = Salesforce(
username=creds['username'],
password=creds['password'],
security_token=creds['security_token'],
instance_url=creds['instance_url']
)
if soql_query and soql_query.strip():
query = soql_query.strip()
result = sf.query_all(query)
records = result['records']
if records:
fields = [k for k in records[0].keys() if k != 'attributes']
else:
fields = []
else:
sf_object = getattr(sf, target_table)
if field_list and field_list.strip():
fields = [f.strip() for f in field_list.split(',')]
else:
desc = sf_object.describe()
fields = [f['name'] for f in desc['fields']]
query = "SELECT " + ", ".join(fields) + " FROM " + target_table
if where_clause and where_clause.strip():
query += " WHERE " + where_clause
result = sf.query_all(query)
records = result['records']
for r in records:
row = {f: str(r.get(f)) if r.get(f) isnotNoneelseNonefor f in fields}
yield (row,)
$$;
Beispiel-Aufrufe der UDTF
-- select all columns:
SELECT DATA FROM TABLE(DB_MAHT_DBT.UTIL.FCT_LOAD_SALESFORCE_TABLE('Account'));
-- select specific columns:
SELECT DATA FROM TABLE(DB_MAHT_DBT.UTIL.FCT_LOAD_SALESFORCE_TABLE('Account', 'Id, Name, Industry, BillingCity'));
-- using a filter:
SELECT DATA FROM TABLE(DB_MAHT_DBT.UTIL.FCT_LOAD_SALESFORCE_TABLE('Account', 'Id, Name, Industry', 'Industry = \'Electronics\''));
-- using custom SOQL query:
SELECT DATA FROM TABLE(DB_MAHT_DBT.UTIL.FCT_LOAD_SALESFORCE_TABLE('a', 'b', 'c', 'SELECT Id, Name FROM Account WHERE Id = \'001g500000KiFO3AAN\''));
In diesem Beispiel gibt die UDTF eine VARIANT-Spalte zurück. Dabei handelt es sich um Snowflakes flexiblen Datentyp für semistrukturierte Daten und JSON-ähnliche Objekte. Nachgelagerte Modelle können einzelne Werte mit der Doppelpunktnotation von Snowflake auslesen, beispielsweise mit data:Id::STRING.
Beispiel-Ausgabe der UDTF
Verwendung der Snowpark-UDTF in einem dbt-Modell
Hier greifen die Komponenten ineinander. Ein dbt-Modell kann die UDTF wie jede andere SQL-Tabellenfunktion aufrufen, während dbt die Ausführung des Modells übernimmt. Die extrahierten Datensätze werden als Snowflake-Tabelle materialisiert und können direkt in die bestehende Transformationspipeline einfließen.
{{ config(
materialized='table') }}
SELECT
data::variant AS data,
'{{ run_started_at }}'::timestamp(0) AS prj$load_dts
FROM TABLE(DB_MAHT_DBT.UTIL.FCT_LOAD_SALESFORCE_TABLE('Account'))
Fazit: Eine vollständige ELT-Pipeline innerhalb von dbt
dbt wurde für Datentransformationen entwickelt und bietet dafür eine solide Grundlage. In Verbindung mit Snowpark User-Defined Table Functions von Snowflake kann es auch als Einstiegspunkt für die Extraktion dienen.
Mit diesem Muster können Teams eine ELT-Pipeline vom Quellsystem bis zum nutzungsbereiten Data Mart innerhalb eines dbt-Projekts und einer Snowflake-Umgebung aufbauen. Ein zusätzliches Extraktionstool oder separate Infrastruktur ist dafür nicht erforderlich.
Das breite Ökosystem an Python-Bibliotheken kann diesen Ansatz für viele Arten von Quellsystemen nutzbar und skalierbar machen, darunter:
Relationale Datenbanken wie PostgreSQL, Oracle und SQL-Server
SaaS-Plattformen wie Salesforce, HubSpot und Jira
Datei- und Kollaborationssysteme wie SharePoint
REST-APIs
Für Teams, die bereits dbt und Snowflake einsetzen, kann es sich lohnen, dieses Muster zu prüfen, bevor ein weiteres Extraktionstool in den Technologie-Stack aufgenommen wird.
Stehst Du aktuell auch vor der Herausforderung, Deine Datenintegration effizienter aufzustellen? Dann lass uns sprechen. Gemeinsam finden wir die passende Lösung für Dein Unternehmen.
Du hast Fragen? Kontaktiere uns
Your contact person
Helene Fuchs
Domain Lead Data Platform & Data Management
Your contact person
Pia Ehrnlechner
Domain Lead Data Platform & Data Management
Wer ist b.telligent?
b.telligent – das ist Data Analytics, AI, Customer Engagement und Data Visualisation. Das ist Deutschland, Österreich, die Schweiz und Rumänien. Doch das Entscheidende ist unser Team: Menschen mit echter Leidenschaft für Daten, die gemeinsam innovative Lösungen schaffen und Unternehmen nachhaltig voranbringen.
Wir zeigen auf, wie Du mit Snowflake Intelligence schneller an verwertbare Insights und Entscheidungen kommst. Und: Wie Du Use Cases realisierst und den ROI der Datenplattform steigerst: mit wenig Aufwand und ohne technisches Know-how.
Openflow von Snowflake macht die Datenintegration schneller, einfacher und effizienter. In diesem Artikel zeigen wir, wie sich diese Vorteile in der Praxis auswirken. Anhand eines echten Beispiels werden Strategien für die einfache und leistungsstarke Verarbeitung großer Mengen kleiner, eingehender Dateien vorgestellt.
Mit Openflow vereinfacht Snowflake die Datenintegration grundlegend: Extraktion und Laden erfolgen als Bestandteil der Snowflake Plattform – ganz ohne externe ETL-Tools. Damit sinkt der Integrationsaufwand deutlich, und das komplette Pipeline-Management wird erheblich schlanker und effizienter.