Catatan
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba masuk atau mengubah direktori.
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba mengubah direktori.
Gunakan pernyataan CREATE FLOW untuk membuat alur atau pengisian ulang bagi tabel dalam sebuah pipeline.
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
Parameter-parameternya
flow_name
Nama alur yang akan dibuat.
KOMENTAR
Deskripsi opsional untuk alur.
-
Perintah
AUTO CDC ... INTOyang menentukan alur, dengancreate_auto_cdc_flow_spec. Anda harus menyertakanAUTO CDC ... INTOpernyataan, atauINSERT INTOpernyataan. GunakanAUTO CDC ... INTOsaat kueri sumber menggunakan semantik data perubahan.Untuk informasi selengkapnya, lihat AUTO CDC INTO (pipeline).
target_table
Tabel yang akan diperbarui. Ini harus berupa tabel streaming.
INSERT KE
Menentukan suatu kueri tabel yang disisipkan ke dalam tabel target.
ONCEJika opsi tidak disediakan, kueri harus berupa kueri streaming. Gunakan kata kunci STREAM untuk menggunakan semantik streaming untuk membaca dari sumbernya. Jika pembacaan mengalami perubahan atau penghapusan pada rekaman yang ada, akan menghasilkan kesalahan. Paling aman untuk membaca dari sumber statis atau yang hanya bisa ditambahkan. Untuk memasukkan data yang memiliki komit perubahan, Anda dapat menggunakan Python dan opsiskipChangeCommitsuntuk menangani kesalahan.INSERT INTOsaling eksklusif denganAUTO CDC ... INTO. GunakanAUTO CDC ... INTOsaat data sumber menyertakan fungsionalitas change data capture (CDC). GunakanINSERT INTOketika sumber tidak melakukannya.Untuk informasi selengkapnya tentang data streaming, lihat Mengubah data dengan alur.
GANTI MENGGUNAKAN ( column_name [, ...] ) URUTAN OLEH sequence_column
Important
Fitur ini ada di Beta. Memerlukan Databricks Runtime 18.2 ke atas.
Mendefinisikan alur sebagai
REPLACE USINGaliran, yang menggantikan semua baris dalam tabel target yang sesuai dengan kolom kunci yang ditentukan dan membiarkan semua baris lainnya tidak disentuh. GunakanREPLACE USINGketika sumber Anda adalah serangkaian snapshot parsial yang dikunci berdasarkan kolom.SEQUENCE BYmengatur pembaruan sehingga urutan tertinggi untuk sebuah kunci menang, bahkan ketika pembaruan datang tidak sesuai urutan.Tentukan setidaknya satu kolom kunci dan tepat satu
SEQUENCE BYkolom. Kueri harus berupa kueri streaming, danBY NAMEwajib dilakukan.REPLACE USINGtidak dapat dikombinasikan denganONCEatau denganAUTO CDC ... INTO.Untuk informasi lebih lanjut, lihat Penggantian snapshot parsial dengan alur GANTI MENGGUNAKAN.
SEKALI
Secara opsional tentukan alur sebagai aliran satu kali, seperti isi ulang. Dengan menggunakan
ONCE, alurnya dapat diubah dengan dua cara:- Sumber
queryataucreate_auto_cdc_flow_specbukan tabel streaming. - Alur dijalankan satu kali secara default. Jika alur diperbarui dengan pembaruan lengkap, maka
ONCEalur akan berjalan kembali untuk membuat ulang data.
ONCEtidak dapat digunakan denganREPLACE USING, yang memerlukan sumber streaming.- Sumber
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);