CREATE FLOW (Pipelines)

Verwenden Sie die CREATE FLOW Anweisung, um Datenflüsse oder Rückfüllungen für Tabellen in einer Pipeline zu erstellen.

Syntax

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

Die Parameter

  • flow_name

    Der Name des zu erstellenden Flusses.

  • KOMMENTAR

    Eine optionale Beschreibung für den Fluss.

  • AUTOMATISCHE CDC IN

    Eine AUTO CDC ... INTO Anweisung, die den Ablauf mithilfe einer create_auto_cdc_flow_spec definiert. Sie müssen entweder eine AUTO CDC ... INTO Anweisung oder eine INSERT INTO Anweisung einschließen. Verwenden Sie AUTO CDC ... INTO bei Abfragen, wenn die Quellabfrage die Änderungsdatensemantik nutzt.

    Weitere Informationen finden Sie unter AUTO CDC INTO (Pipelines).

  • target_table

    Die Tabelle, die aktualisiert werden soll. Dies muss eine Streamingtabelle sein.

  • INSERT IN

    Definiert eine Tabellenabfrage, die in die Zieltabelle eingefügt wird. Wenn die Option nicht angegeben wird, muss es sich bei der ONCE Abfrage um eine Streamingabfrage handeln. Verwenden Sie das STREAM-Schlüsselwort, um Streamingsemantik zum Lesen aus der Quelle zu verwenden. Wenn beim Lesen eine Änderung oder löschung in einem vorhandenen Datensatz auftritt, wird ein Fehler ausgelöst. Es ist am sichersten, aus statischen oder nur angefügten Quellen zu lesen. Zum Einlesen von Daten mit Änderungs-Commits können Sie Python und die skipChangeCommits Option zur Fehlerbehandlung verwenden.

    INSERT INTO schließt sich mit AUTO CDC ... INTO gegenseitig aus. Verwenden Sie diese Funktion AUTO CDC ... INTO , wenn die Quelldaten die Funktionalität für die Änderungsdatenerfassung (Change Data Capture, CDC) enthalten. Verwenden Sie INSERT INTO, wenn die Quelle nicht verwendet wird.

    Weitere Informationen zum Streamen von Daten finden Sie unter Transformieren von Daten mit Pipelines.

  • ERSETZEN MIT ( column_name [, ...] ) SEQUENZ VON sequence_column

    Important

    Dieses Feature befindet sich in der Betaversion. Benötigt Databricks Runtime 18.2 und höher.

    Definiert den Fluss als Fluss REPLACE USING , der alle Zeilen in der Zieltabelle ersetzt, die mit den angegebenen Schlüsselspalten übereinstimmen, und alle anderen Zeilen unberührt lässt. Verwenden Sie REPLACE USING , wenn Ihre Quelle eine Reihe von teilweisen Schnappschüssen ist, die nach Spalten verschlüsselt sind. SEQUENCE BY ordnet die Updates so, dass die höchste Reihenfolge eines Schlüssels gewinnt, selbst wenn Updates nicht in der richtigen Reihenfolge eintreffen.

    Geben Sie mindestens eine Schlüsselspalte und genau eine SEQUENCE BY Spalte an. Die Abfrage muss eine Streaming-Abfrage sein und BY NAME ist erforderlich. REPLACE USING kann weder mit ONCE noch mit AUTO CDC ... INTOkombiniert werden.

    Weitere Informationen finden Sie unter Teilweise Snapshot-Ersetzung mit ERSETZEN UNTER Verwendung von Flows.

  • EINMAL

    Definieren Sie optional den Fluss als einmaligen Ablauf, z. B. als Rückfüllvorgang. Die Verwendung von ONCE ändert den Ablauf auf zwei Arten:

    • Die Quelle query oder create_auto_cdc_flow_spec ist keine Streamingtabelle.
    • Der Fluss wird standardmäßig einmal ausgeführt. Wenn die Pipeline mit einer vollständigen Aktualisierung aktualisiert wird, wird der ONCE Flow erneut ausgeführt, um die Daten neu zu erstellen.

    ONCE Kann nicht mit REPLACE USINGverwendet werden, was eine Streaming-Quelle erfordert.

Examples

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