append_flow

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:
  • Geri dönüş değeri. streaming-query. Bu durumda, bir akış DataFrame değil, bir yığın DataFrame olmalıdır.
  • 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.
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_checkpoint yalnızca append_flow desteklenir.