Ugrás a fő tartalomra

DataBricks alapok


DATABRICKS alapok

Ugyanaz az adat. Nagyobb lehetőségek.




Felépítés • Átalakítás • Elemzés • Automatizálás • Skálázás

1. Munkakörnyezet alapjai (Workspace Basics)

Gyakori felületi elemek

  • Workspace (Munkaterület): Jegyzetfüzetek (Notebooks), Git-tárolók (Repos), Feladatok (Jobs) stb. kezelése.

  • Compute (Számítási kapacitás): Klaszterek kezelése (általános célú és feladat-alapú).

  • Data (Adatok): Hozzáférés a DBFS-hez, a Unity Catalog-hoz és a relációs táblákhoz.

  • Workflows (Munkafolyamatok): Feladatok ütemezése és orkesztrációja.

  • Repos: Verziókövetési és Git-integrációs felület.

  • SQL Warehouses: Lekérdező motorok BI eszközökhöz és SQL analitikához.

  • ML (Gépi tanulás): Modellek tanítása, kísérletkövetés és modellkiszolgálás (Serving).

2. Klasztertípusok (Cluster Types)

A megfelelő számítási erőforrás kiválasztása

KlasztertípusFelhasználási eset
All-Purpose (Általános célú)Interaktív fejlesztés, jegyzetfüzetek futtatása és adatelemzés.
Job Cluster (Feladat klaszter)Automatizált adatfolyamatok futtatása (költséghatékony, leáll a futás végén).
SQL WarehouseSQL analitika és külső BI eszközök (Power BI, Tableau) kiszolgálása.
Serverless (Szerver nélküli)Nincs manuális klaszterkezelés; a Databricks automatikusan és azonnal skálázza.

3. Kulcsfogalmak (Key Concepts)

Alapvető fogalmak és kifejezések

  • Delta Lake: Megbízható adattárolási réteg ACID tranzakciókkal és időutazás (time travel) funkcióval.

  • Unity Catalog: Egységes adat- és AI-irányítás, biztonsági szabályozás és adathozzáférés-kezelés.

  • DBFS (Databricks File System): Elosztott fájlrendszer-absztrakció a felhőtárhelyek felett.

  • Volumes (Kötetek): Fájlok (nem táblázatos adatok) tárolása és kezelése a Unity Catalogban.

  • Metastore: Központosított metaadat-tárhely a táblákhoz és sémákhoz.

  • Delta Table: Nyílt forráskódú, Parquet alapú, optimalizált tranzakciós táblaformátum.

  • Notebook: Interaktív fejlesztői környezet kód futtatásához és vizualizációhoz.

  • Job: Munkafolyamatok ütemezése és kötegelt futtatása.

  • Repo: Kétirányú szinkronizáció külső Git-szolgáltatókkal (GitHub, GitLab, Azure DevOps).

  • Widget: Paraméterezési lehetőség jegyzetfüzetekhez interaktív bemeneti mezőkkel.

4. Spark (Python) alapok

Gyakori PySpark műveletek

Python
# DataFrame létrehozása / beolvasása CSV-ből
df = spark.read.csv("/path/file.csv", header=True, inferSchema=True)

# Adatok megjelenítése (interaktív táblázat Databricks alatt)
display(df)
# vagy egyszerű szöveges konzolos nézet:
df.show()

# Oszlopok kiválasztása és sorok szűrése
df.select("id", "name").filter("age > 18").show()

# Új oszlop hozzáadása függvény használatával
from pyspark.sql.functions import col, year
df = df.withColumn("year", year(col("date")))

# Csoportosítás és aggregálás
df.groupBy("category").count().show()

5. SQL a Databricks-ben

Hasznos SQL parancsok

SQL
-- Adatbázis (Schema) létrehozása és kiválasztása
CREATE DATABASE IF NOT EXISTS sales;
USE sales;

-- Delta formátumú tábla létrehozása
CREATE TABLE employees (
    id INT,
    name STRING,
    dept STRING
) USING DELTA;

-- Adatok lekérdezése szűréssel
SELECT * FROM employees
WHERE dept = 'Data Engineering';

-- Időutazás (Time Travel) korábbi verzió vagy időbélyeg alapján
SELECT * FROM employees VERSION AS OF 2;
-- vagy: SELECT * FROM employees TIMESTAMP AS OF '2026-01-01 00:00:00';

6. Delta Lake műveletek

Megbízható és skálázható adattárolás

Python
# Írás Delta táblába (felülírással)
df.write.format("delta").mode("overwrite").save("/mnt/data/employees")

# Delta tábla beolvasása
df = spark.read.format("delta").load("/mnt/data/employees")

# Merge (Upsert: frissítés vagy beszúrás) művelet
from delta.tables import DeltaTable

deltaTable = DeltaTable.forPath(spark, "/mnt/data/employees")
deltaTable.alias("t").merge(
    df.alias("s"),
    "t.id = s.id"
).whenMatchedUpdateAll(
).whenNotMatchedInsertAll(
).execute()

# Optimalizálás (kis fájlok tömörítése) és törölt verziók takarítása
deltaTable.optimize().executeCompaction()
deltaTable.vacuum(retentionHours=168)  # 7 napos megőrzési idő (7 * 24 = 168 óra)

7. Databricks segédprogramok (dbutils)

Hasznos beépített parancsok jegyzetfüzetekhez

Python
# Fájlrendszer műveletek
dbutils.fs.ls("/mnt/data/")                                      # Fájlok listázása
dbutils.fs.cp("/mnt/source", "/mnt/target", recurse=True)        # Másolás rekurzívan
dbutils.fs.rm("/mnt/old", recurse=True)                          # Törlés rekurzívan

# Paraméter-widgetek kezelése
dbutils.widgets.text("date", "", "Load Date")                   # Szöveges beviteli mező
load_date = dbutils.widgets.get("date")                          # Érték kiolvasása

# Titkos kulcsok (Secrets) elérése
dbutils.secrets.listScopes()                                     # Titok-tartományok listázása
dbutils.secrets.get(scope="my_scope", key="my_key")             # Titkos érték lekérése

8. Feladatok és munkafolyamatok (Jobs & Workflows)

Adatfolyamatok automatizálása

  • Többlépéses feladatok létrehozása: Jegyzetfüzetek, Python scriptek, JAR vagy SQL feladatok összekapcsolása egyetlen folyamatba.

  • Klaszterhozzárendelés: Meglévő klaszter használata vagy ideiglenes feladat-klaszter (Job cluster) automatikus indítása.

  • Ütemezés beállítása: Időzítés Cron-kifejezéssel vagy a kezelőfelületen (UI).

  • Függőségek definiálása: Feladatok közötti függőségi háló (DAG) kialakítása.

  • Futások monitorozása: Hibák és lefutások naplózása, e-mail vagy webhook értesítések küldése.

  • Strukturált folyamatlánc:

    $$\text{Ingest (Beolvasás / Notebook)} \longrightarrow \text{Transform (Átalakítás / PySpark)} \longrightarrow \text{Load (Betöltés / Delta)}$$


9. Unity Catalog

Adatkezelés, hozzáférés-vezérlés és biztonság

  • Központosított hozzáférés-vezérlés: Hierarchikus jogosultságkezelés (metastore $\rightarrow$ catalog $\rightarrow$ schema $\rightarrow$ table/volume).

  • Finomhangolt jogosultságok: Oszlop- és sorszintű szűrések támogatása.

  • Munkaterületeken átívelő működés: Egyetlen központi katalógus több Databricks Workspace-hez rendelhető.

  • Kötetek (Volumes): Strukturálatlan és félig strukturált fájlok biztonságos tárolása (a DBFS helyett preferált megoldás).

  • Felhő-azonosítás integrációja: Microsoft Entra ID (korábban Azure AD) és SSO integráció.

  • Adat-eredet (Lineage): Teljes életút- és auditnaplózás lekérdezési és táblaszinten.

  • Jogosultsági példa SQL-ben:

    SQL
    GRANT SELECT ON TABLE sales.employees TO `data_analyst_group`;
    

10. Gyakori felhasználási esetek (Common Use Cases)

Gyakorlati forgatókönyvek az iparban

  • Kötegelt adatfolyamat (Batch Pipeline): Adatbeolvasás $\rightarrow$ átalakítás $\rightarrow$ cél Delta táblába írás.

  • Inkrementális betöltés: Változások folyamatos feldolgozása a Delta Lake használatával.

  • Adatminőség-ellenőrzés (Data Quality): Előfeltételek és elvárások érvényesítése (pl. Delta Live Tables Expectations).

  • CDC (Change Data Capture): Forrásrendszerek adatváltozásainak folyamatos átvezetése.

  • Orkesztráció Jobs segítségével: Ütemezett, függőségeken alapuló munkafolyamatok automatizálása.

  • SQL Analitika: Riportkészítés és dashboard-kiszolgálás (Power BI, Tableau).

  • Gépi tanulási folyamatok (ML Pipelines): Adat-előkészítés, Feature Store, modelltanítás és verziókezelés (MLflow).

  • Valós idejű adatfolyamok: Alacsony késleltetésű streamelés a Spark Structured Streaming modullal.

  • Csapatok közötti adatmegosztás: Biztonságos megosztás belső vagy külső partnerekkel a Delta Sharing technológiával.

  • Költségoptimalizálás: Automatikus klaszterleállítás (Auto-terminate) és optimális feladat-alapú erőforrások beállítása.

11. Teljesítmény-optimalizálási tippek (Performance Tips)

A hatékony futtatás szabályai

  • Delta formátum használata: A nyers CSV/JSON fájlok helyett mindig Delta táblákba kell írni.

  • Particionálás: Nagy táblák kulcs szerinti particionálása (kerülve a túlzott al-particionálást).

  • Kis fájlok elkerülése: Rendszeres OPTIMIZE futtatása az apró fájlok összevonására.

  • Megfelelő klaszterméret: A feladat memória- és CPU-igényéhez igazított géptípusok választása.

  • Adatok gyorsítótárazása (Cache): Többször újrahasznált DataFrame-ek gyorsítótárba helyezése (.cache()).

  • Broadcast Join alkalmazása: Kis méretű táblák csatolásakor hálózati adatmozgatás elkerülése (broadcast() hint).

  • Spark UI monitorozása: Szűk keresztmetszetek és adateltérések (data skew) azonosítása.

  • Job klaszterek használata: Éles folyamatokhoz kizárólag feladat-klaszterek futtatása a költségcsökkentés érdekében.

  • Felesleges adatok takarítása: Régi verziók és árva fájlok tisztítása a VACUUM paranccsal.

  • Photon motor engedélyezése: A C++ alapú végrehajtó motor bekapcsolása a gyorsabb SQL és DataFrame lekérdezésekhez.

12. Gyorsbillentyűk és navigáció (Useful Shortcuts & Links)

Időtakarékos vezérlők jegyzetfüzetekhez
MűveletGyorsbillentyű / Útvonal
Cella futtatásaShift + Enter
Minden cella futtatásaCtrl + Shift + Enter
Futás megszakításaEsc
Parancsmód váltásEsc
Keresés és csereCtrl + F
Sor kommentelése / feloldásaCtrl + /
Parancspaletta megnyitásaCtrl + P
Feladatok (Jobs) megtekintéseOldalsáv: Workspace $\rightarrow$ Workflows
Klaszternaplók megtekintéseOldalsáv: Compute $\rightarrow$ Adott klaszter $\rightarrow$ Driver Logs




DATABRICKS HALADÓ ARCHITEKTÚRA ÉS AI


Modern Adatplatform • Medallion • MLOps • Adatkormányzás

Streaming • Minőségbiztosítás • FinOps • Gépi Tanulás

1. Medallion Architektúra (Többrétegű feldolgozás)

Strukturált adatfeldolgozási szintek a Lakehouse-ban

  • Bronze (Nyers réteg): A forrásrendszerek változatlan, append-only formátumú adatait tárolja a teljes előzmény megőrzésével.

  • Silver (Tisztított réteg): Tisztított, transzformált, gazdagított és konzisztens sémájú adatok a vállalati nézetekhez.

  • Gold (Üzleti réteg): Üzleti logikára és aggregációkra optimalizált adatok közvetlen BI és analitikai felhasználásra.

2. Auto Loader (cloudFiles)

Inkrementális, nagy hatékonyságú fájlbetöltés felhőtárhelyről

Python
# Automatikus sémadetektálás és változáskövetés felhőtárhelyről
df = (spark.readStream
      .format("cloudFiles")
      .option("cloudFiles.format", "json")
      .option("cloudFiles.schemaLocation", "/mnt/telemetry/_schema")
      .load("/mnt/telemetry/raw_data/"))

# Folyamatos írás Bronze Delta táblába
(df.writeStream
   .format("delta")
   .option("checkpointLocation", "/mnt/telemetry/_checkpoints")
   .start("/mnt/telemetry/bronze"))

3. Delta Live Tables (DLT)

Deklaratív adatfolyamok és automatikus minőség-ellenőrzés

  • Deklaratív kódolás: A rendszer automatikusan kezeli a függőségi gráfot (DAG) és a végrehajtást.

  • Adatminőségi elvárások (Expectations): Automatikus rekord-érvényesítés hibakezelési szabályokkal:

    • EXPECT (feltétel): Hibás sorok rögzítése, de a folyamat fut tovább.

    • EXPECT (feltétel) ON VIOLATION DROP ROW: Érvénytelen sorok automatikus eldobása.

    • EXPECT (feltétel) ON VIOLATION FAIL UPDATE: A teljes folyamat azonnali megszakítása hiba esetén.

Python
import dlt
from pyspark.sql.functions import col

@dlt.table(comment="Tisztított ügyféladatok minőség-ellenőrzéssel")
@dlt.expect_or_drop("valid_email", "email IS NOT NULL AND email LIKE '%@%'")
def customers_silver():
    return dlt.read_stream("customers_bronze").filter(col("status") == "ACTIVE")

4. Structured Streaming & CDC

Valós idejű feldolgozás és változáskövetés

  • Strukturált adatfolyam: Micro-batch alapon működő, hibaálló adatfeldolgozás checkpoint mechanizmussal.

  • CDC feldolgozás: Forrásadatbázisok módosításainak leképezése MERGE vagy apply_changes segítségével.

  • Vízjel (Watermarking): Későn érkező adatok kezelése és az állapotmemória (state store) automatikus tisztítása:

Python
# Későn érkező események kezelése 10 perces időablakkal
streaming_df = (df.withWatermark("event_time", "10 minutes")
                  .groupBy("event_type")
                  .count())

5. Liquid Clustering & Speciális Delta funkciók

Új generációs táblaoptimalizálás és verziókezelés

  • Liquid Clustering: Helyettesíti a hagyományos particionálást és a Z-Order-t; rugalmas, többdimenziós elrendezést biztosít adatújraírás nélkül.

  • Tábla klónozás: Költséghatékony tesztkörnyezetek és pillanatképek létrehozása.

  • Változásnaplózás (CDF): Soronkénti módosítások lekérdezése downstream folyamatokhoz.

SQL
-- Új tábla létrehozása Liquid Clustering technológiával
CREATE TABLE sales_events (
    id BIGINT,
    customer_id STRING,
    event_date DATE
)
USING DELTA
CLUSTER BY (event_date, customer_id);

-- Shallow Clone létrehozása azonnali, tárhelymentes teszteléshez
CREATE TABLE sales_test SHALLOW CLONE sales_events;

6. Biztonság és Kormányzás a Unity Catalogban

Sorszintű szűrés és dinamikus adatmaszkolás

  • Dinamikus adatmaszkolás: Érzékeny mezők kitakarása szerepkör alapján.

  • Sorszintű szűrés (Row Filtering): Felhasználói csoportokhoz kötött láthatóság.

  • Delta Sharing: Nyílt forráskódú adatmegosztási protokoll másolás és zárt felhőfüggőség nélkül.

SQL
-- Dinamikus maszkoló függvény definiálása és hozzárendelése
CREATE OR REPLACE FUNCTION email_mask(email STRING)
RETURN IF(IS_ACCOUNT_GROUP_MEMBER('admin_group'), email, '****@****.com');

ALTER TABLE customers ALTER COLUMN email SET MASK email_mask;

7. MLOps és MLflow

Gépi tanulási életciklus-kezelés

  • Kísérletkövetés (Tracking): Paraméterek, mérőszámok és kimeneti modellek automatikus naplózása.

  • Modellregiszter (Unity Catalog): Modellek jóváhagyási folyamatának és állapotainak kezelése.

  • Feature Store: Újrafelhasználható jellemzők központosított tárolása és kiszolgálása.

Python
import mlflow

with mlflow.start_run():
    mlflow.log_param("max_depth", 5)
    mlflow.log_metric("accuracy", 0.94)
    mlflow.spark.log_model(model, "random_forest_model")

8. Generatív AI és RAG alapok

Nagy nyelvi modellek és vektoros keresés a Databricks platformon

  • Mosaic AI Model Serving: Nyílt forráskódú és zárt LLM-ek skálázható kiszolgálása REST végponton.

  • Vector Search: Beágyazások (embeddings) tárolása és alacsony késleltetésű hasonlósági keresés.

  • Lakehouse Monitoring: Generált kimenetek és adatminőség automatikus drift-analízise.

SQL
-- Beépített AI függvény hívása strukturált SQL lekérdezésben
SELECT 
    ticket_id, 
    ai_classify(comment, ARRAY('Panasz', 'Érdeklődés', 'Dicséret')) AS kategoriak
FROM customer_support;

9. CI/CD és Fejlesztői Eszközök

Kód-alapú infrastruktúra és automatizáció

  • Databricks Asset Bundles (DABs): Teljes Lakehouse projektek (Jobs, DLT, infra) leírása YAML alapon és élesítése CLI-vel.

  • Databricks Connect: Helyi fejlesztői környezetek (VS Code, PyCharm) összekötése távoli Databricks klaszterekkel.

  • Git Provider Szinkronizáció: Branch-alapú fejlesztés és tesztelés közvetlenül a Workspace-en belül.


10. Költségoptimalizálás és FinOps

Erőforrások és költségkeretek felügyelete

  • Spot példányok használata: Akár 70–80%-os költségmegtakarítás a hibatűrő worker csomópontokon.

  • Klaszter házirendek (Cluster Policies): Méretbeli és típusbeli korlátozások beállítása fejlesztői csoportoknak.

  • Rendszertáblák (System Tables): Költség- és futási metrikák elemzése natív SQL felületről.

SQL
-- DBU felhasználás és költségek vizsgálata rendszertáblákból
SELECT 
    workspace_id,
    sku_name,
    SUM(usage_quantity) AS total_dbu
FROM system.billing.usage
WHERE usage_date >= CURRENT_DATE() - INTERVAL 30 DAYS
GROUP BY workspace_id, sku_name;
 









Megjegyzések