Ga naar hoofdinhoud

Lake feed (change feed)

Naast het laden naar de SQL-database kan Yres elke tabel ook naar de Azure Data Lake schrijven. Vanaf v1.56 gebeurt dat als een append-only change feed: per laadrun landt er één Parquet-bestand met alleen de mutaties van die run, waarbij elke rij een markering I (insert), U (update) of D (delete) draagt.

Daarmee is de Data Lake niet langer een reeks losse momentopnames, maar een volwaardige mutatiestroom: je leidt er zowel de actuele stand als de volledige historie uit af, en je kunt hem rechtstreeks als invoer voor een MERGE naar een Delta-tabel gebruiken.

De database blijft de bron van waarheid

De feed is afgeleide data. De SCD2-historie in Azure SQL (STAGE → HIS) blijft leidend; de lake is een parallelle landing voor analytics-, data-science- en lakehouse-scenario's.

Wat er verandert ten opzichte van vroeger​

Vóór v1.56 kopieerde de stap Load DL simpelweg de volledige inhoud van de stagingtabel naar de Data Lake. Elke run schreef dus opnieuw álles weg — ook ongewijzigde rijen — en verwijderingen waren onzichtbaar: de rij verdween gewoon uit de volgende dump.

Vroeger (STAGE-dump)Nu (change feed)
Inhoud per runde volledige stagingtabelalleen de mutaties van die run
Run zonder wijzigingenschreef toch een bestandschrijft geen bestand
Ongewijzigde rijenwerden elke run opnieuw geschrevenworden nooit verstuurd
Verwijderde rijenonzichtbaarexpliciete D-rij (tombstone)
Historiealleen de laatst gedumpte standvolledig af te leiden uit alle bestanden
Groeischaalt met herlaadvolumeschaalt met mutatievolume
Pad<Bron>/<Schema>/<jaar>/<maand>lake/<Bron>/<Schema>/<Tabel>/Year=…/Month=…
Bestandsnaam<Tabel>-<tijdstip> (zonder extensie)<Tabel>-g<Generatie>-<PipelineRunId>.parquet
Herstart van een runleverde een extra bestand opoverschrijft het eigen bestand (binnen dezelfde generatie)
Wat dit betekent voor bestaande afnemers

De feed schrijft naar een nieuw pad. De bestanden die eerder onder <Bron>/<Schema>/<jaar>/<maand> zijn geland, blijven ongemoeid staan — er wordt niets gemigreerd of opgeruimd. Rapporten of notebooks die op het oude pad wezen én op een volledige momentopname per bestand rekenden, moeten worden aangepast: een feedbestand bevat immers alleen de mutaties. Gebruik daarvoor het leespatroon hieronder.

Wat er in de Data Lake landt​

datalake-yres
└── lake/<Bron>/<Schema>/<Tabel>/Year=<jjjj>/Month=<mm>/<Tabel>-g<Generatie>-<PipelineRunId>.parquet
  • Eén bestand per laadrun per tabel. Runs zonder mutaties schrijven niets.
  • Year= / Month= zijn hive-style partitiemappen, zodat elke query-engine op periode kan snoeien.
  • De bestandsnaam draagt het generatienummer (g<Generatie>, bijgehouden in LoadManagement.LakeFeedGeneration) en het ADF pipeline run-id: draai je dezelfde run opnieuw binnen dezelfde generatie, dan overschrijft hij zijn eigen bestand. Dubbele rijen door een herstart zijn dus uitgesloten.

Naast de gewone databkolommen en ETL_Date draagt elke rij vier framework-kolommen:

KolomBetekenis
KeyHashSHA2_512-hash over de sleutelkolommen — de stabiele identiteit van de rij.
RowHashSHA2_512-hash over de gevolgde kolommen van déze versie.
YresActionI = nieuwe sleutel, U = nieuwe versie van een bestaande sleutel, D = sleutel verwijderd.
YresDateStartHet moment waarop deze versie actueel werd; bij een D-rij het moment van verwijderen.
D-rijen dragen geen data

Een tombstone bevat de hashes en YresDateStart, maar de businesskolommen zijn leeg (NULL). Hij zegt "deze sleutel bestaat niet meer", niet "deze sleutel had deze waarden".

D-rijen verschijnen alleen bij de laadtypen die verwijderingen kunnen detecteren — IMAGE, DELTAIMAGE, OVERWRITE en RELOAD. Zie Load types. Het loadtype ADDITIONAL kent geen mutatiebegrip: dat levert elke run de volledige stagingtabel als I-rijen aan (pure append).

Aanzetten​

De feed hangt aan de bestaande kolom DataPlatform op de tabelconfiguratie ([LoadManagement].[UsedTables]):

WaardeGevolg
DWHalleen de SQL-database (standaard)
DLalleen de lake feed
DWH,DLallebei

Bij een DL-only-tabel (DL zonder DWH) maakt Yres bewust géén history-tabel in de database aan — de data leeft volledig in de Data Lake. Bevragen vanuit SQL kan toch: Yres onderhoudt daarvoor automatisch leesobjecten in het [DL]-schema (zie hieronder).

Zet je DL weer uit, dan blijven de al geschreven Parquet-bestanden staan; een health check wijst je op de achtergebleven administratie in de database (zie Bewaking).

Deploy-volgorde

De ADF-pipelines roepen procedures aan die met de database meekomen. Werk daarom altijd eerst de database bij (DACPAC) en publiceer daarna pas de ADF-factory. Andersom faalt de lake-tak.

Hoe Yres de mutaties bepaalt​

Om te weten wát er in een run veranderd is, houdt Yres per lake-tabel een slanke administratie bij in het schema [LAKE] (instelbaar met de setting SchemaLAKE). Die tabel bevat geen businessdata — alleen de hashes, de datums en de eventuele deltakolom.

Bron ──Copy──► STAGE.<Tabel>
│
├─ Prepare lake load → [LoadManagement].[spLoadLake]
│ └→ spHIS_InsertAndUpdate @LakeMode = 1 (SCD2-merge op de slanke [LAKE]-tabel)
│
├─ Load DWH → [LoadManagement].[spLoadDWH] (de gewone SCD2-merge naar HIS)
│
├─ Lookup lake feed → [LoadManagement].[spGetLakeFeed]
│ └→ mutatiequery + aantal mutaties
│
└─ Write lake feed → Copy → Parquet in de Data Lake (alleen als er mutaties zijn)
└→ Maintain lake view → [LoadManagement].[spMaintainLakeExternal] (Scope=VIEW)

Het is dus dezelfde, bewezen SCD2-merge die de historie in de database opbouwt, hier toegepast op een administratietabel zonder inhoud. Wat de merge als nieuw of gewijzigd markeert, is precies wat de feed verstuurt; wat hij afsluit zonder tegenhanger in de staging, wordt een D-rij. Levert een run nul mutaties op, dan slaat ADF de kopieerstap over en ontstaat er geen leeg bestand. Na een geslaagde kopie werkt de stap Maintain lake view (spMaintainLakeExternal met scope VIEW) de external view over de feed-bestanden bij.

De stap Write lake feed loopt parallel aan Load DWH, niet erna: de lake-uitvoer vertraagt het laden van het datawarehouse dus niet.

De feed lezen​

De actuele stand haal je uit de feed door per sleutel de laatste versie te nemen en tombstones weg te laten. Dit patroon werkt in elke engine:

WITH ranked AS (
SELECT *, ROW_NUMBER() OVER (PARTITION BY KeyHash ORDER BY YresDateStart DESC) AS rn
FROM <feed>
)
SELECT * FROM ranked WHERE rn = 1 AND YresAction <> 'D';

De volledige historie is simpelweg álle rijen; de einddatum van een versie leid je af met LEAD(YresDateStart) OVER (PARTITION BY KeyHash ORDER BY YresDateStart).

PlatformHoe je leest
Microsoft FabricOPENROWSET(BULK '…/lake/<Bron>/<Schema>/<Tabel>/**', FORMAT='parquet') in een Warehouse — geen Spark nodig. Voor Direct Lake vouw je de feed met Spark naar een Delta-tabel.
Databricksread_files(…, format => 'parquet'), of met Auto Loader incrementeel vouwen naar Delta: MERGE op KeyHash, D-rijen als delete.
Azure SQL DatabaseVoor DL-only-tabellen bouwt Yres dit zelf: kant-en-klare views in het [DL]-schema (zie hieronder). Handmatig kan het ook, via data virtualization (OPENROWSET over een external data source). Vereist een managed identity op de SQL-server met Storage Blob Data Reader — dezelfde inrichting als de archief-unionviews.
Azure SQL: data virtualization is preview

Lezen vanuit Azure SQL werkt, maar is een preview-feature van Azure SQL Database. Let daarbij op: gebruik altijd een external data source (een kale URL in BULK eist alsnog een credential), gebruik het adls://-schema (niet https://), en geef kolomtypen expliciet op. Wijst het pad naar een map zonder bestanden, dan volgt een foutmelding in plaats van een lege resultaatset.

DL-only-tabellen bevragen vanuit SQL: het [DL]-schema​

Voor DL-only-tabellen genereert en onderhoudt [LoadManagement].[spMaintainLakeExternal] automatisch drie soorten leesobjecten in het [DL]-schema (instelbaar met de setting SchemaDL):

ObjectWat het is
[DL].[<Tabel>_Feed_g<N>]External table per schemageneratie — leest de Parquet-bestanden van die generatie rechtstreeks uit de Data Lake. Bouwsteen; normaal bevraag je deze niet zelf.
[DL].[<Tabel>_Feed]Union-view over alle generaties — de complete, ruwe change feed als één tabel.
[DL].[<Tabel>]HIS-vormige view — leidt uit de feed de vertrouwde SCD2-vorm af, zodat je een DL-only-tabel precies zo bevraagt als een gewone history-tabel.

_Feed of niet? Welke view je gebruikt​

De twee views serveren dezelfde onderliggende bestanden, maar beantwoorden een andere vraag:

  • [DL].[<Tabel>_Feed] — "wat is er gebeurd?" Eén rij per mutatie, precies zoals die in de lake geland is: YresAction (I/U/D), YresDateStart, KeyHash/RowHash. Je ziet álle versies van elke sleutel, inclusief de D-tombstones (met lege businesskolommen). Dit is de invoer voor mutatieverwerking: een MERGE naar een eigen tabel, een incrementele fold, een audit op wat een run precies deed.
  • [DL].[<Tabel>] — "wat is de stand (en de historie)?" Dezelfde feed, maar gepresenteerd als SCD2-history-tabel: YresDateStart wordt ETL_Date, de einddatum van elke versie wordt afgeleid (ETL_EndDate, open versies krijgen 2999-01-01), en per sleutel is de nieuwste versie IsCurrent = 1. Tombstones zie je hier niet als rij — een verwijderde sleutel is gewoon afwezig in de actuele stand, maar zijn laatste versie is er wél netjes door afgesloten, precies zoals een delete dat in een echte HIS-tabel doet. De actuele stand is dus simpelweg WHERE IsCurrent = 1.

Vuistregel: verwerk je mutaties → _Feed; bevraag je de tabel → [DL].[<Tabel>]. Rapportages en ad-hoc-queries horen vrijwel altijd op de HIS-vormige view; alleen consumenten die zelf iets met de mutatiestroom opbouwen hebben _Feed nodig.

Twee details van de HIS-vormige view:

  • Ze bevat extra kolommen YresYear en YresMonth (uit de partitiemappen in de lake), zodat je ziet uit welke periode-map een versie komt.
  • Voor een ADDITIONAL-tabel is er geen versiebegrip: elke rij is en blijft IsCurrent = 1 (pure append), alleen eventuele tombstones uit een eerder laadtype worden weggefilterd. Bij dit laadtype snoeit een filter op YresYear/YresMonth bovendien op bestandsniveau. Bij de andere laadtypen kan dat niet: de versie-afleiding heeft per sleutel de héle geschiedenis nodig (de tombstone die een versie afsluit kan in een andere maand liggen), dus daar leest de view altijd alle bestanden.

Wijzigt het bronschema (een nieuwe generatie), dan komen de objecten automatisch mee: de externe tabellen worden per generatie bijgehouden tijdens de load en de views ververst na elke geslaagde kopie — er is geen handwerk nodig. Een generatie die een kolom nog niet kende, toont daar NULL — in beide views.

Eenmalige inrichting per omgeving: de instellingen SchemaDL en LakeLocation (adls://<container>@<account>.dfs.core.windows.net), een managed identity op de SQL-server met Storage Blob Data Reader op de Data Lake, en database compatibility level 130 of hoger — de health checks wijzen je erop als er iets ontbreekt. Data virtualization is een preview-feature van Azure SQL Database.

Opnieuw opbouwen​

Raakt de administratie uit de pas — bijvoorbeeld na handmatig ingrijpen of doordat er even geen feed geschreven werd — dan wis je hem voor één tabel of voor alles:

EXEC [LoadManagement].[spResetLakeIndex] @Target = 'Bron_Schema_Tabel'; -- leeg = alle lake-tabellen

De eerstvolgende load verstuurt dan de complete actuele dataset opnieuw als I-rijen. Consumenten die de actuele stand via het patroon hierboven bepalen, merken daar niets van; een Delta-fold merget de verse rijen gewoon overheen.

Bewaking​

Een reeks health checks bewaakt de feed en de leesobjecten:

CheckSignaleert
2.12Een tabel staat op DL en heeft succesvolle loads, maar er is geen administratie in het [LAKE]-schema — er wordt voor die tabel dus geen feed geproduceerd. Meestal draait de ADF-factory nog een oudere versie.
2.13Er staat nog een [LAKE]-administratietabel voor een tabel die niet meer op DL staat. De check levert een opruimscript; de Parquet-bestanden blijven ongemoeid.
2.14Een DL-only-tabel is geladen, maar mist (een deel van) zijn leesobjecten in het [DL]-schema — de feed is dan niet volledig vanuit SQL te bevragen.
2.15Er staan nog [DL]-leesobjecten voor een tabel die geen actieve DL-only-tabel meer is. De check levert een opruimscript; de Parquet-bestanden blijven ongemoeid.
2.16Er zijn actieve DL-only-tabellen, maar SchemaDL of LakeLocation is nog niet geconfigureerd.
2.17Het database compatibility level is lager dan 130, terwijl external tables/OPENROWSET dat vereisen. Inclusief kant-en-klaar herstelscript.
2.18Recente meldingen dat het onderhoud van de leesobjecten is mislukt of overgeslagen — met de details in de monitoring.

Groei​

De feed groeit mee met het aantal mutaties, niet met het aantal herladingen: een tabel die dagelijks volledig herladen wordt maar nauwelijks wijzigt, levert nauwelijks bestanden op. Bij hoge mutatievolumes is het gebruikelijk de feed aan consumentzijde periodiek te vouwen naar een Delta-tabel of een snapshot; Yres comprimeert de feed zelf niet.

Zie ook​