Not
Bu sayfaya erişim yetkilendirme gerektiriyor. Oturum açmayı veya dizinleri değiştirmeyi deneyebilirsiniz.
Bu sayfaya erişim yetkilendirme gerektiriyor. Dizinleri değiştirmeyi deneyebilirsiniz.
Important
Bu özellik Beta sürümündedir.
Dekoratör, @dp.replace_flow boru hattınızda bir akış tablosu için REPLACE USING akışı oluşturur. Her güncellemede, akış hedef tablodaki replace_using anahtar sütunlarla eşleşen tüm satırlar yerine geçer ve diğer tüm satırlar dokunulmaz. İşlev bir Apache Spark akış veri çerçevesi döndürmelidir.
Bkz. Kısmi anlık değişimi ile REPLACE USING akışları.
Kaynağınız sütunlara göre anahtarlanmış kısmi anlık görüntüler dizisi olduğunda kullanın @dp.replace_flow . Hedef tabloyu ve akışı tek bir ifadede tanımlamak için ve sequence_by@dp.table'a geçinreplace_using.
Syntax
from pyspark import pipelines as dp
dp.create_streaming_table("<target-table-name>") # Required only if the target table doesn't exist.
@dp.replace_flow(
target = "<target-table-name>",
replace_using = ["<key-column>", "<key-column>"],
sequence_by = "<sequence-column>",
name = "<flow-name>", # optional, defaults to function name
comment = "<comment>", # optional
spark_conf = {"<key>" : "<value>", "<key>" : "<value>"}) # optional
def <function-name>():
return (<streaming-query>)
Parametre
| Parametre | Türü | Description |
|---|---|---|
| function | function |
Gerekli. Kullanıcı tanımlı bir sorgudan Apache Spark akış DataFrame'i döndüren işlev. |
target |
str |
Gerekli. Akışın hedefi olan akış tablosunun adı. |
replace_using |
list |
Gerekli. Hangi hedef satırların değiştirileceğini belirleyen anahtar sütunlar. En az bir sütun belirtin. Ana sütunlar tekrarlanamaz ve her anahtar sütunun türü sıralanabilir olmalıdır. |
sequence_by |
str veya Column |
Gerekli. Güncellemeleri emreten sütun. Her tuş için en yüksek dizis kazanır ve alt sıralı bir sıra hedefte zaten bulunan daha yüksek bir sıranın üzerine yazmaz. |
name |
str |
Akış adı. Sağlanmadıysa, varsayılan olarak işlev adını kullanır. |
comment |
str |
Akış açıklaması. |
spark_conf |
dict |
Bu sorgunun yürütülmesi için Spark yapılandırmalarının listesi. |
Examples
from pyspark import pipelines as dp
# Keep the latest row for each order from a stream of partial snapshots
dp.create_streaming_table("orders_current")
@dp.replace_flow(
target = "orders_current",
replace_using = ["order_id"],
sequence_by = "updated_at"
)
def orders_flow():
return spark.readStream.table("order_updates")
Bir kayıt sütun kombinasyonuyla tanımlandığında birden fazla anahtar sütunu kullanın:
from pyspark import pipelines as dp
dp.create_streaming_table("accounts_current")
@dp.replace_flow(
target = "accounts_current",
replace_using = ["region", "account_id"],
sequence_by = "updated_at"
)
def accounts_flow():
return spark.readStream.table("account_updates")
Sınırlamalar
- Bir akış tablosu tek
REPLACE USINGbir akışı destekler ve ek akış, otomatik CDC akışı veyaREPLACE WHEREakış gibi başka bir akış türüyle birleşemezREPLACE USING. - Sorgu bir akış sorgusu olmalıdır.
@dp.replace_flowyayın dışı bir kaynağı reddeder. -
REPLACE USINGakışlar için Databricks Runtime 18.2 ve üzeri gerektirir.