Använda fristående direktuppspelningstabeller

En fristående strömningstabell är en tabell som är registrerad i Unity Catalog med extra stöd för direktuppspelning eller inkrementell databearbetning, som definieras utanför en Lakeflow-pipeline. En pipeline skapas automatiskt för varje strömmande tabell. Du kan använda strömmande tabeller för inkrementell datainläsning från Kafka och molnobjektlagring.

Du kan skapa och uppdatera fristående strömmande tabeller antingen från ett Databricks SQL Warehouse eller från en anteckningsbok som körs på serverlös generell beräkning. Mer information om skillnaderna mellan de två beräkningsalternativen finns i Krav för fristående pipelines.

För att skapa och uppdatera fristående strömningstabeller med Python från en notebook, se Använd Python med fristående pipelines.

Anmärkning

Mer information om hur du använder Delta Lake-tabeller som strömmande källor och utdata finns i strömmande läsningar och skrivningar för Delta Lake-tabeller.

Kravspecifikation

Mer information om beräkningsalternativ, behörigheter och andra krav för att skapa, uppdatera och fråga fristående strömningstabeller finns i Krav för fristående pipelines.

Skapa strömmande tabeller

En strömmande tabell definieras av en SQL-fråga i Databricks SQL. När du skapar en strömmande tabell används de data som för närvarande finns i källtabellerna för att skapa strömningstabellen. Därefter uppdaterar du tabellen, vanligtvis enligt ett schema, för att hämta eventuella tillagda data i källtabellerna för att lägga till i strömningstabellen.

När du skapar en direktuppspelningstabell betraktas du som ägare till tabellen.

Om du vill skapa en strömmande tabell från en befintlig tabell använder du -instruktionenCREATE STREAMING TABLE, som i följande exempel:

CREATE OR REFRESH STREAMING TABLE sales
  SCHEDULE EVERY 1 hour
  AS SELECT product, price FROM STREAM raw_data;

I det här fallet skapas strömningstabellen sales från specifika kolumner i raw_data tabellen, med ett schema för uppdatering varje timme. Frågan som används måste vara en direktuppspelningsfråga . Använd nyckelordet STREAM för att använda strömmande semantik för att läsa från källan.

Beräkning som används för uppdatering

När du skapar en strömmande tabell med instruktionen CREATE OR REFRESH STREAMING TABLE börjar den första datauppdateringen och populationen omedelbart. Dessa åtgärder använder inte Databricks SQL Warehouse-beräkning. I stället förlitar sig strömmande tabeller på serverlösa pipelines för både skapande och uppdatering. En dedikerad serverlös pipeline skapas och hanteras automatiskt av systemet för varje strömmande tabell.

Läsa in filer med Auto Loader

Om du vill skapa en strömmande tabell från filer i en volym använder du Auto Loader. Använd Auto Loader för de flesta datainmatningsuppgifter från molnobjektlagring. Automatisk inläsning och pipelines är utformade för att inkrementellt och idempotent läsa in ständigt växande data när de kommer till molnlagringen.

Om du vill använda Auto Loader i Databricks SQL använder du read_files funktionen. Följande exempel visar hur du använder Auto Loader för att läsa en volym JSON-filer till en strömmande tabell:

CREATE OR REFRESH STREAMING TABLE sales
  SCHEDULE EVERY 1 hour
  AS SELECT * FROM STREAM read_files(
    "/Volumes/my_catalog/my_schema/my_volume/path/to/data",
    format => "json"
  );

Om du vill läsa data från molnlagring kan du också använda Automatisk inläsning:

CREATE OR REFRESH STREAMING TABLE sales
  SCHEDULE EVERY 1 hour
  AS SELECT *
  FROM STREAM read_files(
    'abfss://myContainer@myStorageAccount.dfs.core.windows.net/analysis/*/*/*.json',
    format => "json"
  );

Mer information om automatisk inläsning finns i Vad är automatisk inläsning?. Mer information om hur du använder automatisk inläsning i SQL finns i Läsa in data från objektlagring.

Strömmande inmatning från andra källor

Exempel på inmatning från andra källor, inklusive Kafka, se Läsa in data i pipelines.

Tillämpa ändringsdatainsamling (CDC) med automatiska CDC-flöden

Använd FLOW AUTO CDC-satsen för att bearbeta Change Data Capture (CDC)-register från en källa till en strömmande tabell. Tidigare användes instruktionen MERGE INTO ofta för bearbetning av CDC-poster på Azure Databricks. Kan dock MERGE INTO ge felaktiga resultat på grund av poster som inte är sekvenserade eller kräver komplex logik för att ordna om poster. Se Ändra datainsamling och ögonblicksbilder.

AUTO CDC förenklar CDC genom att automatiskt hantera out-of-order-poster. Du anger nycklar för att identifiera poster, en sekvenskolumn för beställning och om resultat ska lagras som SCD typ 1 (direktuppdateringar) eller SCD typ 2 (historikspårning).

I följande exempel skapas en strömmande tabell som tillämpar CDC-ändringar med SCD typ 1.

CREATE OR REFRESH STREAMING TABLE target
  FLOW AUTO CDC
  FROM stream(cdc_data.users)
  KEYS (userId)
  SEQUENCE BY sequenceNum
  STORED AS SCD TYPE 1;

I följande exempel används SCD typ 2 för att behålla en historik över ändringar:

CREATE OR REFRESH STREAMING TABLE target
  FLOW AUTO CDC
  FROM stream(cdc_data.users)
  KEYS (userId)
  APPLY AS DELETE WHEN operation = "DELETE"
  SEQUENCE BY sequenceNum
  COLUMNS * EXCEPT (operation, sequenceNum)
  STORED AS SCD TYPE 2;

Fullständig information om alternativ och beteende för Automatisk CDC finns i API:er för AUTOMATISK CDC: Förenkla insamling av ändringsdata med pipelines. Den fullständiga syntaxreferensen finns i CREATE STREAMING TABLE.

Tillämpa selektiv batchersättning med REPLACE-flöden WHERE

FLOW REPLACE WHERE Använd -satsen för att omkomputera och skriva över en målunderuppsättning av en strömmande tabell utan att bearbeta hela tabellhistoriken på nytt. REPLACE WHERE flöden är väl lämpade för inkrementell batchbearbetning av sammanfogningar och aggregeringar, sent inkommande data, ombearbetning uppströms, schemaändringar och bakåtfyllning.

Fullständig information om REPLACE WHERE flöden, inklusive krav, åsidosättningar av predikat och inkrementell uppdatering, finns i REPLACE-flöden WHERE för fristående strömningstabeller.

Applicera delvis snapshot-ersättning med REPLACE USING flows

Viktigt!

Ersätt med flödena är i betaversion.

Använd klausulen FLOW REPLACE USING för att hålla en strömningstabell synkroniserad med en ström av partiella snapshots. Vid varje uppdatering ersätter ett flöde REPLACE USING alla rader som matchar de angivna nyckelkolumnerna och lämnar alla andra rader oförändrade. En SEQUENCE BY kolumn ordnar uppdateringarna så att den högsta sekvensen för en nyckel alltid vinner, även när uppdateringar kommer i fel ordning. Till exempel:

CREATE OR REFRESH STREAMING TABLE payments_current
FLOW REPLACE USING (payment_id) SEQUENCE BY payment_date BY NAME
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

BY NAME måste anges. Den matchar kolumner efter namn i stället för position.

REPLACE USING fungerar på samma sätt för fristående strömningstabeller som i Lakeflow-pipelines. Information om hur det fungerar, ordningsföljd, förväntningar, begränsningar och exempel finns i Partiell ögonblicksbildsersättning med REPLACE USING-flöden. Följande skillnader gäller för fristående strömningstabeller:

  • Definiera flödet i SQL. Skapa flödet REPLACE USING med den inbyggda SQL-klausulen FLOW REPLACE USINGCREATE OR REFRESH STREAMING TABLE. Den fristående CREATE FLOW satsen är en Lakeflow-pipelinekonstruktion och används inte för fristående strömningstabeller.
  • Datorkapaciteten hanteras åt dig. Fristående strömningstabeller körs på serverlösa pipelines som hanteras av systemet och kräver Databricks Runtime 18.2 eller senare. Du väljer inte mellan klassisk och serverless beräkning.

Mata bara in ny data

Som standard read_files läser funktionen alla befintliga data i källmappen när tabellen skapas och bearbetar sedan nyligen ankommande poster med varje uppdatering.

Om du vill undvika att mata in data som redan finns i källmappen när tabellen skapas anger du includeExistingFiles alternativet till false. Det innebär att endast data som tas emot i mappen när tabellen har skapats har bearbetats. Till exempel:

CREATE OR REFRESH STREAMING TABLE sales
  SCHEDULE EVERY 1 hour
  AS SELECT *
  FROM STREAM read_files(
    '/path/to/files',
    includeExistingFiles => false
  );

Runtime-versionen

Strömningstabeller körs alltid på den senaste versionen av Databricks SQL-runtime. Tabellegenskapen pipelines.channel , som tidigare användes för att välja en preview eller current runtime-kanal, stöds inte längre och har ingen effekt. Om en befintlig definition inkluderar denna egenskap ignoreras den säkert och du behöver inte ta bort den.

Dölj känsliga data

Du kan använda strömmande tabeller för att dölja känsliga data från användare som kommer åt tabellen. En metod är att definiera frågan så att den utesluter känsliga kolumner eller rader helt och hållet. Du kan också använda kolumnmasker eller radfilter baserat på behörigheterna för den frågande användaren. Du kan till exempel dölja tax_id kolumnen för användare som inte finns i gruppen HumanResourcesDept. Det gör du genom att använda syntaxen ROW FILTER och MASK när du skapar strömningstabellen. Mer information finns i Radfilter och kolumnmasker.

Uppdatera en strömningstabell

Strömmande tabeller skapar och använder automatiskt serverlösa pipelines för att bearbeta uppdateringsåtgärder. Förnyelsen hanteras av pipelinen och uppdateringen övervakas av Databricks SQL-lagret som används för att skapa den strömmande tabellen. Strömmande tabeller kan uppdateras med hjälp av en pipeline som körs enligt ett schema.

Även om du har en schemalagd uppdatering kan du anropa en manuell uppdatering när som helst. Uppdateringar hanteras av samma pipeline som skapades automatiskt tillsammans med strömningstabellen.

Så här uppdaterar du en strömmande tabell:

REFRESH STREAMING TABLE sales;

Du kan kontrollera status för den senaste uppdateringen med DESCRIBE TABLE EXTENDED.

Anmärkning

Du kan behöva uppdatera strömningstabellen innan du använder frågor om tidsresor.

Information om hur du schemalägger en uppdatering finns i Schemalägg uppdateringar. Schemalagda uppdateringar kan ha uppdateringsmeddelanden och du kan ange prestandaläget för uppdateringen.

Så här fungerar uppdateringen

En uppdatering av strömningstabellen utvärderar endast nya rader som har anlänt efter den senaste uppdateringen och lägger endast till nya data.

Varje uppdatering använder den aktuella definitionen av strömningstabellen för att bearbeta dessa nya data. Om du ändrar en definition för en strömmande tabell beräknas inte befintliga data automatiskt om. Om en ändring inte är kompatibel med befintliga data (till exempel om du ändrar en datatyp) misslyckas nästa uppdatering med ett fel.

I följande exempel förklaras hur ändringar i en strömmande tabelldefinition påverkar uppdateringsbeteendet:

  • Om du tar bort ett filter bearbetas inte tidigare filtrerade rader.
  • Att ändra kolumnprojektioner påverkar inte hur befintliga data bearbetades.
  • Kopplingar med statiska ögonblicksbilder använder ögonblicksbildens tillstånd vid det första bearbetningstillfället. Data som anländer sent och som skulle ha stämt överens med den uppdaterade ögonblicksbilden ignoreras. Detta kan leda till att fakta tas bort om dimensionerna är sena.
  • Om du ändrar CAST för en befintlig kolumn resulterar det i ett fel.

Om dina data ändras på ett sätt som inte kan stödjas i den befintliga strömningstabellen kan du utföra en fullständig uppdatering.

Uppdatera en strömningstabell fullständigt

Fullständiga uppdateringar bearbetar om alla data som är tillgängliga i källan med den senaste definitionen. Vi rekommenderar inte att du anropar fullständiga uppdateringar på källor som inte behåller hela datahistoriken eller har korta kvarhållningsperioder, till exempel Kafka, eftersom den fullständiga uppdateringen trunkerar befintliga data. Du kanske inte kan återställa gamla data om data inte längre är tillgängliga i källan.

Till exempel:

REFRESH STREAMING TABLE sales FULL;

Schemalägga och övervaka uppdateringar

Du kan uppdatera en strömmande tabell automatiskt enligt ett schema eller när överordnade data ändras, och du kan konfigurera tidsgränser för uppdateringar, meddelanden och prestandalägen. Se Schemalagda uppdateringar.

Kontrollera åtkomsten till strömmande tabeller

Strömningstabeller har stöd för omfattande åtkomstkontroller för att stödja datadelning och samtidigt undvika att exponera potentiellt privata data. En ägare av direktuppspelningstabellen eller en användare med behörigheten MANAGE kan bevilja SELECT behörigheter till andra användare. Användare med SELECT åtkomst till strömningstabellen behöver SELECT inte åtkomst till tabellerna som refereras av strömningstabellen. Den här åtkomstkontrollen möjliggör datadelning samtidigt som åtkomsten till underliggande data kontrolleras.

Du kan också ändra ägaren till en strömmande tabell.

Tilldela behörigheter till en strömningstabell

Om du vill bevilja åtkomst till en strömningstabell använder du -instruktionenGRANT:

GRANT <privilege_type> ON <st_name> TO <principal>;

privilege_type Kan vara:

  • SELECT – användaren kan SELECT strömma tabellen.
  • REFRESH – användaren kan REFRESH strömma tabellen. Uppdateringar körs med ägarens behörigheter.

I följande exempel skapas en strömningstabell och användare får behörighet att välja och uppdatera:

CREATE OR REFRESH STREAMING TABLE st_name AS SELECT * FROM source_table;

-- Grant read-only access:
GRANT SELECT ON st_name TO read_only_user;

-- Grant read and refresh access:
GRANT SELECT ON st_name TO refresh_user;
GRANT REFRESH ON st_name TO refresh_user;

Mer information om hur du beviljar behörigheter för skyddsbara objekt i Unity Catalog finns i Referens för Behörigheter för Unity Catalog.

Återkalla behörigheter från en streaming-tabell

Om du vill återkalla åtkomst från en strömningstabell använder du -instruktionenREVOKE:

REVOKE privilege_type ON <st_name> FROM principal;

När SELECT behörigheter i en källtabell återkallas från strömningstabellägaren eller någon annan användare som har beviljats MANAGE eller SELECT behörigheter i strömningstabellen, eller om källtabellen tas bort, kan ägaren eller användaren som beviljats åtkomst till strömningstabellen fortfarande köra frågor mot strömningstabellen. Följande beteende inträffar dock:

  • Den strömmande tabellägaren eller andra som har förlorat åtkomsten till en strömmande tabell kan inte längre REFRESH den strömmande tabellen, och strömningstabellen blir föråldrad med tiden.
  • Om det automatiseras med ett schema, misslyckas nästa planerade REFRESH eller körs inte.

I följande exempel återkallas behörigheten SELECT från read_only_user:

REVOKE SELECT ON st_name FROM read_only_user;

Ändra ägaren till en streamningstabell

En användare med MANAGE behörigheter för en fristående direktuppspelningstabell kan ange en ny ägare via Katalogutforskaren. Den nya ägaren kan vara sig själv eller ett tjänstehuvudnamn där de har Service Principal User-rollen.

  1. Från din Azure Databricks-arbetsyta klickar du på dataikonen.Katalog för att öppna Katalogutforskaren.

  2. Välj den strömmande tabell som du vill uppdatera.

  3. I det högra sidofältet, under Om den här strömningstabellen, letar du upp Ägare och klickar på Pennikonen för att redigera.

    Anmärkning

    Om du får ett meddelande som säger att du ska uppdatera ägaren genom att ändra användaren i inställningen Kör som i pipelineinställningarna, är streamingtabellen definierad i en Lakeflow-pipeline och inte som en fristående tabell. Meddelandet innehåller en länk till pipelineinställningarna: där du kan ändra Kör som användare.

  4. Välj en ny ägare för strömningstabellen.

    Ägare har automatiskt behörigheterna MANAGE och SELECT på strömmande tabeller som de äger. Om du anger ett huvudnamn för tjänsten som ägare för en strömmande tabell som du äger, och du inte uttryckligen har SELECT eller MANAGE behörighet på strömningstabellen, skulle den här ändringen leda till att du förlorar all åtkomst till strömningstabellen. I det här fallet uppmanas du att uttryckligen ange dessa privilegier.

    Välj både Bevilja HANTERA och Bevilja SELECT behörigheter för att ange dem i Spara.

  5. Klicka på Spara för att ändra ägaren.

Ägaren till strömningstabellen uppdateras. Alla framtida uppdateringar körs med den nya ägarens identitet.

När ägaren förlorar behörighet till källtabeller

Om du ändrar ägaren och den nya ägaren inte har åtkomst till källtabellerna (eller om behörigheterna på de underliggande källtabellerna återkallas) kan användarna fortfarande köra frågor mot strömningstabellen. Observera följande:

  • De kan inte REFRESH strömma tabellen.
  • Nästa schemalagda uppdatering av strömningstabellen misslyckas.

Att förlora åtkomsten till källdata förhindrar uppdateringar, men förhindrar inte omedelbart den befintliga strömmande tabellen från att läsas.

ta bort poster från en streamingtabell permanent

Viktigt!

Stöd för REORG-instruktionen med strömmande tabeller finns i offentlig förhandsversion.

Anmärkning

  • Om du använder en REORG-instruktion med en strömmande tabell krävs Databricks Runtime 15.4 och senare.
  • Även om du kan använda -instruktionen REORG med en strömningstabell krävs den bara när du tar bort poster från en strömmande tabell med borttagningsvektorer aktiverade. Kommandot har ingen effekt när det används med en strömmande tabell utan att borttagningsvektorer har aktiverats.

För att fysiskt ta bort poster från den underliggande lagringen för en strömmande tabell med borttagningsvektorer aktiverade, till exempel för GDPR-efterlevnad, måste ytterligare åtgärder vidtas för att säkerställa att en VACUUM åtgärd körs på strömningstabellens data.

Så här tar du bort poster fysiskt från underliggande lagring:

  1. Uppdatera poster eller ta bort poster från strömningstabellen.
  2. Kör en REORG-instruktion mot strömningstabellen och ange parametern APPLY (PURGE). Till exempel REORG TABLE <streaming-table-name> APPLY (PURGE);.
  3. Vänta tills datakvarhållningsperioden för strömningstabellen har passerat. Standardperioden för datakvarhållning är sju dagar, men den kan konfigureras med egenskapen delta.deletedFileRetentionDuration tabell. Se Konfigurera datalagring för frågor rörande tidsresor.
  4. REFRESH strömningstabellen. Se även Uppdatera en streamingtabel. Inom 24 timmar efter REFRESH-åtgärden körs pipelineunderhållsuppgifter automatiskt, inklusive VACUUM-åtgärden som krävs för att säkerställa att poster tas bort permanent.

Övervaka körningar med hjälp av frågehistorik

Du kan använda sidan för frågehistorik för att komma åt frågeinformation och frågeprofiler som kan hjälpa dig att identifiera frågor och flaskhalsar som fungerar dåligt i pipelinen som används för att köra uppdateringar av strömningstabellen. En översikt över vilken typ av information som är tillgänglig i frågehistorik och frågeprofiler finns i Frågehistorik och Frågeprofil.

Viktigt!

Den här funktionen finns som allmänt tillgänglig förhandsversion. Arbetsyteadministratörer kan styra åtkomsten till den här funktionen från sidan Förhandsversioner . Se Hantera förhandsversioner av Azure Databricks.

Alla instruktioner som rör strömmande tabeller visas i frågehistoriken. Du kan använda listrutan Statement för att välja valfritt kommando och granska relaterade frågor. Alla CREATE instruktioner följs av en REFRESH instruktion som körs asynkront på en pipeline. Instruktionerna REFRESH innehåller vanligtvis detaljerade frågeplaner som ger insikter om hur du optimerar prestanda.

Använd följande steg för att komma åt REFRESH instruktioner i användargränssnittet för frågehistorik:

  1. Klicka på Ikonen Historik. Öppna användargränssnittet för frågehistorik i det vänstra sidofältet.
  2. Markera kryssrutan REFRESH från filterlistrutan Statement.
  3. Klicka på namnet på frågeuttrycket för att visa sammanfattningsinformation som frågans varaktighet och aggregerade mått.
  4. Klicka på Se frågeprofil för att öppna frågeprofilen. Mer information om hur du navigerar i frågeprofilen finns i Frågeprofil .
  5. Du kan också använda länkarna i avsnittet Frågekälla för att öppna den relaterade frågan eller pipelinen.

Du kan också komma åt frågeinformation med hjälp av länkar i SQL-redigeraren eller från en notebook-fil som är kopplad till ett SQL-lager.

Få åtkomst till strömmande tabeller från externa klienter

Om du vill komma åt strömmande tabeller från externa Delta Lake- eller Iceberg-klienter som inte stöder öppna API:er kan du använda kompatibilitetsläge. Kompatibilitetsläget skapar en skrivskyddad version av din strömmande tabell som kan nås av alla Delta Lake- eller Iceberg-klienter.

Ytterligare resurser