Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
A streamelő táblák olyan táblák, amely támogatja a streamelést vagy a növekményes adatfeldolgozást. A streamelési táblákat csővezetékek támasztják alá. A streamelési táblák minden frissítésekor a forrástáblákhoz hozzáadott adatok hozzá lesznek fűzve a streamelési táblához. A streamelési táblákat manuálisan vagy ütemezés szerint frissítheti.
Ha többet szeretne megtudni a frissítések végrehajtásáról vagy ütemezéséről, olvassa el a Folyamatfrissítés futtatása című témakört.
Szemantika
CREATE [OR REFRESH] [PRIVATE] STREAMING TABLE
table_name
[ table_specification ]
[ table_clauses ]
[ {flow_clause | AS query} ]
table_specification
( { column_identifier column_type [column_properties] } [, ...]
[ column_constraint ] [, ...]
[ , table_constraint ] [...] )
column_properties
{ NOT NULL | GENERATED ALWAYS AS ( expr ) | GENERATED { ALWAYS | BY DEFAULT } AS IDENTITY [ ( [ START WITH start | INCREMENT BY step ] [ ...] ) ] | DEFAULT default_expression | COMMENT column_comment | column_constraint | MASK clause } [ ... ]
table_clauses
{ USING DELTA
PARTITIONED BY (col [, ...]) |
CLUSTER BY clause |
LOCATION path |
COMMENT view_comment |
TBLPROPERTIES clause |
WITH { ROW FILTER clause } } [ ... ]
} [ ... ]
flow_clause
FLOW { { INSERT [ONCE] BY NAME query } |
{ AUTO CDC auto_cdc_flow_spec } |
{ REPLACE WHERE predicate BY NAME query } |
{ REPLACE USING ( column_name [, ...] ) SEQUENCE BY sequence_column BY NAME query } }
Paraméterek
REFRESH
Ha meg van adva, létrehozza a táblát, vagy frissít egy meglévő táblát és annak tartalmát.
PRIVÁT
Létrehoz egy privát streamelési táblát.
- Ezek nem lettek hozzáadva a katalógushoz, és csak a definiáló folyamaton belül érhetők el
- Ugyanazzal a névvel rendelkezhetnek, mint egy meglévő objektum a katalógusban. A folyamaton belül, ha egy privát streamelési tábla és egy katalógusbeli objektum neve megegyezik, a névre mutató hivatkozás a privát streamelési táblára mutat.
- A privát streamelési táblák csak a folyamat teljes élettartama alatt vannak megőrzve, nem csak egyetlen frissítéssel.
A privát streamelési táblákat korábban a
TEMPORARYparaméterrel hozták létre.table_name
Az újonnan létrehozott tábla neve. A teljes táblanévnek egyedinek kell lennie.
táblázat_specifikáció
Ez az opcionális záradék határozza meg az oszlopok listáját, azok típusait, tulajdonságait, leírását és oszlopkorlátjait.
-
Az oszlopneveknek egyedinek kell lenniük, és le kell képezniük a lekérdezés kimeneti oszlopait.
-
Megadja az oszlop adattípusát. Nem minden adattípust, amelyet az Azure Databricks támogat, támogatnak a streamelési táblák.
column_comment
Egy tetszőleges
STRINGliterál, amely az oszlopot írja le. Ezt a beállítást a beállítássalcolumn_typeegyütt kell megadni. Ha az oszloptípus nincs megadva, a program kihagyja az oszlop megjegyzését.MINDIG GENERÁLT ( expr )
A záradék megadásakor az oszlop értékét a megadott
exprérték határozza meg.A táblázatnak
DEFAULT COLLATIONUTF8_BINARYkell lennie.exprállhat literálokból, táblán belüli oszlopazonosítókból és determinista, beépített SQL-függvényekből vagy operátorokból, kivéve:- Összesítő függvények
- Analitikai ablakfüggvények
- Rangsorolási ablakfüggvények
- Értékadó generátorfunkciók tábla szerint
- Olyan oszlopok, amelyek nem
UTF8_BINARYtípusú rendezéssel rendelkeznek
Emellett
exprnem tartalmazhat lekérdezést.GENERÁLVA { ALWAYS | ALAPÉRTELMEZÉS SZERINT } IDENTITÁSKÉNT [ ( [ KEZDÉS ] [ NÖVEKMÉNY LÉPÉSENKÉNT ] ) ] ]
A következőkre vonatkozik:
Databricks SQL
Databricks Runtime 10.4 LTS és újabbIdentitásoszlopot definiál. Amikor a táblába ír, és nem ad meg értékeket az identitásoszlophoz, a rendszer automatikusan hozzárendel egy egyedi és statisztikailag növekvő értéket (illetve csökkenőt, ha a
stepnegatív). Ez a záradék csak Delta-táblák esetében támogatott. Ez a záradék csak BIGINT adattípusú oszlopokhoz használható.Az automatikusan hozzárendelt értékek a következővel
startkezdődnek és haladnak a következővelstep: . A hozzárendelt értékek egyediek, de nem garantáltan egybefüggőek. Mindkét paraméter megadása nem kötelező, az alapértelmezett érték pedig 1.stepnem lehet0.Ha az automatikusan hozzárendelt értékek túllépik az identitásoszlop típusának tartományát, a lekérdezés sikertelen lesz.
Használat esetén
ALWAYSnem adhat meg saját értékeket az identitásoszlophoz.A következő műveletek nem támogatottak:
-
PARTITIONED BYazonosító oszlop -
UPDATEazonosító oszlop
Megjegyzés:
Az identitásoszlop táblán való deklarálása letiltja az egyidejű tranzakciókat. Csak olyan esetekben használjon identitásoszlopokat, amikor nincs szükség egyidejű írásra a céltáblába.
-
ALAPÉRTELMEZETT default_expression
A következőkre vonatkozik:
Databricks SQL
Databricks Runtime 11.3 LTS és újabbMeghatároz egy
DEFAULTértéket az oszlophoz, amelyet aINSERT,UPDATEésMERGE ... INSERThasznál, amikor az oszlop nincs meghatározva.Ha nincs megadva alapértelmezett érték, a
DEFAULT NULLalkalmazásra kerül a null értékű oszlopokhoz.default_expressionlehetnek literálok, és beépített SQL-függvények vagy operátorok, kivéve:- Összesítő függvények
- Analitikai ablakfüggvények
- Rangsorolási ablakfüggvények
- Értékadó generátorfunkciók tábla szerint
Emellett
default_expressionnem tartalmazhat lekérdezést.A
DEFAULT,CSV,JSON,PARQUETésORCforrások támogatottak.-
Információs elsődleges kulcsot vagy információs idegenkulcs-korlátozást ad hozzá a streamelési tábla oszlopához.
-
Olyan oszlopmaszk funkciót ad hozzá, amely lehetővé teszi a bizalmas adatok anonimizálását.
Lásd: Sorszűrők és oszlopmaszkok.
CONSTRAINT expectation_name EXPECT (expectation_expr) [ SÉRTÉS ESETÉN { HIBA UPDATE | SOR ELHAGYÁSA } ]
Adatminőségi elvárásokat ad hozzá a streamelési táblához. Ezek az adatminőségi elvárások idővel nyomon követhetők és elérhetők a streamelési tábla eseménynaplójában. A
FAIL UPDATEelvárás miatt a feldolgozás sikertelen lesz a tábla létrehozásakor és a tábla frissítésekor is. HaDROP ROWa várakozás nem teljesül, az egész sor elvetésre kerül. Lásd: Adatminőség kezelése folyamatelvárásokkal.expectation_exprállhat literálokból, táblán belüli oszlopazonosítókból és determinista, beépített SQL-függvényekből vagy operátorokból, kivéve:-
Összesítő függvények
- Analitikai ablakfüggvények
- Rangsorolási ablakfüggvények
- Értékadó generátorfunkciók tábla szerint
Emellett
exprnem tartalmazhat lekérdezést.-
Összesítő függvények
-
tábla_korlátozás
Séma megadásakor megadhatja az elsődleges és az idegen kulcsokat. A korlátozások tájékoztató jellegűek, és nincsenek kényszerítve. Tekintse meg a CONSTRAINT záradékot az SQL nyelvi hivatkozásában.
Megjegyzés:
A táblakorlátozások meghatározásához a folyamatnak Unity Catalog-kompatibilis folyamatnak kell lennie.
táblázat_feltételek
Opcionálisan megadhatja a tábla particionálási, megjegyzési és felhasználó által definiált tulajdonságait. Minden al záradék csak egyszer adható meg.
A DELTA HASZNÁLATA
Megadja az adatformátumot. Az egyetlen lehetőség a DELTA.
Ez a záradék opcionális, és alapértelmezés szerint a DELTA értéket veszi fel.
PARTÍCIÓVAL
A tábla particionálásához használandó egy vagy több oszlop választható listája. Kölcsönösen kizárja egymást a
CLUSTER BY-vel.A folyékony klaszterezés rugalmas, optimalizált megoldást biztosít a csoportosításhoz. Fontolja meg a
CLUSTER BYhasználatátPARTITIONED BYhelyett a folyamatvezetékek számára.CLUSTER BY
Engedélyezze a "liquid clustering" funkciót a táblában, és határozza meg a fürtözési kulcsként használni kívánt oszlopokat. Használjon automatikus folyékony fürtözést,
CLUSTER BY AUTOés a Databricks intelligensen választja ki a fürtözési kulcsokat a lekérdezési teljesítmény optimalizálásához. Kölcsönösen kizárja egymást aPARTITIONED BY-vel.HELYSZÍN
A táblaadatok opcionális tárolási helye. Ha nincs beállítva, a rendszer alapértelmezés szerint a folyamat tárolási helyére lesz beállítva.
MEGJEGYZÉS
A táblázat leírásához választható szó szerinti
STRINGkifejezés.TBLPROPERTIES
A tábla táblatulajdonságainak választható listája.
VAL ROW FILTER
Sorszűrő függvényt ad hozzá a táblához. A tábla jövőbeli lekérdezései azoknak a soroknak a részhalmazát kapják meg, amelyekre a függvény IGAZ értéket ad. Ez a részletes hozzáférés-vezérléshez hasznos, mert lehetővé teszi, hogy a függvény megvizsgálja a behívó felhasználó identitását és csoporttagságát, hogy eldöntse, szűr-e bizonyos sorokat.
Lásd
ROW FILTERzáradék.ÁRAMLÁS
Igény szerint beágyazott folyamatot definiál a tábla létrehozásával. A folyamat állapotalapú lekérdezés, amely frissíti a tábla tartalmát. Ha
FLOWnincs megadva, használhatja helyette, vagy külön definiálhatjaAS querya folyamatokat.CREATE FLOWA következő folyamattípusok egyikét adhatja meg:INSERT NÉV SZERINT
Adatokat szúr be a táblába oszlopnév alapján. Ha a
ONCEbeállítás nincs megadva, a lekérdezésnek streamelési lekérdezésnek kell lennie. Használja aSTREAMkulcsszót stream szemantikával, hogy olvasson a forrásból. Ha az olvasás egy meglévő rekord módosítását vagy törlését tapasztalja, hibaüzenet jelenik meg. A legbiztonságosabb, ha statikus vagy csak hozzáfűző forrásokból olvas.Megjegyzés:
FLOW INSERT BY NAMEa egyenértékű a használatávalAS query. A következő két utasítás viselkedése azonos:CREATE OR REFRESH STREAMING TABLE raw_data AS SELECT * FROM STREAM read_files('abfss://my_path'); CREATE OR REFRESH STREAMING TABLE raw_data FLOW INSERT BY NAME SELECT * FROM STREAM read_files('abfss://my_path');EGYSZER
Igény szerint a folyamatot egyszeri folyamatként, például visszatöltésként definiálja. Ha
ONCEmeg van adva, a lekérdezés nem streamelési lekérdezés, és a folyamat alapértelmezés szerint egyszer fut. Ha a tábla teljes frissítéssel frissül, aONCEfolyamat újra fut az adatok újbóli létrehozásához.ONCEcsak a folyamatokraINSERT BY NAMEvonatkozik.AUTO CDCFontos
Elérhető a Databricks Runtime 17.3-ban és a
PREVIEWPipelines-csatornában.AUTO CDCOlyan folyamatot definiál, amely egy forrás adatrögzítési (CDC) rekordjait dolgozza fel a táblába. Akkor használhatóAUTO CDC, ha a forrásadatok CDC szemantikát tartalmaznak. Lásd az AUTO CDC API-k: Egyszerűsítse a változáskövető adatrögzítést a csővezetékekkel.REPLACE WHEREpredikate BY NAME query
REPLACE WHEREOlyan folyamatot definiál, amely csak az egyezőpredicatesorokat írja át újra és írja felül, így az összes többi sort érintetlenül hagyja. Illesztések és aggregációk növekményes kötegelt feldolgozásához, későn érkező adatokhoz, sémafejlődéshez és visszatöltésekhez használhatóREPLACE WHERE. ABY NAMEhasználata kötelező. Lásd: Batch-feldolgozás CSERE WHERE folyamatokkal.CSERÉLD KIHASZNÁLVA ( column_name [, ...] ) SORREND sequence_column NÉV SZERINTI lekérdezés
Fontos
Ez a funkció bétaverzióban érhető el. Databricks Runtime 18.2 vagy annál magasabb verziót igényel.
Olyan
REPLACE USINGfolyamatot definiál, amely lecseréli az összes megadott kulcsoszlophoz tartozó sort, és a többi sort érintetlenül hagyja. Használd,REPLACE USINGha a forrásod részleges pillanatképek sorozata, amelyeket oszlopok szerint kulcsolnak.SEQUENCE BYRendeli a frissítéseket, így a kulcs legmagasabb sorrendje nyer, még akkor is, ha a frissítések sorrendben érkeznek. A forrásnak streamelési forrásnak kell lennie. ABY NAMEhasználata kötelező. Lásd Részleges snapshot helyettesítés CSERÉLJ USING flow-okkal.
-
Ez a záradék feltölti a táblát a
queryadataival. Ennek a lekérdezésnek egy streamelési lekérdezésnek kell lennie. A STREAM kulcsszóval stream-szemantikát használhat a forrásból való olvasáshoz. Ha az olvasás egy meglévő rekord módosítását vagy törlését tapasztalja, hibaüzenet jelenik meg. A legbiztonságosabb, ha statikus vagy csak hozzáfűző forrásokból olvas. A módosítási véglegesítéseket tartalmazó adatok betöltéséhez hozzáadhatja az olvasási lehetőséget askipChangeCommitshibák kezeléséhez.Ha egy
query-t és egytable_specification-et együtt ad meg, atable_specification-ben megadott táblasémának tartalmaznia kell aqueryáltal visszaadott összes oszlopot, ellenkező esetben hibaüzenetet kap. Atable_specificationmegadott, dequeryáltal nem visszaadott oszlopoknullértékeket ad vissza lekérdezéskor.További információért az adatfolyamokról, tekintse meg a Folyamatokkal történő adatátalakítást.
Olvasási beállítások
A lekérdezésben megadhatja az olvasási beállításokat az adatok forrásból való beolvasásának konfigurálásához. Megadhatja
skipChangeCommitspéldául, hogy a forrásadatokban szereplő módosítási véglegesítéseket átugorja. Az olvasási beállítások térképként vannak megadva aWITHlekérdezés záradékában. Például:SELECT * FROM STREAM source_table WITH (SKIPCHANGECOMMITS=TRUE, STARTINGVERSION=X)Ez
=TRUEnem kötelező, így az alábbihoz hasonló logikai beállítást is megadhat:SELECT * FROM STREAM source_table WITH (SKIPCHANGECOMMITS)Megjegyzés:
Az olvasási lehetőségek csak a Databricks Runtime 17.3-at vagy újabb verzióját támogatják.
Az alábbi olvasási lehetőségek támogatottak a Delta esetében, az egyes lehetőségek részleteiért lásd a Delta Lake-tábla streamelt olvasásait és írásait.
maxFilesPerTriggermaxBytesPerTriggerstartingVersionstartingTimestampreadChangeFeedwithEventTimeOrderskipChangeCommits
Szükséges engedélyek
A folyamat futtató felhasználójának a következő engedélyekkel kell rendelkeznie:
-
SELECTengedély az alaptáblákra, amelyekre a streamelési tábla hivatkozik. -
USE CATALOGjogosultsággal a szülőkatalóguson, valamintUSE SCHEMAjogosultsággal a szülősémán. -
CREATE MATERIALIZED VIEWjogosultság a streaming tábla sémájára.
Ahhoz, hogy a felhasználó frissíthesse a streamelési táblában definiált folyamatot, a következőt kell megkövetelnie:
-
USE CATALOGjogosultsággal a szülőkatalóguson, valamintUSE SCHEMAjogosultsággal a szülősémán. - A streamelési tábla tulajdonjoga vagy
REFRESHjogosultság a streamelési táblán. - A streamelési tábla tulajdonosának jogosultsággal
SELECTkell rendelkeznie a streamelési tábla által hivatkozott alaptáblák felett.
Ahhoz, hogy egy felhasználó le tudja kérdezni az eredményként kapott streamelési táblát, a következőre van szükség:
-
USE CATALOGjogosultsággal a szülőkatalóguson, valamintUSE SCHEMAjogosultsággal a szülősémán. -
SELECTjogosultság a streamelési tábla felett.
Korlátozások
- Csak a táblatulajdonosok frissíthetik a streamelő táblákat a legújabb adatok lekéréséhez.
-
ALTER TABLEparancsok nem engedélyezettek a streamelési táblákon. A tábla definícióját és tulajdonságait aCREATE OR REFRESHvagy ALTER STREAMING TABLE utasítással kell módosítani. - A táblaséma DML-parancsokkal (például
INSERT INTO) történő fejlesztése és aMERGEnem támogatott. - A streamelési táblákban a következő parancsok nem támogatottak:
CREATE TABLE ... CLONE <streaming_table>COPY INTOANALYZE TABLERESTORETRUNCATEGENERATE MANIFEST[CREATE OR] REPLACE TABLE
- A tábla átnevezése vagy a tulajdonos megváltoztatása nem támogatott.
Példák
-- Define a streaming table from a volume of files:
CREATE OR REFRESH STREAMING TABLE customers_bronze
AS SELECT * FROM STREAM read_files("/databricks-datasets/retail-org/customers/*", format => "csv")
-- Define a streaming table from a streaming source table:
CREATE OR REFRESH STREAMING TABLE customers_silver
AS SELECT * FROM STREAM(customers_bronze)
-- Use automatic liquid clustering to let Databricks choose the clustering columns:
CREATE OR REFRESH STREAMING TABLE customers_bronze_auto
CLUSTER BY AUTO
AS SELECT * FROM STREAM read_files("/databricks-datasets/retail-org/customers/*", format => "csv")
-- Define a table with a row filter and column mask:
CREATE OR REFRESH STREAMING TABLE customers_silver (
id int COMMENT 'This is the customer ID',
name string,
region string,
ssn string MASK catalog.schema.ssn_mask_fn COMMENT 'SSN masked for privacy'
)
WITH ROW FILTER catalog.schema.us_filter_fn ON (region)
AS SELECT * FROM STREAM(customers_bronze)
-- Define a streaming table with an identity column:
CREATE OR REFRESH STREAMING TABLE customers_with_id (
customer_id BIGINT GENERATED ALWAYS AS IDENTITY,
name string,
region string
)
AS SELECT name, region FROM STREAM(customers_bronze)
-- Define a streaming table that you can add flows into:
CREATE OR REFRESH STREAMING TABLE orders;
-- Define a streaming table with an inline append flow:
CREATE OR REFRESH STREAMING TABLE raw_data
FLOW INSERT BY NAME SELECT * FROM STREAM read_files('abfss://my_path');
-- Define a streaming table with an inline AUTO CDC flow:
CREATE OR REFRESH STREAMING TABLE target
FLOW AUTO CDC
FROM stream(cdc_data.users)
KEYS (userId)
SEQUENCE BY sequenceNum
STORED AS SCD TYPE 1;
-- Define a streaming table with an inline REPLACE USING flow that keeps the latest
-- row for each payment_id:
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);