CREATE FLOW (işlem hatları)

CREATE FLOW ifadesini bir işlem hattındaki tablolar için akışlar veya geri yüklemeler oluşturmak için kullanın.

Sözdizimi

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

Parametreler

  • flow_name

    Oluşturulacak akışın adı.

  • YORUM

    Akış için isteğe bağlı bir açıklama.

  • OTOMATIK CDC INTO

    AUTO CDC ... INTO ile akışı tanımlayan bir create_auto_cdc_flow_spec ifadesi. Ya bir AUTO CDC ... INTO ifadesi ya da bir INSERT INTO ifadesi eklemelisiniz. Kaynak sorgu değişiklik verisi semantiğini kullandığında kullanın AUTO CDC ... INTO .

    Daha fazla bilgi için bkz. AUTO CDC INTO (işlem hatları).

  • target_table

    Güncelleştirilecek tablo. Bu bir Akış tablosu olmalıdır.

  • INSERT İÇİNE

    Hedef tabloya eklenen tablo sorgusunu tanımlar. Seçenek ONCE, verilmezse sorgu bir akış sorgusu olmalıdır. Kaynaktan okumak üzere akış semantiğini kullanmak için STREAM anahtar sözcüğünü kullanın. Okuma işlemi var olan bir kayıtta bir değişiklik veya silme işlemiyle karşılaşırsa bir hata oluşur. Statik veya yalnızca ekleme kaynaklarından okumak en güvenlidir. Değişiklik taahhütleri içeren verileri içeri çekmek için Python'ı ve skipChangeCommits seçeneğini hataları işlemek için kullanabilirsiniz.

    INSERT INTO ile AUTO CDC ... INTObirbirini dışlar. Kaynak veriler değişiklik verisi yakalama (CDC) işlevselliği içerdiğinde kullanın AUTO CDC ... INTO . Kaynakta kullanılmadığında INSERT INTO kullanın.

    Akış verileri hakkında daha fazla bilgi için bkz. İşlem hatları ile veri dönüştürme.

  • YERINE ÇALIN ( column_name [, ...] ) DIZISI sequence_column

    Important

    Bu özellik Beta sürümündedir. Databricks Runtime 18.2 ve üzeri gerektirir.

    Akışı bir akış REPLACE USING olarak tanımlar; hedef tablodaki tüm satırları belirtilen anahtar sütunlara uydurur ve diğer tüm satırlar dokunulmadan bırakır. Kaynağınız sütunlara göre anahtarlanmış kısmi anlık görüntüler dizisi olduğunda kullanın REPLACE USING . SEQUENCE BY güncellemeleri sipariş eder, böylece anahtar için en yüksek sıra, güncellemeler sırasız gelse bile kazanır.

    En az bir anahtar sütunu ve tam olarak bir SEQUENCE BY sütun belirtin. Sorgu akış sorgusu olmalı ve BY NAME zorunludur. REPLACE USINGile veya AUTO CDC ... INTObirleştirilemezONCE.

    Daha fazla bilgi için, REPLACE USING akışlarıyla kısmi anlık görüntü değiştirme bölümünü inceleyebilirsiniz.

  • BİR DEFA

    İsteğe bağlı olarak akışı bir kerelik akış olarak (örneğin, bir geri doldurma) tanımlayın. Kullanımı ONCE , akışı iki şekilde değiştirir:

    • Kaynak query veya create_auto_cdc_flow_spec bir akış tablosu değil.
    • Akış varsayılan olarak bir kez çalıştırılır. Eğer işlem hattı eksiksiz bir yenilemeyle güncellenirse, ONCE akış verileri yeniden oluşturmak için tekrar çalıştırılır.

    ONCE ile REPLACE USINGkullanılamaz, bu da bir akış kaynağı gerektirir.

Örnekler

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