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.
Tabel streaming adalah tabel dengan dukungan untuk pemrosesan data streaming atau inkremental. Tabel streaming didukung oleh pipa. Setiap kali tabel streaming di-refresh, data yang ditambahkan ke tabel sumber ditambahkan ke tabel streaming. Anda dapat me-refresh tabel streaming secara manual atau sesuai jadwal.
Untuk mempelajari selengkapnya tentang cara melakukan atau menjadwalkan refresh, lihat Menjalankan pembaruan alur.
Syntax
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 } }
Parameter-parameternya
REFRESH
Jika ditentukan, membuat tabel, atau memperbarui tabel yang sudah ada dan kontennya.
SWASTA
Membuat tabel streaming privat.
- Mereka tidak ditambahkan ke katalog dan hanya dapat diakses dalam alur yang menentukan
- Mereka dapat memiliki nama yang sama dengan objek yang ada di katalog. Dalam alur, jika tabel streaming privat dan objek dalam katalog memiliki nama yang sama, referensi ke nama diselesaikan ke tabel streaming privat.
- Tabel streaming privat hanya dipertahankan selama masa operasional alur, tidak hanya satu pembaruan.
Tabel streaming privat sebelumnya dibuat dengan
TEMPORARYparameter .table_name
Nama tabel yang baru dibuat. Nama tabel yang sepenuhnya memenuhi syarat harus unik.
spesifikasi_tabel
Klausa opsional ini menentukan daftar kolom, jenis, properti, deskripsi, dan batasan kolomnya.
-
Nama kolom harus unik dan sesuai dengan kolom output kueri.
-
Menentukan jenis data kolom. Tidak semua jenis data yang didukung oleh Azure Databricks didukung oleh tabel streaming.
column_comment
STRINGopsional yang menjelaskan kolom. Opsi ini harus ditentukan bersama dengancolumn_type. Jika jenis kolom tidak ditentukan, komentar kolom akan dilewati.DIHASILKAN SELALU SEBAGAI ( expr )
Ketika Anda menentukan klausa ini, nilai kolom ini ditentukan oleh
expryang ditentukan .DEFAULT COLLATIONdi tabel harusUTF8_BINARY.exprdapat terdiri dari literal, pengidentifikasi kolom dalam tabel, dan fungsi atau operator SQL bawaan deterministik kecuali:- Fungsi Agregat
- Fungsi jendela analitik
- Fungsi jendela peringkat
- Fungsi penghasil nilai tabel
- Kolom dengan kolasi selain
UTF8_BINARY
Juga
exprtidak boleh berisi subkueri apa pun.DIHASILKAN { SELALU | SECARA BAWAAN } SEBAGAI IDENTITAS [ ( [ MULAI DENGAN mulai ] [ DITAMBAH DENGAN langkah ] ) ]
Berlaku untuk:
Databricks SQL
Databricks Runtime 10.4 LTS ke atasMenentukan kolom identitas. Saat Anda menulis ke tabel dan tidak memberikan nilai untuk kolom identitas, kolom tersebut akan secara otomatis diberi nilai yang unik dan bertambah secara statistik (atau berkurang jika
stepnegatif). Klausa ini hanya didukung untuk tabel Delta. Klausa ini hanya dapat digunakan untuk kolom dengan jenis data BIGINT.Nilai yang ditetapkan secara otomatis dimulai dengan
startdan meningkat sesuai denganstep. Nilai yang ditetapkan unik tetapi tidak dijamin berurutan. Kedua parameter bersifat opsional, dan nilai defaultnya adalah 1.steptidak bisa menjadi0.Jika nilai yang ditetapkan secara otomatis berada di luar rentang tipe kolom identitas, kueri akan gagal.
Saat
ALWAYSdigunakan, Anda tidak dapat memberikan nilai Anda sendiri untuk kolom identitas.Operasi berikut tidak didukung:
-
PARTITIONED BYkolom identitas -
UPDATEkolom identitas
Nota
Mendeklarasikan kolom identitas pada tabel menonaktifkan transaksi bersamaan. Hanya gunakan kolom identitas dalam kasus penggunaan di mana penulisan bersamaan ke tabel target tidak diperlukan.
-
default_expression DEFAULT
Berlaku untuk:
Databricks SQL
Databricks Runtime 11.3 LTS ke atasMenentukan nilai
DEFAULTuntuk kolom yang digunakan padaINSERT,UPDATE, danMERGE ... INSERTsaat kolom tidak ditentukan.Jika tidak ada default yang dibuat,
DEFAULT NULLditerapkan untuk kolom yang dapat bernilai null.default_expressiondapat terdiri dari fungsi atau operator SQL literal dan bawaan kecuali:- Fungsi Agregat
- Fungsi jendela analitik
- Fungsi jendela peringkat
- Fungsi penghasil nilai tabel
Juga
default_expressiontidak boleh berisi subkueri apa pun.DEFAULTdidukung untuk sumberCSV,JSON,PARQUET, danORC.-
Menambahkan kunci primer informasi atau batasan kunci asing informasional ke kolom dalam tabel streaming.
-
Menambahkan fungsi masker kolom untuk menganonimkan data sensitif.
CONSTRAINT nama_ekspetasi EXPECT (ekspresi_ekspetasi) [ PADA PELANGGARAN { GAGAL UPDATE | HAPUS BARIS } ]
Menambahkan ekspektasi kualitas data ke tabel streaming. Harapan kualitas data ini dapat dilacak dari waktu ke waktu dan diakses melalui log peristiwa tabel streaming. Ekspektasi
FAIL UPDATEmenyebabkan pemrosesan gagal saat membuat tabel serta me-refresh tabel.DROP ROWEkspektasi menyebabkan seluruh baris dihapus jika ekspektasi tidak terpenuhi. Lihat Mengelola kualitas data dengan ekspektasi alur kerja.expectation_exprdapat terdiri dari literal, pengidentifikasi kolom dalam tabel, dan fungsi atau operator SQL bawaan deterministik kecuali:-
Fungsi Agregat
- Fungsi jendela analitik
- Fungsi jendela peringkat
- Fungsi penghasil nilai tabel
Juga
exprtidak boleh berisi subkueri apa pun.-
Fungsi Agregat
-
pembatasan_tabel
Saat menentukan skema, Anda dapat menentukan kunci primer dan asing. Batasan bersifat informasi dan tidak diberlakukan. Lihat klausa CONSTRAINT dalam referensi bahasa SQL.
Nota
Untuk menentukan batasan tabel, alur Anda harus berupa alur yang mendukung Katalog Unity.
table_clauses
Tentukan pemartisian, komentar, dan properti yang ditentukan pengguna secara opsional untuk tabel. Setiap sub klausul hanya dapat ditentukan satu kali.
MENGGUNAKAN DELTA
Menetapkan format data. Satu-satunya opsi adalah DELTA.
Klausa ini bersifat opsional, dan secara default menjadi DELTA.
DISEGMENKAN OLEH
Daftar opsional dari satu atau beberapa kolom yang akan digunakan untuk pemartisian dalam tabel. Bersifat saling eksklusif dengan
CLUSTER BY.Pengklusteran cairan memberikan solusi yang fleksibel dan dioptimalkan untuk pengklusteran. Pertimbangkan untuk menggunakan
CLUSTER BYalih-alihPARTITIONED BYuntuk alur.CLUSTER BY
Aktifkan pengklusteran cair pada tabel dan tentukan kolom yang akan digunakan sebagai kunci pengklusteran. Gunakan pengklusteran cairan otomatis dengan
CLUSTER BY AUTO, dan Databricks dengan cerdas memilih kunci pengklusteran untuk mengoptimalkan performa kueri. Bersifat saling eksklusif denganPARTITIONED BY.LOKASI
Lokasi penyimpanan opsional untuk data tabel. Jika tidak diatur, sistem default ke lokasi penyimpanan alur.
KOMENTAR
Literal
STRINGyang opsional untuk menjelaskan tabel.TBLPROPERTIES
Daftar properti tabel opsional untuk tabel.
DENGAN ROW FILTER
Menambahkan fungsi filter baris ke tabel. Kueri di masa mendatang untuk tabel tersebut akan mendapatkan bagian dari baris-baris yang fungsinya bernilai BENAR. Ini berguna untuk kontrol akses terperintah, karena memungkinkan fungsi untuk memeriksa identitas dan keanggotaan grup pengguna yang memanggil untuk memutuskan apakah akan memfilter baris tertentu.
Lihat klausa
ROW FILTER.ALIRAN
Secara opsional menentukan alur sebaris dengan pembuatan tabel. Alur adalah kueri stateful yang me-refresh konten tabel. Jika
FLOWtidak ditentukan, Anda dapat menggunakanAS querysebagai gantinya, atau menentukan alur secara terpisah denganCREATE FLOW. Anda dapat menentukan salah satu jenis alur berikut:INSERT MENURUT NAMA
Menyisipkan data ke dalam tabel menurut nama kolom.
ONCEJika opsi tidak disediakan, kueri harus berupa kueri streaming. Gunakan kata kunciSTREAMuntuk memanfaatkan semantik streaming dalam membaca dari sumber. 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.Nota
FLOW INSERT BY NAMEsetara dengan menggunakanAS query. Dua pernyataan berikut memiliki perilaku yang identik: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');SEKALI
Secara opsional mendefinisikan alur sebagai alur satu kali, seperti isi ulang. Ketika
ONCEdisediakan, kueri bukan kueri streaming, dan alur berjalan satu kali secara default. Jika tabel di-refresh dengan refresh penuh,ONCEalur berjalan lagi untuk membuat ulang data.ONCEhanya berlaku untukINSERT BY NAMEalur.AUTO CDCPenting
Tersedia di Databricks Runtime 17.3 ke atas dan
PREVIEWsaluran Alur.AUTO CDCMenentukan alur yang memproses mengubah rekaman pengambilan data (CDC) dari sumber ke dalam tabel. GunakanAUTO CDCsaat data sumber menyertakan semantik CDC. Lihat API CDC Otomatis: Menyederhanakan penangkapan perubahan data menggunakan pipeline.GANTI WHEREpredikat menurut kueri NAMA
REPLACE WHEREMenentukan alur yang mengolah ulang dan hanya menimpa baris yang cocokpredicate, membiarkan semua baris lain tidak tersentuh. GunakanREPLACE WHEREuntuk pemrosesan batch inkremental gabungan dan agregasi, data yang terlambat tiba, evolusi skema, dan backfills.BY NAMEdiperlukan. Lihat Pemrosesan batch dengan alur REPLACEWHERE.GANTI MENGGUNAKAN ( column_name [, ...] ) Kueri URUTAN BERDASARKAN sequence_column BERDASARKAN NAMA
Penting
Fitur ini ada di Beta. Memerlukan Databricks Runtime 18.2 ke atas.
Mendefinisikan
REPLACE USINGalur yang menggantikan semua baris yang sesuai dengan kolom kunci yang ditentukan dan membiarkan semua baris lainnya tidak tersentuh. 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. Sumber harus berupa sumber streaming.BY NAMEdiperlukan. Lihat Penggantian snapshot parsial dengan alur GANTI MENGGUNAKAN.
-
Klausa ini mengisi tabel menggunakan data dari
query. Kueri ini 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 menyerap data yang memiliki penerapan perubahan, Anda dapat menambahkanskipChangeCommitsopsi baca untuk menangani kesalahan.Saat Anda menentukan
querydantable_specificationbersama-sama, skema tabel yang ditentukan dalamtable_specificationharus berisi semua kolom yang dikembalikan olehquery, jika tidak, Anda akan mendapatkan kesalahan. Kolom apa pun yang ditentukan dalamtable_specificationtetapi tidak dikembalikan olehquerymengembalikan nilainullsaat dikueri.Untuk informasi selengkapnya tentang data streaming, lihat Mengubah data dengan alur.
Opsi Baca
Anda dapat menentukan opsi baca dalam kueri untuk mengonfigurasi cara data dibaca dari sumber. Misalnya, Anda dapat menentukan
skipChangeCommitsuntuk melewati penerapan perubahan apa pun dalam data sumber. Opsi baca ditentukan sebagai peta dalamWITHklausa kueri. Contohnya:SELECT * FROM STREAM source_table WITH (SKIPCHANGECOMMITS=TRUE, STARTINGVERSION=X)=TRUEbersifat opsional, sehingga Anda juga dapat menentukan opsi boolean seperti ini:SELECT * FROM STREAM source_table WITH (SKIPCHANGECOMMITS)Nota
Opsi baca hanya didukung untuk Databricks Runtime 17.3 ke atas.
Opsi baca di bawah ini didukung untuk Delta, untuk detail tentang setiap opsi, lihat Baca dan tulis streaming tabel Delta Lake.
maxFilesPerTriggermaxBytesPerTriggerstartingVersionstartingTimestampreadChangeFeedwithEventTimeOrderskipChangeCommits
Memerlukan izin
Pengguna 'run-as' untuk pipeline harus memiliki izin berikut:
-
SELECThak akses atas tabel dasar yang direferensikan oleh tabel streaming. -
USE CATALOGhak istimewa pada katalog induk dan hak istimewaUSE SCHEMApada skema induk. -
CREATE MATERIALIZED VIEWhak akses pada skema untuk tabel streaming.
Agar pengguna dapat memperbarui alur tempat tabel streaming ditentukan, mereka memerlukan:
-
USE CATALOGhak istimewa pada katalog induk dan hak istimewaUSE SCHEMApada skema induk. - Kepemilikan atau hak istimewa pada tabel streaming
REFRESH. - Pemilik tabel streaming harus memiliki
SELECThak istimewa atas tabel dasar yang dirujuk oleh tabel streaming.
Agar pengguna dapat mengkueri tabel streaming yang dihasilkan, mereka memerlukan:
-
USE CATALOGhak istimewa pada katalog induk dan hak istimewaUSE SCHEMApada skema induk. -
SELECThak akses istimewa atas tabel streaming.
Keterbatasan
- Hanya pemilik tabel yang dapat memperbarui tabel streaming untuk mendapatkan data terbaru.
-
ALTER TABLEperintah tidak diperbolehkan di tabel streaming. Definisi dan properti tabel harus diubah melalui pernyataanCREATE OR REFRESHatau ALTER STREAMING TABLE. - Mengembangkan skema tabel melalui perintah DML seperti
INSERT INTO, danMERGEtidak didukung. - Perintah berikut ini tidak didukung pada tabel streaming:
CREATE TABLE ... CLONE <streaming_table>COPY INTOANALYZE TABLERESTORETRUNCATEGENERATE MANIFEST[CREATE OR] REPLACE TABLE
- Mengganti nama tabel atau mengubah pemilik tidak didukung.
Examples
-- 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);