INGEST-lagret
INGEST-lagret hämtar data ur källsystemen och landar den i datasjön. Det hanterar vad, när och hur i extraktionen.
Den här sidan är referensen för vad lagret gör. För den skärmvisa guiden till att konfigurera det, se Arbeta med Ingest.
Ingest-lagret ansvarar för att hämta data från källsystem och landa den i datasjön. Det hanterar vad, när och hur för dataextraktion.
Anslutningstyper
Innan du kan skapa exporter måste en anslutning konfigureras för ditt källsystem. Anslutningstypen avgör vilka fält anslutningen bär; inom en typ kör körningen exakt en operatorklass per källsystem, och operatorn avgör vilka av fälten som faktiskt läses.
| Anslutningstyp | Källa | Operatorer |
|---|---|---|
| Database (DB) | SQL-databaser som nås med SQL | PostgreSQL, SQL Server, Oracle, MySQL |
| API | HTTP-ändpunkter som returnerar JSON | Generisk REST, Heartpace, Salesforce Service Cloud, Talkdesk, SharePoint (Graph), Microsoft Teams, Google Drive / Sheets |
| File | Fillagringssystem | Lokal disk, Amazon S3, Azure Data Lake Gen2, SFTP |
| Custom | Meddelandeköer och skräddarsydda källor | AWS SQS, Azure Service Bus, Redis, samt kundspecifika kopplingar |
| Manual | Ingen automatiserad extraktion | Data uppladdad manuellt |
Köoperatorerna levereras under Custom-kontraktet, så en kös koordinater färdas som key=value-anslutningsegenskaper i stället för som egna fält. Inget i köbeteendet ändras av det.
Tre anslutningsinställningar formar varje export som görs från anslutningen:
| Inställning | Effekt |
|---|---|
AddServiceColVal | När den är satt får varje exporterad post en __service-kolumn med detta värde. Tomt värde stänger av kolumnen. |
ColumnDelimiter | Fältavgränsare när CDC skriver CSV-utdata. Ignoreras för JSON. |
LineDelimiter | Postavgränsare som skickas till CDC. Alla operatorer utom den generiska databasoperatorn tvingar radbrytning oavsett detta värde. |
Den fullständiga fältreferensen för anslutningar — autentiseringsuppgifter, autentiseringsmetoder, lagringsinställningar och kökoordinater — finns i utvecklarguiden för Ingest.
Exporttyper
Varje källsystem kan ha en eller flera källexporter. Exporttypen matchar anslutningstypen, och fält som operatorn ignorerar visas men märks som oanvända.
Databasexporter
Extrahera data med SQL-frågor från anslutna databaser. Frågan strömmas i block om 25 000 rader och skrivs som radbrytningsseparerad JSON.
| Inställning | Beskrivning |
|---|---|
databaseName | Referens till den anslutna databasen. Skrivskyddad, hämtas från anslutningen. |
fromClause | Tabell, vy, funktion eller lagrad procedur att fråga. Måste vara schemakvalificerad (dbo.Customers) — metadatauppslaget delar på punkten. |
whereClause | Valfritt SQL-villkor. Infogas ordagrant, så det måste innehålla nyckelordet WHERE. Ett inkrementellt filter läggs till med AND. |
sqlOverride | Fullständig anpassad SQL-fråga — ersätter den genererade frågan helt. Tidstoken substitueras fortfarande, men det inkrementella filtret läggs inte till; bygg in det i frågan själv. |
inclTables | Ytterligare tabeller vars kolumnmetadata hämtas för typkonvertering. Standard är from-satsen. Använd den för att lista varje objekt bakom en vy eller procedur, så att kolumntyperna löses upp rätt. |
inclColumns | Sammanfogas med kommatecken rakt in i SELECT-satsen, så alias och uttryck fungerar. ["*"] väljer allt. |
filterColumn | Aktiverar inkrementell laddning. Jämförelsen är numerisk, så kolumnen måste vara ett stigande tal. |
lastFilterValue | Vattenmärket från föregående körning, som postas tillbaka till plattformen efter varje körning. 0 vid första körningen. |
Två lägen:
- Standardläge — bygg frågan från
fromClause+whereClause+ kolumnval - SQL Override-läge — skriv en komplett anpassad SQL-sats som förbigår den genererade frågan
Databasexporter beräknar också en __checksum per post, vilket aktiverar den snabba CDC-jämförelsen.
API-exporter
Extrahera data från HTTP-ändpunkter. Utdata är alltid radbrytningsseparerad JSON — outputFormat ignoreras här.
| Inställning | Beskrivning |
|---|---|
apiUrl | Bas-URL för API:et. Skrivskyddad, hämtas från anslutningen. |
apiPath | Ändpunktens sökväg, som läggs till bas-URL:en. |
filterCondition | Nyckel=värde-par som skickas som förfrågningsparametrar — en query-sträng, eller en JSON-body när anslutningens method-header är post. |
Filtervillkor är strängarrayer som ["user_status=active", "date_from=2024-01-01"]. Tidstoken substitueras i värdena, och ett värde som innehåller klamrar tolkas som JSON, så en JSON-body kan konfigureras som ett filtervillkor.
Fel: en 401 utlöser en förnyelse av token och ett omförsök. Övriga statuskoder från 400 och uppåt får arbetslasten att misslyckas, liksom ett svar utan detekterbar innehållstyp.
Operatorspecialiseringar hanterar pagineringar och jobbmodeller som den generiska operatorn inte klarar:
| Operator | Pagineringsmodell |
|---|---|
| Heartpace | limit/offset, styrt av svarets meta-block |
| Salesforce Service Cloud | Följer nextRecordsUrl tills resultatmängden rapporterar sig färdig |
| Talkdesk | Asynkrona rapportjobb, som pollas var tionde sekund i upp till 50 minuter |
| SharePoint (Graph) | Rekursiv mappgenomgång, sorterad nyast först med tidigt avbrott |
| Teams, Google Drive | En enda förfrågan som returnerar en values-array vars första rad är rubriken |
Filexporter
Extrahera data från fillagringssystem. Matchande filer kopieras lokalt och konverteras sedan med DuckDB eller skickas vidare orörda.
| Inställning | Beskrivning |
|---|---|
filePath | Källkatalogens prefix, som läggs till anslutningens bassökväg. Tidstoken substitueras här, så en partitionerad källa kan adresseras med [__SCHEDULEDTIME_AS_HIVE__]. |
fileName | En delsträngsmatchning mot filnamnet — inte ett glob-mönster. |
filterCondition | Filurval och läsaralternativ (nedan). |
outputFormat | json eller csv konverterar filerna; binary laddar upp dem orörda och hoppar över CDC. Konverterad utdata skrivs gzippad. |
Nycklar för filurval:
| Nyckel | Standard | Beskrivning |
|---|---|---|
delete | False | Tar bort källfilerna efter lyckad kopiering. Det går inte att ångra. |
use_filename_timestamps | False | True matchar tidsstämpeln i filnamnet; False filtrerar på senast ändrad-tid. |
file_timestamp_format | [YYYY][MM][DD] | Tidsstämpelmönstret inuti filnamnen. Den finaste delen som finns med sätter också sökgranulariteten. |
start_timestamp | nu | ISO-datumtid — tidsfönstrets början. |
end_timestamp | nu | ISO-datumtid — tidsfönstrets slut. |
find_partitions | False | Endast Amazon S3. Listar underprefix under bassökvägen och söker i varje. Inte partitionsmedveten — för en year=/month=/day=-struktur, använd ett Hive-token i filePath i stället. |
Nycklar för läsaren:
| Nyckel | Standard | Beskrivning |
|---|---|---|
excel_sheet, excel_range | – | Bladnamn och cellområde för .xlsx-filer. |
csv_has_header | true | Första raden är rubrikrad. |
csv_sample_size | 10000 | Rader som samplas för typdetektering. |
json_sample_size | 10000 | Rader som samplas för schemahärledning. |
json_format | auto | auto, newline_delimited, array och DuckDB:s övriga JSON-format. |
json_union_by_name | true | Slår ihop scheman över filer i stället för att kräva identiska. |
file_encoding | utf-8 | Teckenkodning för CSV. |
Formatet detekteras i ordningen .xlsx → .parquet → JSON → CSV; filer som inte matchar något av dem faller tillbaka på binär kopiering. Alla matchande filer läses i en och samma omgång, så de måste dela schema om inte json_union_by_name täcker skillnaden.
Köexporter
Omvandlar en meddelandeström till batchfiler. En körning tömmer kön och avslutas sedan, så exporten ligger kvar på det vanliga schemat i stället för att köra en egen loop.
Själva kön konfigureras på anslutningen; exporten lägger bara till batchgränserna, via filterCondition:
| Nyckel | Standard | Beskrivning |
|---|---|---|
max_batch_bytes | 31457280 (30 MB) | Stäng batchfilen vid denna okomprimerade storlek. |
max_batch_seconds | 900 (15 min) | Stäng batchfilen efter denna tid i klocktimmar. |
max_batches_per_run | 0 (obegränsat) | Begränsar en het kö så att den inte svälter ut andra arbetslaster. |
Beteende värt att känna till:
- En tom kö skriver ingenting. Ingen nollbytefil skapas och ingen nedströmsbearbetning körs.
- CDC förbigås. Händelser konsumeras en gång och uppdateras aldrig, så sätt
enableCDC: 0för att slippa en meningslös baslinjefil. - Leverans är minst-en-gång. Kvittering sker efter att batchen skrivits till disk, så en krasch mitt i en körning leder till omleverans snarare än förlust.
- Felformade meddelanden skickas vidare, med en
__parseError-markering, i stället för att bli liggande som gift i kön. - Batchfiler namnges
<alias>_<suffix>_batch<NNNNNN>.txt.
Familjen är gjord för högvolymströmmar av små händelser. Landa stora payloads i objektlagring och läs in dem som en filexport, och låt kön bära bara pekaren.
Anpassade exporter
Minimal konfiguration för proprietära integrationer. Anslutningens egenskaper är den enda konfigurationskanalen, filterCondition substitueras och skickas till kopplingen, och outputFormat blir utdatafilens filändelse.
Schemaläggning
Alla exporttyper stöder schemaläggning:
| Inställning | Beskrivning | Exempel |
|---|---|---|
scheduleCycle | Frekvens | never, minute, hourly, daily, monthly. Intervall under 15 minuter rekommenderas inte. |
scheduleCycleInterval | Multiplikator | 2 med hourly = varannan timme |
scheduleDoNotStartBeforeTime | Tidigaste tillåtna start | 02:00:00 (starta inte före 02:00) |
never stänger av en export utan att radera den.
Exportstrategier
Att välja rätt exportstrategi är avgörande för att balansera datafärskhet mot kostnad och volym.
Fullständig export
Extraherar hela datasetet varje gång. Enkelt men dyrt för stora tabeller.
När den ska användas:
- Små referens-/dimensionstabeller
- Källsystem som inte stöder ändringsspårning
- Initiala laddningar eller periodiska fullständiga uppdateringar
Så konfigurerar du:
- Inga särskilda inställningar behövs — detta är standardbeteendet
- Låt
stopAtRow: -1stå kvar för ingen radgräns
Inkrementell export (vattenmärke)
Läser bara de rader källan har lagt till sedan förra körningen. Endast databaskällor.
När den ska användas:
- Stora transaktionstabeller där fullständiga exporter blir för dyra
- Tabeller med en tillförlitligt stigande nyckel eller sekvens
Så konfigurerar du:
- Sätt
filterColumntill en stigande numerisk kolumn — jämförelsen är numerisk, så en text- eller datumkolumn fungerar inte - Låt
lastFilterValuevara; agenten postar tillbaka det nya vattenmärket efter varje körning, med start från0
Så fungerar det:
- Den genererade frågan får ett villkor av typen
filterColumn > lastFilterValue - Kolumnens maxvärde följs per block och returneras som det nya vattenmärket
- En körning som kapats av
stopAtRowflyttar inte fram vattenmärket, så inget hoppas över nästa gång sqlOverrideförbigår detta helt — filtret läggs inte till i en anpassad fråga
Change Data Capture (CDC)
Jämför den här körningens extrakt mot föregående och märker varje post som ändrad, raderad eller oförändrad. Finns för alla källtyper som producerar rader.
När den ska användas:
- Källor utan egen ändringsindikator
- Alla källor där nedströms behöver veta om raderingar
Så konfigurerar du:
- Sätt
enableCDC: 1på exporten (0exporterar allt varje körning) - Markera fälten som ska följas med
includeInCDC: 1i sourcefile-strukturen - Använd ett
validFrom-fält för att identifiera vilken version av en post som är aktuell - Sätt
sortOrderpå versionsfältet (positivt = fallande, senaste först)
Så fungerar det:
- Föregående extrakt behålls som en baslinjefil med namnet
<alias>_latestversion.txtoch jämförs i DuckDB - Poster skrivs med en
__cdc-kolumn med värdetchange,deleteellerequal - Där en
__checksumfinns — databaskällor — ersätts den fullständiga jämförelsen av en indexerad anti-join, vilket är betydligt snabbare på breda tabeller - Om kolumnuppsättningen ändrats mellan körningar returneras hela den nya laddningen, märkt som ny
- Utdata delas vid ungefär 2 GB per fil; varje del laddas upp
Ingen baslinjefil betyder fullständig laddning, så den första körningen efter att CDC slagits på exporterar alltid allt. Varje körning väver dessutom in ett litet urval av den nya datan som en medveten omsådd, så CDC-utdata är aldrig ett rent delta.
stopAtRow stänger av CDC för den körningen och kastar den kapade utdatan i stället för att behålla den som baslinje — ett kapat extrakt jämfört mot en fullständig baslinje skulle läsa varje rad det inte skrev som en radering. Det är ett testreglage, inte en strypventil.
Binära filexporter och köexporter förbigår CDC helt.
Glidande fönster
En tidsbegränsad extraktion som flyttas framåt vid varje körning, eftersom fönstret uttrycks som tidstoken som renderas om varje körning i stället för som fasta datum.
När den ska användas:
- Filinläsning där filnamn eller partitionssökvägar bär tidsstämplar
- API:er med datumintervallparametrar
- Scenarier där du behöver de senaste N dygnen eller timmarna
Så konfigurerar du (filexporter):
- Sätt
use_filename_timestamps: Trueom tidsstämpeln finns i filnamnet, ochfile_timestamp_formatså att det matchar — annars matchas fönstret mot senast ändrad-tid - Uttryck fönstret med token i stället för datum:
filterCondition: ["start_timestamp=[__SCHEDULEDTIME_MINUS_1_DAYS__]", "end_timestamp=[__SCHEDULEDTIME__]"]
- För en partitionerad källa, adressera partitionen från
filePathmed[__SCHEDULEDTIME_AS_HIVE__]
Så konfigurerar du (API-exporter): Använd samma token i förfrågningsparametrarna:
filterCondition: ["date_from=[__LASTEXECUTION__]", "date_to=[__SCHEDULEDTIME__]"]
En körning läser ett fönster. [__SCHEDULEDTIME__] är den tidslucka körningen tillhör snarare än klocktiden, så ett omförsök läser samma fönster — men ett missat dygn innebär att den schemaläggningen körs om, inte att fönstret vidgas.
Kolumnval
Styr vilka kolumner som hamnar i den exporterade datan.
Inkluderings- och exkluderingslägen
Inkluderingsläge (standard):
Som standard gäller inclColumns: ["*"] — alla kolumner inkluderas. För att bara ta med vissa kolumner:
inclColumns: ["CustomerID", "Name", "Email", "CreatedDate"]
exclColumns: []
Exkluderingsläge: Ta med allt utom vissa kolumner:
inclColumns: ["*"]
exclColumns: ["InternalNotes", "TempFlag", "DebugData"]
Regler:
- Om
inclColumnsinnehåller"*"inkluderas alla fält exclColumnsfungerar som en svartlista ovanpå inkluderingslistan — nekande vinner- UI:t hindrar dig från att lämna noll inkluderade fält — det faller tillbaka på
["*"] - Kolumnnamn är skiftlägeskänsliga och måste matcha källan exakt
På databasexporter beter sig de två listorna olika. inclColumns fogas in i SELECT-satsen och är därmed det enda sättet att ta bort en kolumn; exclColumns stänger bara av typkonvertering och tar inte bort något ur utdatan.
Datatransformationer
Substitutioner
Substitutioner härleder eller ger standardvärde åt en kolumn medan exporten skrivs. Var och en definieras som { column, alias, replace, type }:
| Egenskap | Beskrivning |
|---|---|
column | Källkolumnen som substitutionen läser |
alias | Valfri målkolumn. Utan den skrivs källkolumnen över |
replace | Uttrycket eller literalen som används som ersättning |
type | transform — replace är ett uttryck som beräknas mot cellvärdet, och tillämpas bara när värdet inte är null. null — replace är standardvärdet som används när värdet är null |
Typiska användningar är att ersätta null-värden och att maska känsliga värden innan de lämnar källan.
replace beräknas som kod av agenten. Behandla exportdefinitioner som betrodd konfiguration, inte som användarindata.
Utdataformat och filnamn
Utdataformat (outputFormat):
| Format | Beskrivning |
|---|---|
json | Radbrytningsseparerad JSON. Standardvärdet, och det enda format de vanliga API-operatorerna skriver |
csv | Avgränsad text, med anslutningens ColumnDelimiter |
binary | Filer laddas upp orörda. Hoppar över konvertering och CDC; endast filkällor |
TABLE, ICEBERG och JSON är målformat som väljs på sourcefile i DLS — de är inte utdataformat för export.
Dynamiska filnamn (suffix):
Suffixet som läggs till utdatafilens namn är antingen en bokstavlig sträng eller ett tidstoken:
| Token | Värde |
|---|---|
[__SCHEDULEDTIME__] | Tidsluckan körningen tillhör — oförändrad av en sen start eller ett omförsök |
[__EXECUTIONTIME__] | Faktisk körningstid, i agentens tidszon |
[__LASTEXECUTION__] | Föregående lyckade körnings schemalagda tid (1900-01-01T00:00:00 vid första körningen) |
Samma token fungerar i förfrågningsparametrar, filsökvägar, where-satser och ändpunkterna i ett filurvalsfönster. De tar ett valfritt tidsavdrag och ett valfritt format:
[__<TIME>[_MINUS_<N>_<UNIT>][_AS_<FORMAT>]__]
<UNIT> är SECONDS, MINUTES, HOURS, DAYS eller WEEKS; det finns inget _PLUS_, eftersom en export läser ett fönster som redan inträffat. <FORMAT> är TIMESTAMP, DATETIME, DATE, DATECOMPACT, TIME, ISO, EPOCH, EPOCHMS eller HIVE, eller ett bokstavligt mönster byggt av [YYYY] [YY] [MM] [DD] [HH] [MI] [SS]. Ett värde utan _AS_-segment ärver exportens dateTimeFormat, som är TIMESTAMP som standard — eller DATETIME för databaskällor, vilket passar SQL-literaler.
HIVE renderar en partitionssökväg som year=2026/month=8/day=13 i stället för en datumtid, för att läsa källor som är upplagda så. Det hör hemma i en filsökväg och kan inte användas som exportövergripande format. Se utvecklarguiden för Ingest för varianterna med granularitet och nollutfyllnad.
Radbegränsning:
Sätt stopAtRow för att begränsa en exports storlek: -1 är ingen gräns, och valfritt positivt tal stoppar körningen där. Vad som räknas beror på källan:
| Källa | Enhet |
|---|---|
| API, anpassad | Rader |
| Databas | Block om 25 000 rader |
| Fil, konverterad | Skrivna rader — alla matchande filer laddas fortfarande ned |
| Fil, binär | Filer — det enda fallet där enheten inte är rader |
| Kö | Meddelanden, räknade över hela körningen |
Eftersom en kapad körning också stänger av CDC är stopAtRow till för att testa en ny export, inte för att begränsa en i produktion.
Tekniska kolumner
Varje export lägger till egna kolumner vid sidan av källdatan:
| Kolumn | Läggs till när | Värde |
|---|---|---|
__service | AddServiceColVal är satt på anslutningen | Det värdet |
__scheduledAt | Alltid | Körningens schemalagda tid, med exporttiden som reserv |
__exportedAt | Alltid | När exporten började skriva |
__checksum | Databaskällor | Hash av posten, som möjliggör den snabba CDC-vägen |
__cdc | CDC är påslaget | change, delete eller equal |
__parseError | Ett kömeddelande inte kunde tolkas | Tolkningsfelet, med den råa bodyn bevarad |
Destinationssökvägar
Exporterad data landar i konfigurerade zonsökvägar:
| Zon | Källa | Syfte |
|---|---|---|
| Landing Zone | Från Settings (landingZoneName) | Temporär staging |
| Raw Zone | Från Settings (rawZoneName) | Permanent arkiv |
| Trusted Zone | Från Settings (trustedZoneName) | Validerad data |
Standard sökvägsmönster: [system]/[filename]/[YYYY]/[MM]/[DD]/
Sourcefile Structure
När data har ingesterats beskrivs den av en sourcefile — schemadefinitionen för inkommande data.
Nyckelegenskaper
| Egenskap | Beskrivning |
|---|---|
sourceFilename | Unik identifierare |
system | Källsystemreferens |
fileType | CSV, JSON eller XML |
fileEncoding | UTF-8, UTF-16, Windows-1252 |
columnDelimiter | Avgränsare för CSV-filer |
targetMethod | TRANSACTION, APPEND, CHANGES ONLY, LATEST VERSION eller OVERWRITE — se DLS-fältreferensen |
targetFormat | TABLE, ICEBERG eller JSON |
targetNormalization | NONE, LISTS eller LISTS AND OBJECTS |
fileCompressionType | Komprimeringstyp: gzip eller ingen |
enableEncryption | Krypteringsflagga |
Hierarkisk filstruktur
För nästlad data (t.ex. JSON med arrayer) stöder sourcefile en hierarkisk struktur:
Level 0 (OBJECT): Root
├── customer.id (keyFieldIndicator=1)
├── customer.name
└── customer.loadDate (validFrom=1)
Level 1 (LIST): orders (useForSplittingRecords=1)
├── order.orderId
├── order.amount
└── order.status
Varje nivå kan vara ett OBJECT (enskild post) eller LIST (array som delas upp i rader).
Fältnivåegenskaper
| Egenskap | Beskrivning |
|---|---|
fieldKey | Källkolumnnamn (skiftlägeskänsligt) |
fieldAlias | Valfritt omdöpt kolumnnamn |
dataType | Varchar, Integer, Decimal, Boolean, Time, Date, Timestamp |
keyFieldIndicator | 1 = del av primärnyckel |
keyOrder | Position i sammansatt nyckel |
fieldOrder | Kolumnordning i mål |
excludeField | 1 = hoppa över detta fält helt |
excludeFromProfiling | 1 = kör inte kvalitetskontroller |
sensitive | 1 = PII/känslig data-flagga |
includeInCDC | 1 = spåra ändringar för detta fält |
validFrom | 1 = SCD Typ 2-versioneringsfält (max ett per nivå) |
fieldDomain | Styrningsklassificeringskategori |
Bästa praxis
Välja exportstrategi
| Scenario | Rekommenderad strategi |
|---|---|
| Liten dimensionstabell (<100K rader) | Fullständig export med CDC, dagligen |
| Stor transaktionstabell med stigande nyckel | Inkrementellt vattenmärke på filterColumn |
| Stor tabell utan ändringsindikator | Fullständig export med CDC — checksummevägen håller jämförelsen billig |
| Fildrop med tidsstämplade filnamn | Glidande fönster på start_timestamp / end_timestamp |
Partitionerad sjö (year=/month=/day=) | Hive-token i filePath, en partition per körning |
| API med paginering + datumfilter | Glidande fönster via förfrågningsparametrar |
| Meddelandekö | Batchgränser på anslutningen, CDC avstängt |
| Initial dataladdning | Fullständig export, byt sedan till inkrementell |
Tips för kolumnval
- Börja med
inclColumns: ["*"]och användexclColumnsför att ta bort oönskade kolumner - På databasexporter tar du i stället bort kolumner ur
inclColumns—exclColumnstar inte bort dem där - Exkludera stora text-/blob-kolumner som inte behövs nedströms
- Exkludera interna/felsökningskolumner som skapar brus
- Kom ihåg: kolumnnamn är skiftlägeskänsliga
Schemaläggningsöverväganden
- Använd
scheduleDoNotStartBeforeTimeför att undvika körning under högtrafiktimmar - För beroende exporter, förskjut scheman (t.ex. dimensionstabeller före faktatabeller)
- Använd
monthly-cykel för långsamt föränderlig referensdata - Använd
hourlymed litet intervall för nästan-realtidsbehov
Nästa steg
- DLS-lagret — vad som händer med en leverans när den landat
- Arbeta med Ingest — konfigurationsguiden fält för fält
- Arkitektur — hur lagren hänger ihop