Skip to main content

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.

AnslutningstypKällaOperatorer
Database (DB)SQL-databaser som nås med SQLPostgreSQL, SQL Server, Oracle, MySQL
APIHTTP-ändpunkter som returnerar JSONGenerisk REST, Heartpace, Salesforce Service Cloud, Talkdesk, SharePoint (Graph), Microsoft Teams, Google Drive / Sheets
FileFillagringssystemLokal disk, Amazon S3, Azure Data Lake Gen2, SFTP
CustomMeddelandeköer och skräddarsydda källorAWS SQS, Azure Service Bus, Redis, samt kundspecifika kopplingar
ManualIngen automatiserad extraktionData uppladdad manuellt
note

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ällningEffekt
AddServiceColValNär den är satt får varje exporterad post en __service-kolumn med detta värde. Tomt värde stänger av kolumnen.
ColumnDelimiterFältavgränsare när CDC skriver CSV-utdata. Ignoreras för JSON.
LineDelimiterPostavgrä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ällningBeskrivning
databaseNameReferens till den anslutna databasen. Skrivskyddad, hämtas från anslutningen.
fromClauseTabell, vy, funktion eller lagrad procedur att fråga. Måste vara schemakvalificerad (dbo.Customers) — metadatauppslaget delar på punkten.
whereClauseValfritt SQL-villkor. Infogas ordagrant, så det måste innehålla nyckelordet WHERE. Ett inkrementellt filter läggs till med AND.
sqlOverrideFullstä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.
inclTablesYtterligare 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.
inclColumnsSammanfogas med kommatecken rakt in i SELECT-satsen, så alias och uttryck fungerar. ["*"] väljer allt.
filterColumnAktiverar inkrementell laddning. Jämförelsen är numerisk, så kolumnen måste vara ett stigande tal.
lastFilterValueVattenmä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ällningBeskrivning
apiUrlBas-URL för API:et. Skrivskyddad, hämtas från anslutningen.
apiPathÄndpunktens sökväg, som läggs till bas-URL:en.
filterConditionNyckel=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:

OperatorPagineringsmodell
Heartpacelimit/offset, styrt av svarets meta-block
Salesforce Service CloudFöljer nextRecordsUrl tills resultatmängden rapporterar sig färdig
TalkdeskAsynkrona 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 DriveEn 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ällningBeskrivning
filePathKä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__].
fileNameEn delsträngsmatchning mot filnamnet — inte ett glob-mönster.
filterConditionFilurval och läsaralternativ (nedan).
outputFormatjson eller csv konverterar filerna; binary laddar upp dem orörda och hoppar över CDC. Konverterad utdata skrivs gzippad.

Nycklar för filurval:

NyckelStandardBeskrivning
deleteFalseTar bort källfilerna efter lyckad kopiering. Det går inte att ångra.
use_filename_timestampsFalseTrue 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_timestampnuISO-datumtid — tidsfönstrets början.
end_timestampnuISO-datumtid — tidsfönstrets slut.
find_partitionsFalseEndast 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:

NyckelStandardBeskrivning
excel_sheet, excel_range–Bladnamn och cellområde för .xlsx-filer.
csv_has_headertrueFörsta raden är rubrikrad.
csv_sample_size10000Rader som samplas för typdetektering.
json_sample_size10000Rader som samplas för schemahärledning.
json_formatautoauto, newline_delimited, array och DuckDB:s övriga JSON-format.
json_union_by_nametrueSlår ihop scheman över filer i stället för att kräva identiska.
file_encodingutf-8Teckenkodning 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:

NyckelStandardBeskrivning
max_batch_bytes31457280 (30 MB)Stäng batchfilen vid denna okomprimerade storlek.
max_batch_seconds900 (15 min)Stäng batchfilen efter denna tid i klocktimmar.
max_batches_per_run0 (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: 0 fö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ällningBeskrivningExempel
scheduleCycleFrekvensnever, minute, hourly, daily, monthly. Intervall under 15 minuter rekommenderas inte.
scheduleCycleIntervalMultiplikator2 med hourly = varannan timme
scheduleDoNotStartBeforeTimeTidigaste tillåtna start02: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: -1 stå 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:

  1. Sätt filterColumn till en stigande numerisk kolumn — jämförelsen är numerisk, så en text- eller datumkolumn fungerar inte
  2. Låt lastFilterValue vara; agenten postar tillbaka det nya vattenmärket efter varje körning, med start från 0

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 stopAtRow flyttar inte fram vattenmärket, så inget hoppas över nästa gång
  • sqlOverride fö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:

  1. Sätt enableCDC: 1 på exporten (0 exporterar allt varje körning)
  2. Markera fälten som ska följas med includeInCDC: 1 i sourcefile-strukturen
  3. Använd ett validFrom-fält för att identifiera vilken version av en post som är aktuell
  4. Sätt sortOrder på versionsfältet (positivt = fallande, senaste först)

Så fungerar det:

  • Föregående extrakt behålls som en baslinjefil med namnet <alias>_latestversion.txt och jämförs i DuckDB
  • Poster skrivs med en __cdc-kolumn med värdet change, delete eller equal
  • Där en __checksum finns — 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
note

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.

warning

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):

  1. Sätt use_filename_timestamps: True om tidsstämpeln finns i filnamnet, och file_timestamp_format så att det matchar — annars matchas fönstret mot senast ändrad-tid
  2. Uttryck fönstret med token i stället för datum:
filterCondition: ["start_timestamp=[__SCHEDULEDTIME_MINUS_1_DAYS__]", "end_timestamp=[__SCHEDULEDTIME__]"]
  1. För en partitionerad källa, adressera partitionen från filePath med [__SCHEDULEDTIME_AS_HIVE__]

Så konfigurerar du (API-exporter): Använd samma token i förfrågningsparametrarna:

filterCondition: ["date_from=[__LASTEXECUTION__]", "date_to=[__SCHEDULEDTIME__]"]
note

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 inclColumns innehåller "*" inkluderas alla fält
  • exclColumns fungerar 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
warning

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 }:

EgenskapBeskrivning
columnKällkolumnen som substitutionen läser
aliasValfri målkolumn. Utan den skrivs källkolumnen över
replaceUttrycket eller literalen som används som ersättning
typetransform — 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.

warning

replace beräknas som kod av agenten. Behandla exportdefinitioner som betrodd konfiguration, inte som användarindata.

Utdataformat och filnamn

Utdataformat (outputFormat):

FormatBeskrivning
jsonRadbrytningsseparerad JSON. Standardvärdet, och det enda format de vanliga API-operatorerna skriver
csvAvgränsad text, med anslutningens ColumnDelimiter
binaryFiler 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:

TokenVä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ällaEnhet
API, anpassadRader
DatabasBlock om 25 000 rader
Fil, konverteradSkrivna rader — alla matchande filer laddas fortfarande ned
Fil, binärFiler — 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:

KolumnLäggs till närVärde
__serviceAddServiceColVal är satt på anslutningenDet värdet
__scheduledAtAlltidKörningens schemalagda tid, med exporttiden som reserv
__exportedAtAlltidNär exporten började skriva
__checksumDatabaskällorHash av posten, som möjliggör den snabba CDC-vägen
__cdcCDC är påslagetchange, delete eller equal
__parseErrorEtt kömeddelande inte kunde tolkasTolkningsfelet, med den råa bodyn bevarad

Destinationssökvägar​

Exporterad data landar i konfigurerade zonsökvägar:

ZonKällaSyfte
Landing ZoneFrån Settings (landingZoneName)Temporär staging
Raw ZoneFrån Settings (rawZoneName)Permanent arkiv
Trusted ZoneFrå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
EgenskapBeskrivning
sourceFilenameUnik identifierare
systemKällsystemreferens
fileTypeCSV, JSON eller XML
fileEncodingUTF-8, UTF-16, Windows-1252
columnDelimiterAvgränsare för CSV-filer
targetMethodTRANSACTION, APPEND, CHANGES ONLY, LATEST VERSION eller OVERWRITE — se DLS-fältreferensen
targetFormatTABLE, ICEBERG eller JSON
targetNormalizationNONE, LISTS eller LISTS AND OBJECTS
fileCompressionTypeKomprimeringstyp: gzip eller ingen
enableEncryptionKrypteringsflagga
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
EgenskapBeskrivning
fieldKeyKällkolumnnamn (skiftlägeskänsligt)
fieldAliasValfritt omdöpt kolumnnamn
dataTypeVarchar, Integer, Decimal, Boolean, Time, Date, Timestamp
keyFieldIndicator1 = del av primärnyckel
keyOrderPosition i sammansatt nyckel
fieldOrderKolumnordning i mål
excludeField1 = hoppa över detta fält helt
excludeFromProfiling1 = kör inte kvalitetskontroller
sensitive1 = PII/känslig data-flagga
includeInCDC1 = spåra ändringar för detta fält
validFrom1 = SCD Typ 2-versioneringsfält (max ett per nivå)
fieldDomainStyrningsklassificeringskategori

Bästa praxis​

Välja exportstrategi
ScenarioRekommenderad strategi
Liten dimensionstabell (<100K rader)Fullständig export med CDC, dagligen
Stor transaktionstabell med stigande nyckelInkrementellt vattenmärke på filterColumn
Stor tabell utan ändringsindikatorFullständig export med CDC — checksummevägen håller jämförelsen billig
Fildrop med tidsstämplade filnamnGlidande fönster på start_timestamp / end_timestamp
Partitionerad sjö (year=/month=/day=)Hive-token i filePath, en partition per körning
API med paginering + datumfilterGlidande fönster via förfrågningsparametrar
MeddelandeköBatchgränser på anslutningen, CDC avstängt
Initial dataladdningFullständig export, byt sedan till inkrementell
Tips för kolumnval
  • Börja med inclColumns: ["*"] och använd exclColumns för att ta bort oönskade kolumner
  • På databasexporter tar du i stället bort kolumner ur inclColumns — exclColumns tar 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 scheduleDoNotStartBeforeTime fö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 hourly med litet intervall för nästan-realtidsbehov

Nästa steg​