CREATE STREAMING TABLE (csővezetékek)

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 TEMPORARY paramé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.

    • column_identifier

      Az oszlopneveknek egyedinek kell lenniük, és le kell képezniük a lekérdezés kimeneti oszlopait.

    • oszlop_típus

      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 STRING literál, amely az oszlopot írja le. Ezt a beállítást a beállítással column_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_BINARY kell 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:

      Emellett expr nem 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:jelölje be az igennel jelölt jelölőnégyzetet Databricks SQL jelölje be az igennel jelölt jelölőnégyzetet Databricks Runtime 10.4 LTS és újabb

      Identitá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 step negatí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 start kezdődnek és haladnak a következővel step: . 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. step nem lehet 0.

      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 ALWAYS nem adhat meg saját értékeket az identitásoszlophoz.

      A következő műveletek nem támogatottak:

      • PARTITIONED BY azonosító oszlop
      • UPDATE azonosí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:igen, be van jelölve Databricks SQL igen, be van jelölve Databricks Runtime 11.3 LTS és újabb

      Meghatároz egy DEFAULT értéket az oszlophoz, amelyet a INSERT, UPDATE és MERGE ... INSERT használ, amikor az oszlop nincs meghatározva.

      Ha nincs megadva alapértelmezett érték, a DEFAULT NULL alkalmazásra kerül a null értékű oszlopokhoz.

      default_expression lehetnek literálok, és beépített SQL-függvények vagy operátorok, kivéve:

      Emellett default_expression nem tartalmazhat lekérdezést.

      A DEFAULT, CSV, JSON, PARQUET és ORC források támogatottak.

    • column_constraint

      Információs elsődleges kulcsot vagy információs idegenkulcs-korlátozást ad hozzá a streamelési tábla oszlopához.

    • Maszk záradék

      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 UPDATE elvárás miatt a feldolgozás sikertelen lesz a tábla létrehozásakor és a tábla frissítésekor is. Ha DROP ROW a 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:

      Emellett expr nem tartalmazhat lekérdezést.

  • 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 BY használatát PARTITIONED BY helyett 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 a PARTITIONED BY-vel.

      Lásd: Táblákhoz folyékony klaszterezés használata.

    • 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 STRING kifejezé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 FILTER zá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 FLOW nincs megadva, használhatja helyette, vagy külön definiálhatja AS query a folyamatokat.CREATE FLOW A 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 ONCE beállítás nincs megadva, a lekérdezésnek streamelési lekérdezésnek kell lennie. Használja a STREAM kulcsszó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 NAME a egyenértékű a használatával AS 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 ONCE meg 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, a ONCE folyamat újra fut az adatok újbóli létrehozásához. ONCE csak a folyamatokra INSERT BY NAME vonatkozik.

      • AUTO CDC

        Fontos

        Elérhető a Databricks Runtime 17.3-ban és a PREVIEW Pipelines-csatornában.

        AUTO CDC Olyan 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 WHERE Olyan 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 . A BY NAME haszná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 USING folyamatot definiál, amely lecseréli az összes megadott kulcsoszlophoz tartozó sort, és a többi sort érintetlenül hagyja. Használd, REPLACE USING ha a forrásod részleges pillanatképek sorozata, amelyeket oszlopok szerint kulcsolnak. SEQUENCE BY Rendeli 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. A BY NAME használata kötelező. Lásd Részleges snapshot helyettesítés CSERÉLJ USING flow-okkal.

  • AS-lekérdezés

    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 a skipChangeCommits hibák kezeléséhez.

    Ha egy query-t és egy table_specification-et együtt ad meg, a table_specification-ben megadott táblasémának tartalmaznia kell a queryáltal visszaadott összes oszlopot, ellenkező esetben hibaüzenetet kap. A table_specification megadott, de query által nem visszaadott oszlopok null é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 skipChangeCommits pé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 a WITH lekérdezés záradékában. Például:

      SELECT * FROM STREAM source_table WITH (SKIPCHANGECOMMITS=TRUE, STARTINGVERSION=X)
      

      Ez =TRUE nem 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.

      • maxFilesPerTrigger
      • maxBytesPerTrigger
      • startingVersion
      • startingTimestamp
      • readChangeFeed
      • withEventTimeOrder
      • skipChangeCommits

Szükséges engedélyek

A folyamat futtató felhasználójának a következő engedélyekkel kell rendelkeznie:

  • SELECT engedély az alaptáblákra, amelyekre a streamelési tábla hivatkozik.
  • USE CATALOG jogosultsággal a szülőkatalóguson, valamint USE SCHEMA jogosultsággal a szülősémán.
  • CREATE MATERIALIZED VIEW jogosultsá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 CATALOG jogosultsággal a szülőkatalóguson, valamint USE SCHEMA jogosultsággal a szülősémán.
  • A streamelési tábla tulajdonjoga vagy REFRESH jogosultság a streamelési táblán.
  • A streamelési tábla tulajdonosának jogosultsággal SELECT kell 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 CATALOG jogosultsággal a szülőkatalóguson, valamint USE SCHEMA jogosultsággal a szülősémán.
  • SELECT jogosultsá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 TABLE parancsok nem engedélyezettek a streamelési táblákon. A tábla definícióját és tulajdonságait a CREATE OR REFRESH vagy 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 a MERGE nem támogatott.
  • A streamelési táblákban a következő parancsok nem támogatottak:
    • CREATE TABLE ... CLONE <streaming_table>
    • COPY INTO
    • ANALYZE TABLE
    • RESTORE
    • TRUNCATE
    • GENERATE 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);