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.
Dekoratör, @dp.append_flow işlem hattı tablolarınız için ekleme akışları veya geri doldurmalar oluşturur. İşlev bir Apache Spark akış veri çerçevesi döndürmelidir. Bkz . Lakeflow işlem hattı akışlarıyla verileri artımlı olarak yükleme ve işleme.
Ekleme akışları akış tablolarını veya havuzlarını hedefleyebilir.
Sözdizimi
from pyspark import pipelines as dp
dp.create_streaming_table("<target-table-name>") # Required only if the target table doesn't exist.
@dp.append_flow(
target = "<target-table-name>",
name = "<flow-name>", # optional, defaults to function name
once = False, # optional
spark_conf = {"<key>" : "<value", "<key" : "<value>"}, # optional
comment = "<comment>", # optional
import_checkpoint = "<checkpoint-path>") # optional
def <function-name>():
return (<streaming-query>) #
Parametreler
| 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. Ekleme akışının hedefi olan tablo veya havuzun adı. |
name |
str |
Akış adı. Sağlanmadıysa, varsayılan olarak işlev adını kullanır. |
once |
bool |
İsteğe bağlı olarak, akışı yedek doldurma gibi tek seferlik bir akış olarak tanımlayın. Kullanımı once=True , akışı iki şekilde değiştirir:
|
comment |
str |
Akış açıklaması. |
spark_conf |
dict |
Bu sorgunun yürütülmesi için Spark yapılandırmalarının listesi |
import_checkpoint |
str |
Mevcut bir Yapılandırılmış Akış kontrol noktasına giden yol, böylece taşınan akış, kaynağı yeniden işlemek yerine son taahhüt edilen ofsetinden devam eder. Bir kontrol noktası içe aktarmak Beta'da. Bkz. Yapılandırılmış Akış Kontrol Noktasını Taşımak. |
Örnekler
from pyspark import pipelines as dp
# Create a sink for an external Delta table
dp.create_sink("my_sink", "delta", {"path": "/tmp/delta_sink"})
# Add an append flow to an external Delta table
@dp.append_flow(name = "flow", target = "my_sink")
def flowFunc():
return <streaming-query>
# Add a backfill
@dp.append_flow(name = "backfill", target = "my_sink", once = True)
def backfillFlowFunc():
return (
spark.read
.format("json")
.load("/path/to/backfill/")
)
# Create a Kafka sink
dp.create_sink(
"my_kafka_sink",
"kafka",
{
"kafka.bootstrap.servers": "host:port",
"topic": "my_topic"
}
)
# Add an append flow to a Kafka sink
@dp.append_flow(name = "flow", target = "my_kafka_sink")
def myFlow():
return read_stream("xxx").select(F.to_json(F.struct("*")).alias("value"))
Yapılandırılmış Akış kontrol noktasını Göç Et
Important
Bir kontrol noktası içe aktarmak Beta'da.
Mevcut bir Yapılandırılmış Akış iş yükünü kaynağı yeniden işlemeden bir boru hattına taşımak için kullanılır import_checkpoint . Bunu checkpointLocation , bulut depolama, Unity Kataloğu hacmi veya DBFS yolu olan Structured Streaming sorgusunuza ayarlayın. İlk boru hattı güncellemesinde, akış o kontrol noktasını pipeline'ın yönetilen depolamasına klonlar. Akış, son bağlanmış ofsetten (örneğin toplamalar, deduplication anahtarları ve watermarklar gibi) olduğu gibi devam eder. Sonraki boru hattı güncellemeleri akışın klonlanmış kontrol noktasını kullanır; orijinal kontrol noktası değiştirilmez.
Akış, create_table veya bir sink ile oluşturulan yönetilen bir tabloyu hedeflemelidir.
Pipeline'ı çalıştırmadan önce orijinal Structured Streaming sorgusunu durdurun. Orijinal Yapılandırılmış Akış sorgusu aktarmadan sonra tekrar kullanılabilir, ancak kontrol noktası durumunu yönetmeniz ve boru hattı ile sorgu aynı anda aynı tabloya yazmadığından emin olmalısınız; bu da tekrarlanan veri üretebilir.
Structured Streaming sorgusunu yeni bir tabloya yazan ve kontrol noktasını içe aktaran bir boru hattı akışı olarak yeniden oluşturun:
from pyspark import pipelines as dp
# Create a new managed table for the pipeline
dp.create_table("target_table")
# Continue from the imported checkpoint instead of reprocessing the source.
@dp.append_flow(
target = "target_table",
import_checkpoint = "/Volumes/my_catalog/my_schema/checkpoints/my_stream",
)
def migrate():
# The same source your original query read from.
return spark.readStream.table("source_table")
Kontrol noktası yalnızca bir kez, ilk boru hattı güncellemesinde içe aktarılır; sonraki güncellemeler ise 'yi görmezden import_checkpointgelir.
Tam yenileme kontrol noktasını yeniden içe aktarmaz; yeni, boş bir kontrol noktasından başlar ve kaynağı yeniden işler. Farklı bir kontrol noktası içe aktarmak için, hedef tablo için daha önce kullanılmamış bir akış adı kullanın; mevcut bir akış adının tekrar kullanılması ithalatı atlar.
Sınırlamalar
- Zaten var olan bir tabloya (örneğin, orijinal Yapılandırılmış Akış sorgu hedefi) bir kontrol noktası aktarmak desteklenmez. Pipeline'ın oluşturduğu yeni bir tabloyu veya bir sink'i hedefleyin.
-
import_checkpointyalnızca append_flow desteklenir.