
Na jakich zasadach opiera się idealne magazynowanie danych?
Skupienie na wartości biznesowej i analizie przy braku powielania kodu. Zarządzanie DWH jako bazą kodu: wersjonowanie, przegląd, automatyczne testowanie i CI. Modularność, rozszerzalność, otwarty kod źródłowy i społeczność. Przyjazna dokumentacja użytkownika oraz wizualizacja zależności (Data Lineage).
O tym wszystkim w szczegółach oraz o roli DBT w ekosystemie Big Data i Analiz — zapraszam pod kat.
Cześć wszystkim
Z tej strony Artemij Kozyr. Już od ponad 5 lat pracuję z magazynami danych, zajmując się budowaniem ETL/ELT oraz analizą danych i wizualizacją. Obecnie pracuję w , uczę w OTUS na kursie , i dzisiaj chcę się z Wami podzielić artykułem, który napisałem w przeddzień startu nowej edycji kursu.
Krótki przegląd
Framework DBT — to wszystko o literze T w akronimie ELT (Extract — Transform — Load).
Wraz z pojawieniem się takich wydajnych i skalowalnych baz danych analitycznych jak BigQuery, Redshift, Snowflake, utraciło sens wykonywanie transformacji poza magazynem danych.
DBT nie wyciąga danych ze źródeł, ale oferuje ogromne możliwości pracy z danymi, które już zostały załadowane do magazynu (w Internal lub External Storage).

Główne zadanie DBT — wziąć kod, skompilować go do SQL, uruchomić polecenia w odpowiedniej kolejności w magazynie.
Struktura projektu DBT
Projekt składa się z dwóch typów katalogów i plików:
- Model (.sql) — jednostka transformacji, wyrażona zapytaniem SELECT
- Plik konfiguracyjny (.yml) — parametry, ustawienia, testy, dokumentacja
Na podstawowym poziomie praca wygląda następująco:
- Użytkownik przygotowuje kod modeli w dowolnym wygodnym IDE
- Za pomocą CLI wywoływane jest uruchomienie modeli, DBT kompiluje kod modeli do SQL
- Skompilowany kod SQL jest wykonywany w magazynie w określonej kolejności (graf)
Tak może wyglądać uruchomienie z CLI:

Wszystko jest SELECT
To jest killer feature frameworka Data Build Tool. Innymi słowy, DBT abstrahuje cały kod związany z materializacją Twoich zapytań w magazynie (warianty poleceń CREATE, INSERT, UPDATE, DELETE ALTER, GRANT, ...).
Każdy model wiąże się z napisaniem jednego zapytania SELECT, które definiuje wynikowy zbiór danych.
Logika transformacji może być wielowarstwowa i konsolidować dane z kilku innych modeli. Przykład modelu, który zbuduje witrynę zamówień (f_orders):
{% set payment_methods = ['credit_card', 'coupon', 'bank_transfer', 'gift_card'] %}
with orders as (
select * from {{ ref('stg_orders') }}
),
order_payments as (
select * from {{ ref('order_payments') }}
),
final as (
select
orders.order_id,
orders.customer_id,
orders.order_date,
orders.status,
{% for payment_method in payment_methods -%}
order_payments.{{payment_method}}_amount,
{% endfor -%}
order_payments.total_amount as amount
from orders
left join order_payments using (order_id)
)
select * from final
Co interesującego możemy tutaj zobaczyć?
Po pierwsze: Użyto CTE (Common Table Expressions) — do organizacji i zrozumienia kodu, który zawiera wiele transformacji i logiki biznesowej
Po drugie: Kod modelu — to mieszanka SQL i języka (templating language).
W przykładzie użyto pętli for do obliczenia sumy dla każdej metody płatności podanej w wyrażeniu set. Używana jest również funkcja ref — możliwość odwoływania się wewnątrz kodu do innych modeli:
- Podczas kompilacji ref zostanie przekształcony w docelowy wskaźnik na tabelę lub widok w magazynie
- ref pozwala na zbudowanie grafu zależności modeli
To właśnie dodaje do DBT niemal nieograniczone możliwości. Najczęściej używane z nich to:
- If / else statements — operatorzy warunkowi
- For loops — pętle
- Variables — zmienne
- Macro — tworzenie makr
Materializacja: Table, View, Incremental
Strategia Materializacji — podejście, według którego wynikowy zestaw danych modelu będzie przechowywany w magazynie.
W podstawowym ujęciu to:
- Table — fizyczna tabela w magazynie
- View — widok, wirtualna tabela w magazynie
Istnieją także bardziej złożone strategie materializacji:
- Incremental — inkrementalny załadunek (dużych tabel faktów); nowe wiersze są dodawane, zmienione — aktualizowane, usunięte — usuwane
- Ephemeral — model nie jest materializowany bezpośrednio, ale uczestniczy jako CTE w innych modelach
- Jakiekolwiek inne strategie, które możesz dodać samodzielnie
Oprócz strategii materializacji otwierają się możliwości optymalizacji pod konkretne magazyny, na przykład:
- Snowflake: Tabele przejrzyste, Zachowanie łączenia, Klasteryzacja tabel, Kopiowanie przywilejów, Bezpieczne widoki
- Redshift: Klucz rozdzielający, Klucz sortujący (przeplatany, złożony), Widoki z późnym powiązaniem
- BigQuery: Partycjonowanie i klasteryzacja tabel, Zachowanie łączenia, Szyfrowanie KMS, Etykiety i tagi
- Spark: Format pliku (parquet, csv, json, orc, delta), partition_by, clustered_by, buckets, incremental_strategy
Obecnie wspierane są następujące magazyny:
- Postgres
- Redshift
- BigQuery
- Snowflake
- Presto (częściowo)
- Spark (częściowo)
- Microsoft SQL Server (adapter społecznościowy)
Udoskonalmy nasz model:
- Możemy sprawić, by jego zapełnienie było inkrementalne (Incremental)
- Dodajmy klucze segmentacji i sortowania dla Redshift
-- Konfiguracja modelu:
-- Inkrementalne zapełnienie, unikalny klucz do aktualizacji rekordów (unique_key)
-- Klucz segmentacji (dist), klucz sortowania (sort)
{{
config(
materialized='incremental',
unique_key='order_id',
dist="customer_id",
sort="order_date"
)
}}
{% set payment_methods = ['credit_card', 'coupon', 'bank_transfer', 'gift_card'] %}
with orders as (
select * from {{ ref('stg_orders') }}
where 1=1
{% if is_incremental() -%}
-- Ten filtr zostanie zastosowany tylko w przypadku uruchomienia inkrementalnego
and order_date >= (select max(order_date) from {{ this }})
{%- endif %}
),
order_payments as (
select * from {{ ref('order_payments') }}
),
final as (
select
orders.order_id,
orders.customer_id,
orders.order_date,
orders.status,
{% for payment_method in payment_methods -%}
order_payments.{{payment_method}}_amount,
{% endfor -%}
order_payments.total_amount as amount
from orders
left join order_payments using (order_id)
)
select * from final
Graf zależności modeli
Jest to także drzewo zależności. Jest to DAG (Directed Acyclic Graph — Skierowany Graf Acykliczny).
DBT buduje graf na podstawie konfiguracji wszystkich modeli projektu, a dokładniej, odniesień ref() wewnątrz modeli do innych modeli. Posiadanie grafu umożliwia wykonanie następujących rzeczy:
- Uruchamianie modeli w odpowiedniej kolejności
- Równolegle tworzenie hurtowni
- Uruchamianie dowolnego podgrafu
Przykład wizualizacji grafu:

Każdy węzeł grafu to model, krawędzie grafu określone są przez wyrażenie ref.
Jakość danych i Dokumentacja
Oprócz tworzenia samych modeli, DBT pozwala przetestować szereg założeń (assertions) dotyczących wynikowego zbioru danych, takich jak:
- Not Null
- Unikalny
- Referencyjna integralność — integralność odniesienia (np. customer_id w tabeli orders odpowiada id w tabeli customers)
- Zgodność z listą dopuszczalnych wartości
Możliwe jest dodanie własnych testów (custom data tests), takich jak na przykład % odchylenia przychodów w porównaniu do danych sprzed dnia, tygodnia, miesiąca. Każde założenie, które można sformułować w postaci zapytania SQL, może stać się testem.
W ten sposób można wychwytywać w hurtowniach magazynów niepożądane odchylenia i błędy w danych.
Jeśli chodzi o dokumentację, DBT zapewnia mechanizmy umożliwiające dodawanie, wersjonowanie i publikowanie metadanych oraz komentarzy na poziomie modeli, a nawet atrybutów.
Oto jak wygląda dodawanie testów i dokumentacji na poziomie pliku konfiguracyjnego:
- name: fct_orders
description: Ta tabela zawiera podstawowe informacje o zamówieniach oraz kilka wyprowadzonych faktów na podstawie płatności
columns:
- name: order_id
tests:
- unique # test na unikalność wartości
- not_null # test na obecność null
description: To jest unikalny identyfikator dla zamówienia
- name: customer_id
description: Klucz obcy do tabeli klientów
tests:
- not_null
- relationships: # test na integralność referencyjną
to: ref('dim_customers')
field: customer_id
- name: order_date
description: Data (UTC), w której złożono zamówienie
- name: status
description: '{{ doc("orders_status") }}'
tests:
- accepted_values: # test na dopuszczalne wartości
values: ['placed', 'shipped', 'completed', 'return_pending', 'returned']
A oto jak ta dokumentacja wygląda na wygenerowanej stronie internetowej:

Makra i Moduły
Cel DBT nie polega tyle na tym, aby stać się zestawem skryptów SQL, ale na dostarczeniu użytkownikom potężnych i bogatych w możliwości narzędzi do budowania własnych transformacji oraz publikacji tych modułów.
Makra to zestawy konstrukcji i wyrażeń, które mogą być wywoływane jako funkcje wewnątrz modeli. Makra pozwalają na ponowne wykorzystanie SQL między modelami i projektami zgodnie z zasadą inżynieryjną DRY (Don’t Repeat Yourself).
Przykład makra:
{% macro rename_category(column_name) %}
case
when {{ column_name }} ilike '%osx%' then 'osx'
when {{ column_name }} ilike '%android%' then 'android'
when {{ column_name }} ilike '%ios%' then 'ios'
else 'other'
end as renamed_product
{% endmacro %}
I jego zastosowanie:
{% set column_name = 'product' %}
select
product,
{{ rename_category(column_name) }} -- wywołanie makra
from my_table
DBT dostarczany jest z menedżerem pakietów (packages), który umożliwia użytkownikom publikowanie i ponowne wykorzystywanie poszczególnych modułów i makr.
Oznacza to możliwość załadowania i używania takich bibliotek jak:
- : praca z Datą/Czasem, Kluczami Zastępczymi, testami schematu, Pivotem/Unpivotem i innymi
- Gotowe szablony witryn dla takich usług jak i
- Biblioteki dla określonych hurtowni danych, na przykład
- — Moduł do logowania pracy DBT
Pełną listę pakietów można znaleźć na .
Jeszcze więcej możliwości
Tutaj opiszę kilka innych interesujących cech i realizacji, które ja i zespół wykorzystujemy do budowy Magazynu Danych w .
Podział środowisk wykonawczych DEV — TEST — PROD
Nawet wewnątrz jednego klastra DWH (w ramach różnych schematów). Na przykład, za pomocą następującego wyrażenia:
with source as (
select * from {{ source('salesforce', 'users') }}
where 1=1
{%- if target.name in ['dev', 'test', 'ci'] -%}
where timestamp >= dateadd(day, -3, current_date)
{%- endif -%}
)
Ten kod dosłownie mówi: dla środowisk dev, test, ci weź dane tylko z ostatnich 3 dni i nie więcej. Oznacza to, że uruchamianie w tych środowiskach będzie znacznie szybsze i wymagało mniej zasobów. Przy uruchomieniu w środowisku prod warunek filtra zostanie zignorowany.
Materializacja z alternatywnym kodowaniem kolumn
Redshift to kolumnowa baza danych, która pozwala ustawić algorytmy kompresji danych dla każdej oddzielnej kolumny. Wybór optymalnych algorytmów może zmniejszyć zajmowaną przestrzeń na dysku o 20-50%.
Makro wykona polecenie ANALYZE COMPRESSION, stworzy nową tabelę z zalecanymi algorytmami kodowania kolumn wskazanymi kluczami segmentacji (dist_key) i sortowania (sort_key), przeniesie do niej dane, a w razie potrzeby usunie starą kopię.
Sygnatura makra:
{{ compress_table(schema, table,
drop_backup=False,
comprows=none|Integer,
sort_style=none|compound|interleaved,
sort_keys=none|List,
dist_style=none|all|even,
dist_key=none|String) }}
Logowanie uruchomień modeli
Do każdego wykonania modelu można przypisać haki (hooks), które będą wykonywane przed uruchomieniem lub natychmiast po zakończeniu tworzenia modelu:
pre-hook: "{{ logging.log_model_start_event() }}"
post-hook: "{{ logging.log_model_end_event() }}"
Moduł logowania pozwoli na rejestrowanie wszystkich niezbędnych metadanych w oddzielnej tabeli, na podstawie której później można przeprowadzać audyty i analizować problemy (bottlenecks).
Tak wygląda dashboard z danych logowania w Looker:

Automatyzacja utrzymania Magazynu
Jeśli korzystasz z jakichś rozszerzeń funkcjonalności używanego Magazynu, takich jak UDF (User Defined Functions), to wersjonowanie tych funkcji, zarządzanie dostępami i automatyczne wdrażanie nowych wydań jest bardzo wygodne do przeprowadzenia w DBT.
Używamy UDF w Pythonie do obliczania wartości hash, domen adresów e-mail oraz dekodowania masek bitowych (bitmask).
Przykład makra, które tworzy UDF w dowolnym środowisku uruchomieniowym (dev, test, prod):
{% macro create_udf() -%}
{% set sql %}
CREATE OR REPLACE FUNCTION {{ target.schema }}.f_sha256(mes "varchar")
RETURNS varchar
LANGUAGE plpythonu
STABLE
AS $$
import hashlib
return hashlib.sha256(mes).hexdigest()
$$
;
{% endset %}
{% set table = run_query(sql) %}
{%- endmacro %}
W Wheely używamy Amazon Redshift, który jest oparty na PostgreSQL. Dla Redshift ważne jest regularne zbieranie statystyk dla tabel oraz zwalnianie miejsca na dysku — odpowiednio polecenia ANALYZE i VACUUM.
W tym celu co noc wykonywane są polecenia z makra redshift_maintenance:
{% macro redshift_maintenance() %}
{% set vacuumable_tables=run_query(vacuumable_tables_sql) %}
{% for row in vacuumable_tables %}
{% set message_prefix=loop.index ~ " of " ~ loop.length %}
{%- set relation_to_vacuum = adapter.get_relation(
database=row['table_database'],
schema=row['table_schema'],
identifier=row['table_name']
) -%}
{% do run_query("commit") %}
{% if relation_to_vacuum %}
{% set start=modules.datetime.datetime.now() %}
{{ dbt_utils.log_info(message_prefix ~ " Vacuuming " ~ relation_to_vacuum) }}
{% do run_query("VACUUM " ~ relation_to_vacuum ~ " BOOST") %}
{{ dbt_utils.log_info(message_prefix ~ " Analyzing " ~ relation_to_vacuum) }}
{% do run_query("ANALYZE " ~ relation_to_vacuum) %}
{% set end=modules.datetime.datetime.now() %}
{% set total_seconds = (end - start).total_seconds() | round(2) %}
{{ dbt_utils.log_info(message_prefix ~ " Finished " ~ relation_to_vacuum ~ " in " ~ total_seconds ~ "s") }}
{% else %}
{{ dbt_utils.log_info(message_prefix ~ ' Skipping relation "' ~ row.values() | join ('"."') ~ '" as it does not exist') }}
{% endif %}
{% endfor %}
{% endmacro %}
DBT Cloud
Istnieje możliwość korzystania z DBT jako usługi (Managed Service). W zestawie:
- Web IDE do rozwijania projektów i modeli
- Konfiguracja zadań i harmonogramowanie
- Łatwy i wygodny dostęp do logów
- Strona internetowa z dokumentacją Twojego projektu
- Integracja CI (Continuous Integration)

Podsumowanie
Przygotowywanie i korzystanie z DWH staje się tak przyjemne i korzystne, jak picie smoothie. DBT składa się z Jinja, rozszerzeń użytkownika (modułów), kompilatora, silnika (executor) oraz menedżera pakietów. Łącząc te elementy, otrzymujesz pełnowartościowe środowisko robocze dla swojego magazynu danych. Dziś z pewnością nie ma lepszego sposobu zarządzania transformacjami w DWH.
Przekonania, którymi kierowali się twórcy DBT, formułują się następująco:
- Kod, a nie GUI, jest najlepszą abstrakcją do wyrażania skomplikowanej logiki analitycznej
- Praca z danymi powinna dostosowywać najlepsze praktyki inżynierii oprogramowania (Software Engineering)
- Najważniejsza infrastruktura do pracy z danymi powinna być kontrolowana przez społeczność użytkowników jak oprogramowanie z otwartym kodem źródłowym
- Nie tylko narzędzia analityczne, ale także kod coraz częściej stanie się własnością społeczności Open Source
Te podstawowe przekonania doprowadziły do powstania produktu, który dziś jest wykorzystywany przez ponad 850 firm, i stanowią podstawę wielu interesujących rozszerzeń, które będą tworzone w przyszłości.
Dla tych, którzy są zainteresowani, istnieje nagranie otwartego wykładu, który prowadziłem kilka miesięcy temu w ramach otwartego wykładu w OTUS — .
Oprócz DBT i magazynów danych, w ramach kursu Data Engineer na platformie OTUS, ja i moi koledzy prowadzimy zajęcia na temat różnych innych aktualnych i nowoczesnych tematów:
- Koncepcje architektoniczne aplikacji Big Data
- Praktyka z Spark i Spark Streaming
- Zgłębianie sposobów i narzędzi do ładowania źródeł danych
- Budowanie analitycznych witryn w DWH
- Koncepcje NoSQL: HBase, Cassandra, ElasticSearch
- Zasady organizacji monitorowania i orkiestracji
- Projekt końcowy: zbieramy wszystkie umiejętności pod mentorską opieką
Linki:
- — Oficjalna dokumentacja
- — Przeglądowy artykuł jednego z autorów DBT
- — YouTube, Nagranie otwartego wykładu OTUS
- — Najbliższy otwarty wykład 15 maja 2020
- — OTUS
- — Spojrzenie w przyszłość pracy z danymi i analityki
- — Ewolucja analityki i wpływ Open Source
- — Zasady budowy CI z wykorzystaniem DBT
- — Praktyka, Krok po kroku instrukcje do samodzielnej pracy
- — Github, kod projektu edukacyjnego
Źródło: habr.com

