CREATE FLOW (csővezetékek)

A CREATE FLOW utasítással folyamatokat vagy visszatöltéseket hozhat létre egy csővezeték tábláihoz.

Szemantika

CREATE FLOW flow_name [COMMENT comment] AS
{
  AUTO CDC [ONCE] INTO target_table create_auto_cdc_flow_spec |
  INSERT [ONCE] INTO target_table BY NAME [ replace_using_spec ] query
}

replace_using_spec
  REPLACE USING ( column_name [, ...] ) SEQUENCE BY sequence_column

Paraméterek

  • flow_name

    A létrehozandó folyamat neve.

  • MEGJEGYZÉS

    A folyamat opcionális leírása.

  • AUTOMATIKUS CDC-BE

    Egy AUTO CDC ... INTO utasítás, amely meghatározza a folyamatot egy create_auto_cdc_flow_spec. Tartalmaznia kell vagy egy AUTO CDC ... INTO utasítást, vagy egy INSERT INTO utasítást. Akkor használható AUTO CDC ... INTO , ha a forrás lekérdezés változásadat-szemantikát használ.

    További információ: AUTO CDC INTO (adattovábbítási csatornák).

  • target_table

    A frissíteni kívánt táblázat. Ennek streamelésre alkalmas táblának kell lennie.

  • INSERT BE

    A céltáblába való beszúrásra szánt táblalekérdezést definiál. Ha a ONCE beállítás nincs megadva, a lekérdezésnek 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 használhatja a Pythont és a skipChangeCommits hibák kezelésére szolgáló lehetőséget.

    INSERT INTO kölcsönösen kizárja a AUTO CDC ... INTO. Akkor használható AUTO CDC ... INTO , ha a forrásadatok módosítási adatrögzítési (CDC) funkciót tartalmaznak. A INSERT INTO elemet akkor használja, ha a forrásban nem szerepel.

    További információért az adatfolyamokról, tekintse meg a Folyamatokkal történő adatátalakítást.

  • CSERÉLD KIHASZNÁLVA ( column_name [, ...] ) SOROZAT sequence_column

    Important

    Ez a funkció bétaverzióban érhető el. Databricks Runtime 18.2 vagy annál magasabb verziót igényel.

    Definiálja az áramlást, REPLACE USING amely a céltáblában lévő összes sort lecseréli, hogy illeszkedjen a megadott kulcsoszlopokhoz, és minden más 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.

    Legalább egy kulcsoszlopot és pontosan egy SEQUENCE BY oszlopot kell megadni. A lekérdezésnek streaming lekérdezésnek kell lennie, és BY NAME kötelező. REPLACE USING Nem lehet kombinálni vagy ONCE kombinálni AUTO CDC ... INTO.

    További információért lásd: Részleges snapshot helyettesítés CSERÉLJ USING flow-okkal.

  • EGYSZER

    Igény szerint a folyamatot egyszeri folyamatként, például visszatöltésként definiálhatja. A ONCE használata kétféleképpen változtatja meg a folyamatot:

    • A forrás query vagy create_auto_cdc_flow_spec nem egy folyamatos tábla.
    • A folyamat alapértelmezés szerint egyszer fut. Ha a folyamat teljes frissítéssel frissül, akkor a ONCE folyamat újra fut az adatok újbóli létrehozásához.

    ONCE Nem használható , REPLACE USINGami streaming forrást igényel.

Példák

-- EXAMPLE 1:
-- Create a streaming table, and add two flows that append data to it:
CREATE OR REFRESH STREAMING TABLE users;

-- first flow into target_table:
CREATE FLOW users_flow AS
INSERT INTO users BY NAME
SELECT * FROM stream(raw_data.users);

-- second flow into target_table:
CREATE FLOW backfill_users AS
INSERT ONCE INTO users BY NAME
SELECT * FROM user_backfill_table;

-- EXAMPLE 2:
-- Create a streaming table, and add a flow that applies CDC changes to it:
CREATE OR REFRESH STREAMING TABLE admins_cdc_target_table;

-- first flow into target_table:
CREATE FLOW admin_cdc_flow AS
AUTO CDC INTO admins_cdc_target_table
FROM stream(cdc_data.admins)
KEYS (userId)
APPLY AS DELETE WHEN
  operation = "DELETE"
SEQUENCE BY sequenceNum
COLUMNS * EXCEPT (operation, sequenceNum)
STORED AS SCD TYPE 2;

-- EXAMPLE 3:
-- Create a streaming table, and add a REPLACE USING flow that keeps the latest
-- row for each payment_id from a stream of partial snapshots:
CREATE OR REFRESH STREAMING TABLE payments_latest;

CREATE FLOW payments_replace_flow AS
INSERT INTO payments_latest BY NAME
REPLACE USING (payment_id) SEQUENCE BY payment_date
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);