transformWithState ile özel bir durum bilgisi olan uygulama oluşturun

transformWithState ile durum bilgisine sahip akış uygulamaları oluşturabilir ve düşük gecikmeli, gerçek zamana yakın çözümler geliştirebilirsiniz. Özel durum bilgisi olan işleçlerle, geleneksel Yapılandırılmış Akış işleme ile mümkün olmayan yeni işletimsel kullanım örnekleri oluşturmanıza olanak tanıyan rastgele durum bilgisi olan mantık oluşturabilirsiniz.

Not

Toplamalar, yinelenenleri kaldırma ve akış birleştirmeleri gibi durum bilgisi olan işlemler için Databricks, özel mantık yerine yerleşik Yapılandırılmış Akış işleçlerinin kullanılmasını önerir. Bkz. Durumlu akış nedir?.

Databricks, isteğe bağlı durum dönüşümleri için, transformWithState ve flatMapGroupsWithState gibi eski işleçler yerine mapGroupsWithState kullanılmasını önerir. Bkz. Eski rastgele durum bilgisi olan işleçler.

Gereksinimler

transformWithState ve transformWithStateInPandas işleçleri aşağıdaki gereksinimlere sahiptir:

  • Databricks Runtime 16.2 ve üzerinde kullanılabilir.
    • Gerçek zamanlı mod için Databricks Runtime 17.3 LTS veya üzerini kullanın. Gerçek zamanlı mod kavramlarına bakınız.
    • Standart erişim modu için Python Databricks Runtime 16.3 ve üzeri sürümleri ile Scala, Databricks Runtime 17.3 ve üzeri sürümleriyle kullanılabilir.
  • RocksDB, Databricks Runtime 17.3 ve üzeri sürümlerin varsayılan durum deposu sağlayıcısıdır.
    • Databricks Runtime 17.2 ve altı için RocksDB durum deposu sağlayıcısını yapılandırmanız gerekir. Databricks, Spark yapılandırmasında RocksDB'nin etkinleştirilmesini önerir.

      spark.conf.set("spark.sql.streaming.stateStore.providerClass", "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")
      

transformWithState nedir?

transformWithState operatörü, Yapılandırılmış Akış sorgusuna özel bir durum bilgili işlemci uygular. transformWithStatekullanmak için durum bilgisi olan özel bir işlemci uygulamanız gerekir. Yapılandırılmış Akış, Python, Scala veya Java kullanarak durum bilgisi olan işlemcinizi oluşturmaya yönelik API'leri içerir.

Gruplandırma anahtarına özel mantık uygulamak için kullanın transformWithState . Aşağıda üst düzey tasarım açıklanmaktadır:

  • Bir veya daha fazla durum değişkeni tanımlayın.
  • Durum bilgileri her gruplandırma anahtarı için kalıcıdır. Kullanıcı tanımlı kodda her durum değişkenine erişebilirsiniz.
  • İşlenen her mikro grup için, anahtara ait tüm satırlar bir yineleyici olarak kullanılabilir.
  • Satırların nasıl yayımlanacağını denetlemek için StatefulProcessorHandle öğesini zamanlayıcılar ve kullanıcı tanımlı koşullarla birlikte kullanın.
  • Durum süre sonu ve durum boyutunu yönetmek için, durum değerleri tek tek yaşam süresi (TTL) tanımlarını destekler.

transformWithState Durum deposunda şema evrimini desteklediğinden, geçmiş durum bilgilerini kaybetmeden üretim uygulamalarınızı yineleyebilir ve güncelleştirebilirsiniz. Durum şemasını güncelleştirdikten sonra, kod dağıtımlarını ve bakımını basitleştiren satırları yeniden işlemeniz gerekmez. Bkz. durum deposunda şema evrimi.

Önemli

Azure Databricks belgelerinde, hem Python hem de Scala uygulamalarını açıklamak için transformWithState kullanılır:

  • PySpark hem satır tabanlı transformWithState API'yi hem de Pandas tabanlı transformWithStateInPandas işleci destekler.
    • transformWithStateInPandas gerçek zamanlı modda desteklenmez. Bunun yerine transformWithStatekullanın. Ayrıntılar için bkz. transformWithState gerçek zamanlı modda.
    • Sıra tabanlı transformWithState API, daha yüksek verimlilik için asenkron işlemeyi asyncio destekler. Sunucusuz hesaplamada asenkron işlem desteklenmez. Bkz. Asenkron İşleme (Beta).
  • Scala yalnızca satır tabanlı transformWithState API'yi destekler.

Scala ve Python uygulamaları transformWithState aynı özelliklere sahiptir, ancak söz diziminde bazı farklılıklar vardır.

Bir StatefulProcessor Tanımlama

StatefulProcessor sınıfını genişleterek ve yöntemlerini gerçekleştirerek durum bilgisi içeren bir işlemci tanımlarsınız.

Spark, StatefulProcessorHandle öğenizin init yöntemine bir StatefulProcessor iletir. Durum değişkenleri oluşturmak ve durum deposuyla etkileşim kurmak için tutamacı kullanın.

transformWithState üç durum türünü destekler: ValueState, ListStateve MapState. Her tür, her gruplandırma anahtarının durumunu farklı bir temel veri yapısı kullanarak depolar.

Özel mantığınızı tanımlamak için aşağıdaki yöntemleri uygulayın:

  • Uygulamanızın verileri nasıl işlediğini, durumu nasıl güncellediğini ve her bir mikro toplu işlem için satırları nasıl ürettiğini denetlemek üzere handleInputRows öğesini uygulayın. Bkz. giriş satırlarını işleme.
  • Gruplandırma anahtarının bir mikro toplu işlemde yeni satırlar alıp almadığına bakılmaksızın zaman tabanlı mantığı çalıştırmak için handleExpiredTimer uygulayın. Bkz . Süresi dolan zamanlayıcıları işleme.
  • İsteğe bağlı olarak, uygulamanız herhangi bir giriş satırını işlemeden önce durumu önceden doldurmak için handleInitialState öğesini uygulayın. Bkz. Başlangıç durumunu işleme.

Aşağıdaki tablo, bu yöntemlerin işlevsel davranışlarını karşılaştırır:

Davranış handleInputRows handleExpiredTimer
Durum değerlerini alma, yerleştirme, güncelleştirme veya temizleme Evet Evet
Zamanlayıcı oluşturma veya silme Evet Evet
Satırları çıktı olarak verme Evet Evet
Geçerli mikro gruptaki satırları yinele Evet Hayır
Geçen süreyi temel alan tetikleyici mantığı Hayır Evet

Karmaşık mantığı gerektiği gibi uygulamak için hem hem handleInputRows de handleExpiredTimer öğesini birleştirebilirsiniz.

Örneğin, her mikro toplu iş için durum değerlerini güncelleştirmek için handleInputRows kullanan ve gelecekte 10 saniyelik bir süreölçer ayarlayan bir uygulama uygulayabilirsiniz. Ek satır işlenmezse, durum deposundaki geçerli değerleri yayımlamak için handleExpiredTimer kullanabilirsiniz. Gruplandırma anahtarı için yeni satırlar işlenirse, mevcut zamanlayıcıyı temizleyebilir ve yeni bir zamanlayıcı ayarlayabilirsiniz.

StatefulProcessorHandle

Sınıfı, PySpark'ta StatefulProcessorHandle kodunuzun durum bilgilerini nasıl kullandığını denetleye işlevlere erişmenizi sağlar.

Bir StatefulProcessor başlatırken, StatefulProcessorHandle öğesini her zaman içe aktarmalı ve bunu handle değişkenine geçirmelisiniz. handle değişkeni, Python sınıfınızdaki yerel değişkeni durum değişkeniyle bağlar.

Not

Scala getHandle yöntemini kullanır.

Özel durum türleri

Tek bir durum bilgisi olan işleçte birden çok durum nesnesi uygulayabilirsiniz.

Tam uygulama mantığınıza göre bir durum türü seçin. Örneğin, ValueState ve user_id'ye göre gruplandırılmış oturumları session_id ile izleyebilirsiniz. Veya, birden fazla oturumdaki koşulları değerlendirmek için, MapState öğesini eşleme anahtarı olarak kullanarak user_id ölçütüne göre gruplanmış bir session_id kullanın.

Durum nesneniz bir StructTypekullanıyorsa, şema yapısındaki her alan için benzersiz adlar tanımlamanız gerekir. Bu adlar, state store okunurken görünür. Bkz. Yapılandırılmış Akış durumu bilgilerini okuyun.

Aşağıdaki bölümlerde tarafından transformWithStatedesteklenen durum türleri açıklanmaktadır:

ValueState

ValueState her gruplandırma anahtarı için bir değer depolar.

Değer durumu, yapı veya demet gibi karmaşık türler içerebilir. için ValueState, değerin tamamını değiştirmek için mantık uygulamanız gerekir.

Değer güncelleştirildiğinde değer durumu için yaşam süresi sıfırlanır. ValueState için bir kaynak anahtarı işleyip depolanan ValueState öğesini güncelleştirmezseniz yaşam süresi sıfırlanmaz.

ListState

ListState her gruplandırma anahtarı için bir liste depolar.

Liste durumu, her biri karmaşık türler içerebilen bir değer koleksiyonudur. Listedeki her değerin kendi yaşam süresi vardır.

Tek tek öğeleri ekleyerek, bir öğe listesi ekleyerek veya putile listenin tamamının üzerine yazarak listeye öğe ekleyebilirsiniz. Yaşam süresini sıfırlamak için bir put işlem kullanmanız gerekir.

MapState

MapState her gruplandırma anahtarı için bir harita depolar. Haritalar, bir Python sözlüğüne (dict ) eşdeğer apache Spark'tır.

Eşleme durumu, her biri bir değere karşılık gelen ayrı anahtarlardan oluşan bir koleksiyondur; bu değerlerin her biri karmaşık türler içerebilir. Bir haritadaki her anahtar-değer çifti kendi yaşam süresine sahiptir.

Belirli bir anahtarın değerini güncelleştirebilir veya anahtarı ve değerini kaldırabilirsiniz. Anahtarını kullanarak tek bir değer döndürebilir, tüm anahtarları listeleyebilir, tüm değerleri listeleyebilir veya eşlemedeki tam anahtar-değer çiftleri kümesiyle çalışmak için bir yineleyici döndürebilirsiniz.

Önemli

Gruplandırma anahtarları, Yapılandırılmış Akış sorgusunun GROUP BY yan tümcesinde belirtilen alanları açıklar. Eşleme durumları, bir gruplandırma anahtarı için herhangi bir sayıda anahtar-değer çifti içerebilir.

Örneğin, sorgunuzda GROUP BY user_id kullanılıyorsa ve her session_id için bir eşleme tanımlamak istiyorsanız, gruplandırma anahtarınız user_id olur ve MapState anahtarı session_id olur:

Python
class SessionTracker(StatefulProcessor):
  def init(self, handle: StatefulProcessorHandle) -> None:
    self.sessions = handle.getMapState("sessions", "session_id string", "count long")

  def handleInputRows(self, key, rows: Iterator[Row], timerValues) -> Iterator[Row]:
    for row in rows:
      session_key = (row["session_id"],)  # session_id is the MapState key
      count = self.sessions.getValue(session_key)[0] if self.sessions.containsKey(session_key) else 0
      new_count = count + 1
      self.sessions.updateValue(session_key, (new_count,))
    yield from []

  def close(self) -> None:
    pass

df.groupBy("user_id").transformWithState(SessionTracker(), ...) # user_id is the grouping key
Scala programlama dili
case class Event(userId: String, sessionId: String)

class SessionTracker extends StatefulProcessor[String, Event, (String, Long)] {
  @transient private var sessions: MapState[String, Long] = _

  override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
    sessions = getHandle.getMapState[String, Long]("sessions", Encoders.STRING, Encoders.scalaLong, TTLConfig.NONE)
  }

  override def handleInputRows(
      key: String,
      rows: Iterator[Event],
      timerValues: TimerValues): Iterator[(String, Long)] = {
    rows.foreach { event =>
      val count = if (sessions.containsKey(event.sessionId)) sessions.getValue(event.sessionId) else 0L
      sessions.updateValue(event.sessionId, count + 1) // sessionId is the MapState key
    }
    Iterator.empty
  }
}

df.as[Event]
  .groupByKey(_.userId) // userId is the grouping key
  .transformWithState(new SessionTracker(), TimeMode.None(), OutputMode.Update())

içinde özel durum değişkeni oluşturma StatefulProcessor

StatefulProcessor'nizi başlatırken, her durum nesnesi için özel mantığınızdaki durum nesneleriyle etkileşim kurmanızı sağlayan bir yerel değişken oluşturursunuz. init sınıfındaki yerleşik StatefulProcessor yöntemini geçersiz kılarak durum değişkenlerini tanımlayın ve başlatın.

getValueState içinde, getListState, getMapState ve StatefulProcessor yöntemlerini kullanarak istediğiniz sayıda durum nesnesi tanımlayabilirsiniz.

Her durum nesnesi aşağıdakilere sahip olmalıdır:

  • Benzersiz bir ad
  • Şema
    • Python'da şemayı belirtmeniz gerekir.
    • Scala'da durum şemasını belirtmek için bir Encoder geçirebilirsiniz.

İsteğe bağlı olarak, yaşam süresi (TTL) süresini milisaniye cinsinden de sağlayabilirsiniz. Eşleme durumu uyguluyorsanız, eşleme anahtarları ve değerler için ayrı bir şema tanımı sağlamanız gerekir.

Not

, StatefulProcessor durum bilgilerini sorgulama, güncelleştirme ve yayma mantığını ayrı ayrı işler. Bkz. Özel mantık içeren yöntemlerde durum değişkenlerinizi kullanma.

Özel mantıkla yöntemlerde durum değişkenlerinizi kullanma

Durum nesnelerinin durum alma, mevcut durum bilgilerini güncelleştirme ve geçerli durumu temizleme yöntemleri vardır.

Her gruplandırma anahtarının ayrılmış durum bilgileri vardır.

  • StatefulProcessor, özel mantığınıza ve belirtilen çıktı şemasına göre satırlar üretir. Bkz. Satırları yayımla.
  • Durum deposundaki statestore değerlere erişmek için okuyucuyu kullanın. Bu okuyucu toplu iş yüklerine yöneliktir ve düşük gecikme süreli iş yüklerine yönelik değildir. Bkz. Yapılandırılmış Akış durumu bilgilerini okuyun.
  • Kullanılarak handleInputRows belirtilen mantık yalnızca anahtarın satırları bir mikro toplu işlemde mevcutsa çalışır. Bkz. giriş satırlarını işleme.
  • Satırları gözlemlemeye bağlı olmadan tetiklenen zamana dayalı mantığı uygulamak için handleExpiredTimer kullanın. Bkz . Süresi dolan zamanlayıcıları işleme.

Not

Durum nesneleri, anahtarları gruplandırarak yalıtılır ve aşağıdaki etkilere neden olabilir:

  • Durum değerleri, farklı bir gruplandırma anahtarıyla ilişkili satırlardan etkilenmez.
  • Değerleri karşılaştırmaya veya gruplandırma anahtarları arasındaki durumu güncellemeye bağımlı olan bir mantığı uygulayamazsınız.

Gruplandırma anahtarındaki değerleri karşılaştırabilirsiniz. Özel mantığınızın kullanabileceği ikinci bir anahtarla mantık uygulamak için bir MapState kullanın. Örneğin, user_id göre gruplandırmak ve ip_address anahtarınız için MapState kullanmak, eşzamanlı kullanıcı oturumlarını izlemenize olanak tanır.

Durumla çalışırken dikkat edilmesi gereken gelişmiş hususlar

Durum bilgisi güncellemeleri hataya dayanıklıdır. Bir görev, bir mikro toplu işlemin işlenmesi tamamlanmadan önce çökerse, yeniden denemede son başarılı mikro toplu işlemdeki değer kullanılır.

En iyi duruma getirilmiş performans için Databricks, belirli bir anahtar için yineleyicideki tüm değerleri işlemenizi ve güncelleştirmeleri tek bir yazmada işlemenizi önerir. Bir durum değişkenine yazdığınızda bu, RocksDB'ye yazma işlemi tetikler.

Durum değerlerinin varsayılan değerleri yoktur. Mantığınız mevcut durum bilgilerinin okunmasını gerektiriyorsa yöntemini kullanın exists .

Null durumu mantığını gerçekleştirmek için MapState değişkenleri, tek tek anahtarları denetlemenize veya tüm anahtarları listelemenize olanak tanır.

Giriş satırlarını işleme

Uygulamanızın handleInputRows satır ve güncelleştirme durum değerlerini nasıl işlediğini tanımlamak için yöntemini kullanın. Bu yöntem, Structured Streaming sorgunuz bir gruplama anahtarına ait satırları her işlediğinde çalıştırılır.

transformWithStateile uygulanan durum bilgisi olan uygulamaların çoğunda, çekirdek mantık handleInputRowskullanılarak tanımlanır.

İşlenen her mikro toplu iş güncelleştirmesi için, belirli bir gruplandırma anahtarı için mikro toplu işteki tüm satırlar yineleyici kullanılarak kullanılabilir. Kullanıcı tanımlı mantık, geçerli mikro toplu işteki tüm satırlar ve statestore'daki değerlerle etkileşime geçebilir.

Süresi dolan zamanlayıcıları işleme

Geçen süreye göre özel mantık gerçekleştirmek için handleExpiredTimer yöntemini kullanın.

Bir gruplandırma anahtarında zamanlayıcılar, zaman damgaları tarafından benzersiz olarak tanımlanır.

Süreölçerin süresi dolduğunda, sonuç uygulamanızda uygulanan mantık tarafından belirlenir. Yaygın desenler şunlardır:

  • Durum değişkeninde depolanan bilgileri yayma.
  • Depolanan durum bilgilerini silme.
  • Yeni zamanlayıcı oluşturuluyor.

Süresi dolan zamanlayıcılar, bir mikro toplu işlemde ilişkili anahtarları için hiçbir satır işlenmemiş olsa bile tetiklenir.

Zaman modunu belirtme

StatefulProcessor öğesini transformWithState öğesine geçirirken, zaman modunu timeMode parametresini kullanarak belirtmeniz gerekir.

Aşağıdaki seçenekler desteklenir:

Zaman modu Açıklama
ProcessingTime Zamanlayıcılar ve TTL desteklenir ve Apache Spark her mikro toplu işi işlediğinde duvar saati zamanına göre değerlendirilir. Verilerdeki zaman damgalarına bakılmaksızın, zamanlayıcıların satırların işlendiği zamanlara göre sabit bir aralıkta tetiklenmelerini istediğinizde kullanın ProcessingTime .
EventTime Zamanlayıcılar desteklenmektedir ve event-time watermark’ına göre değerlendirilmektedir. Apache Spark giriş verilerindeki zaman damgalarını gözlemledikçe filigran ilerler. TTL, EventTime ile desteklenmez. Verileriniz zaman damgaları içerdiğinde ve zamanlayıcıların bu zaman damgalarının ilerlemesine göre tetiklenmelerini istediğinizde kullanın EventTime . EventTime kullanırken, eventTimeColumnName parametresini de belirtmeniz gerekir. Bkz. eventTimeColumnName.
NoTime veya TimeMode.None() Zamanlayıcılar ve TTL desteklenmez. Durum bilgisi olan uygulamanız zamana dayalı mantık gerektirmediğinde kullanın NoTime .

eventTimeColumnName

Zaman modunu kullanırken EventTime parametresi, eventTimeColumnName çıkış şemanızda olay zaman damgasını içeren sütunun adını belirtir. Apache Spark, watermark’ı çıkış akışına iletmek için bu sütunu kullanır; böylece alt akıştaki zaman tabanlı işlemlerin doğru şekilde gerçekleştirilmesini sağlar.

Python

eventTimeColumnNameveya transformWithStateiçin transformWithStateInPandas ek bir bağımsız değişkendir:

q = (
  df.groupBy("key")
    .transformWithState(
      statefulProcessor=MyProcessor(),
      outputStructType=output_schema,
      outputMode="Append",
      timeMode="EventTime",
      eventTimeColumnName="outputTimestamp",
    )
    .writeStream...
)
Scala programlama dili

transformWithState, eventTimeColumnName yerine timeMode kabul eder. Bu yaklaşım her zaman EventTime modunu kullanır:

val q = spark
  .readStream
  .format("delta")
  .load(srcDeltaTableDir)
  .as[(String, String)]
  .groupByKey(x => x._1)
  .transformWithState(
    new MyProcessor(),
    "outputTimestamp",
    OutputMode.Append(),
  )
  .writeStream...

Yerleşik zamanlayıcı değerleri

Databricks, durum bilgisi olan özel uygulamanızda sistem saatinin çağrılmasını önermemektedir, çünkü bu, görev hatası durumunda güvenilir olmayan yeniden denemelere yol açabilir. İşleme süresine veya filigrana erişmeniz gerektiğinde TimerValues sınıfındaki yöntemleri kullanın:

TimerValues Açıklama
getCurrentProcessingTimeInMs İşlem süresinin zaman damgasını mevcut toplu işlem için Unix epoch'undan beri milisaniye cinsinden döndürür.
getCurrentWatermarkInMs Geçerli toplu iş için filigranın zaman damgasını, Unix dönemi başlangıcından itibaren milisaniye cinsinden geri döndürüyor.

Not

İşlem süresi, mikro toplu işlemin Apache Spark tarafından işlenme süresini açıklar. Kafka gibi birçok akış kaynağı sistem işleme süresini de içerir.

Akış sorgularında filigranlar genellikle olay zamanına veya akış kaynağının işleme süresine göre tanımlanır. veri işleme eşiklerini denetlemek için filigranları uygulama bölümüne bakın.

transformWithStateile hem filigranlar hem de pencereler birlikte kullanılabilir. TTL, zamanlayıcılar ve MapState veya ListState işlevselliğinden yararlanarak özel durum bilgisi olan uygulamanızda benzer işlevler uygulayabilirsiniz.

Durum türleri için yaşam süresi (TTL)

Yetersiz bellek hatalarını önlemek ve eski durum türü değerlerini kaldırmak için, transformWithState her durum türü değeri için isteğe bağlı yaşam süresi (TTL) değerini destekler. Süre dolduktan sonra TTL, durum türü değerlerini sessizce çıkartır. TTL, handleExpiredTimer veya herhangi bir özel mantığı çalıştırmaz. Durum süresi dolduğunda kodu çalıştırmak için bunun yerine bir zamanlayıcı kullanın.

Önemli

TTL kullanmazsanız, yetersiz bellek hatalarını önlemek için durumun bellekten atılmasını yönetmeniz gerekir.

Tüm durum türleri için, durum bilgileri güncelleştirilirken TTL sıfırlanır. TTL, her durum türü değeri için uygulanır ve her durum türü için farklı kurallara göre işler:

  • Durum değişkenleri, gruplandırma anahtarları ile sınırlıdır.
  • ValueState nesneler için gruplandırma anahtarı başına yalnızca tek bir değer depolanır. TTL bu değer için geçerlidir.
  • ListState nesneler için liste birçok değer içerebilir. TTL, listedeki her değere bağımsız olarak uygulanır.
    • TTL, bir ListState içindeki tek tek değerlere uygulanıyor olsa da, tek bir değeri güncellemenin tek yolu put yöntemidir; bu yöntem ListState değişkeninin tüm içeriğinin üzerine yazar ve listedeki tüm değerler için TTL'yi sıfırlar.
  • MapState nesneler için her eşleme anahtarının ilişkili bir durum değeri vardır. TTL, bir eşlemedeki her anahtar-değer çiftine bağımsız olarak uygulanır.

Not

Zamanlayıcılar, satır yayımlama da dahil olmak üzere, durumun bellekten atılmasının ötesinde özel mantık tanımlamanıza olanak sağlar. İsteğe bağlı olarak, belirli bir durum değerinin durum bilgilerini temizlemek ve değerleri yaymak veya koşullu mantığı tetiklemek için zamanlayıcıları kullanabilirsiniz. Bkz . Süresi dolan zamanlayıcıları işleme.

Durumlu örnek uygulama

Aşağıdaki örnek, örnek durum değişkenleri dahil olmak üzere özel, durum bilgisi tutan bir işleyici tanımlar: SimpleCounterProcessor SimpleCounterProcessor, her gruplandırma anahtarı için satırları saymak üzere ValueState, ListState ve MapState kullanır.

Python (Pandas)

import pandas as pd
from pyspark.sql import Row
from pyspark.sql.streaming import StatefulProcessor, StatefulProcessorHandle
from pyspark.sql.types import StructType, StructField, IntegerType, StringType
from typing import Iterator

spark.conf.set("spark.sql.streaming.stateStore.providerClass","org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")

output_schema = StructType(
    [
        StructField("id", StringType(), True),
        StructField("countAsString", StringType(), True),
    ]
)

class SimpleCounterProcessor(StatefulProcessor):
  def init(self, handle: StatefulProcessorHandle) -> None:
    value_state_schema = StructType([StructField("count", IntegerType(), True)])
    list_state_schema = StructType([StructField("count", IntegerType(), True)])
    self.value_state = handle.getValueState(stateName="valueState", schema=value_state_schema)
    self.list_state = handle.getListState(stateName="listState", schema=list_state_schema)
    # Schema can also be defined using strings and SQL DDL syntax
    self.map_state = handle.getMapState(stateName="mapState", userKeySchema="name string", valueSchema="count int")

  def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
    # Seed the running total from state so the count accumulates across micro-batches
    count = self.value_state.get()[0] if self.value_state.exists() else 0
    for pdf in rows:
      list_state_rows = [(120,), (20,)] # A list of tuples
      self.list_state.put(list_state_rows)
      self.list_state.appendValue((111,))
      self.list_state.appendList(list_state_rows)
      pdf_count = pdf.count()
      count += pdf_count.get("value")
    self.value_state.update((count,)) # Count is passed as a tuple
    iter = self.list_state.get()
    list_state_value = next(iter)[0]
    value = count
    user_key = ("user_key",)
    if self.map_state.exists():
      if self.map_state.containsKey(user_key):
        value += self.map_state.getValue(user_key)[0]
    self.map_state.updateValue(user_key, (value,)) # Value is a tuple
    yield pd.DataFrame({"id": key, "countAsString": str(count)})

q = (df.groupBy("key")
  .transformWithStateInPandas(
    statefulProcessor=SimpleCounterProcessor(),
    outputStructType=output_schema,
    outputMode="Update",
    timeMode="None",
  )
  .writeStream...
)

Python (satır tabanlı)

from pyspark.sql import Row
from pyspark.sql.streaming import StatefulProcessor, StatefulProcessorHandle
from pyspark.sql.types import StructType, StructField, IntegerType, StringType
from typing import Iterator

spark.conf.set("spark.sql.streaming.stateStore.providerClass", "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")

output_schema = StructType(
  [
    StructField("id", StringType(), True),
    StructField("countAsString", StringType(), True),
  ]
)

class SimpleCounterProcessor(StatefulProcessor):
  def init(self, handle: StatefulProcessorHandle) -> None:
    value_state_schema = StructType([StructField("count", IntegerType(), True)])
    list_state_schema = StructType([StructField("count", IntegerType(), True)])
    self.value_state = handle.getValueState(stateName="valueState", schema=value_state_schema)
    self.list_state = handle.getListState(stateName="listState", schema=list_state_schema)
    self.map_state = handle.getMapState(stateName="mapState", userKeySchema="name string", valueSchema="count int")

  def handleInputRows(self, key, rows: Iterator[Row], timerValues) -> Iterator[Row]:
    # Seed the running total from state so the count accumulates across micro-batches
    count = self.value_state.get()[0] if self.value_state.exists() else 0
    for row in rows:
      list_state_rows = [(120,), (20,)]  # A list of tuples
      self.list_state.put(list_state_rows)
      self.list_state.appendValue((111,))
      self.list_state.appendList(list_state_rows)
      count += 1
    self.value_state.update((count,))  # Count is passed as a tuple
    iter_list = self.list_state.get()
    list_state_value = next(iter_list)[0]
    value = count
    user_key = ("user_key",)
    if self.map_state.exists():
      if self.map_state.containsKey(user_key):
        value += self.map_state.getValue(user_key)[0]
    self.map_state.updateValue(user_key, (value,))  # Value is a tuple
    yield Row(id=key[0], countAsString=str(count))

q = (
  df.groupBy("key")
    .transformWithState(
      statefulProcessor=SimpleCounterProcessor(),
      outputStructType=output_schema,
      outputMode="Update",
      timeMode="None",
    )
    .writeStream...
)

Scala programlama dili

import org.apache.spark.sql.streaming._
import org.apache.spark.sql.{Dataset, Encoder, Encoders , DataFrame}
import org.apache.spark.sql.types._
import org.apache.spark.sql.functions._

spark.conf.set("spark.sql.streaming.stateStore.providerClass","org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")

class SimpleCounterProcessor extends StatefulProcessor[String, (String, String), (String, String)] {
  @transient private var countState: ValueState[Int] = _
  @transient private var listState: ListState[Int] = _
  @transient private var mapState: MapState[String, Int] = _

  private val longEncoder = Encoders.scalaLong
  private val intEncoder = Encoders.scalaInt
  private val stringEncoder = Encoders.STRING

  override def init(
      outputMode: OutputMode,
      timeMode: TimeMode): Unit = {
    countState = getHandle.getValueState[Int]("countState",
      intEncoder, TTLConfig.NONE)
    listState = getHandle.getListState[Int]("listState",
      intEncoder, TTLConfig.NONE)
    mapState = getHandle.getMapState[String, Int]("mapState",
      stringEncoder, intEncoder, TTLConfig.NONE)
  }

  override def handleInputRows(
      key: String,
      inputRows: Iterator[(String, String)],
      timerValues: TimerValues): Iterator[(String, String)] = {
    var count = countState.getOption().getOrElse(0)
    for (row <- inputRows) {
      val listData = Array(120, 20)
      listState.put(listData)
      listState.appendValue(count)
      listState.appendList(listData)
      count += 1
    }
    val iter = listState.get()
    var listStateValue = 0
    if (iter.hasNext) {
      listStateValue = iter.next()
    }
    countState.update(count)
    var value = count
    val userKey = "userKey"
    if (mapState.exists()) {
      if (mapState.containsKey(userKey)) {
        value += mapState.getValue(userKey)
      }
    }
    mapState.updateValue(userKey, value)
    Iterator((key, count.toString))
  }
}

val q = spark
        .readStream
        .format("delta")
        .load("$srcDeltaTableDir")
        .as[(String, String)]
        .groupByKey(x => x._1)
        .transformWithState(
            new SimpleCounterProcessor(),
            TimeMode.None(),
            OutputMode.Update(),
        )
        .writeStream...

Örneği uçtan uca çalıştırın

Not

Bu sayfadaki çalıştırılabilir örnekler, mevcut verilerinizi etkilemeden çalışabilmek için özel main.stateful_examples bir şema içinde tablolar oluşturur. Katalogda main şema oluşturma izniniz yoksa, örneklerdeki katalog ve şemayı tablo oluşturabileceğiniz bir konuma değiştirin.

Yukarıdaki işlemci durumlu mantığı tanımlıyor ama sorgu başlatmıyor. SimpleCounterProcessor öğesini kopyala-yapıştır ile çalıştırmak için, akış kaynağı olarak küçük bir Delta Lake tablosuna başlangıç verisi ekleyin, ardından bellek içi bir hedefe yazan bir sorgu başlatın. Bu örnek, sorgunun başlangıç satırlarını işlemesini ve ardından durmasını sağlamak için Trigger.AvailableNow kullanır. Kaynağı tohumlamak ve sorguyu başlatmak için aşağıdakileri çalıştırın:

import uuid

# Create a dedicated schema for the example tables
spark.sql("CREATE SCHEMA IF NOT EXISTS main.stateful_examples")

# Seed a small Delta table to use as the streaming source
spark.sql("DROP TABLE IF EXISTS main.stateful_examples.tws_counter_source")
spark.createDataFrame(
  [("a", "1"), ("a", "2"), ("a", "3"), ("b", "1"), ("b", "2")],
  "key string, value string",
).write.saveAsTable("main.stateful_examples.tws_counter_source")

df = spark.readStream.table("main.stateful_examples.tws_counter_source")

q = (
  df.groupBy("key")
    .transformWithState(
      statefulProcessor=SimpleCounterProcessor(),
      outputStructType=output_schema,
      outputMode="Update",
      timeMode="None",
    )
    .writeStream.format("memory")
    .queryName("counter_output")
    .option("checkpointLocation", f"/tmp/checkpoint_{uuid.uuid4()}")
    .trigger(availableNow=True)
    .start()
)

q.awaitTermination()

Sorgu tamamlandıktan sonra, her gruplama anahtarının sayısını görüntüleyin:

display(spark.sql("SELECT id, countAsString FROM counter_output ORDER BY id"))

Anahtarın a üç, anahtarın b iki satırı var, bu yüzden sorgu geri döner:

id  countAsString
a   3
b   2

Daha fazla örnek için bkz. örnek durum bilgisi olan uygulamalar.

Not

Python'da durum değerleri demetlerdir. Demetleri put ve update öğelerine geçirin ve get öğesinden demetler bekleyin.

Örneğin, şemanız ValueState tek bir tamsayıysa:

current_value_tuple = value_state.get() # Returns the value state as a tuple
current_value = current_value_tuple[0]  # Extracts the first item in the tuple
new_value = current_value + 1           # Calculate a new value
value_state.update((new_value,))        # Pass the new value formatted as a tuple

ListState içindeki öğeler veya MapState içindeki değerler için de bu yaklaşımı kullanın.

Satırları gönder

Her gruplama anahtarı için handleInputRows’nin nasıl satır üreteceğini tanımlamak üzere handleExpiredTimer ya da transformWithState kullanmanız gerekir. Bkz . Giriş satırlarını işleme ve Süresi dolan zamanlayıcıları işleme.

Özel durum bilgisi olan uygulamalar, durum bilgilerinin nasıl kullanılacağı hakkında hiçbir varsayımda bulunmaz. Belirli bir koşul için uygulama hiçbir satır, bir satır veya çok sayıda satır yaymayabilir.

Not

Birden çok durum değeri uygulayabilir ve satır yayma için birden çok koşul tanımlayabilirsiniz, ancak tüm satırların aynı şemayı kullanması gerekir.

Python (Pandas)

transformWithStateInPandas ile, çıkış şemanızı outputStructType anahtar sözcüğü ile tanımlayın.

pandas DataFrame nesnesini ve yield kullanarak satırları çıkarın.

İsteğe bağlı olarak boş bir DataFrame oluşturabilirsiniz yield . update çıkış modunu kullanır ve boş bir DataFrame çıktı olarak verirseniz, bu işlem gruplandırma anahtarının değerlerini null olarak günceller.

Python (satır tabanlı)

transformWithState ile, çıkış şemanızı outputStructType anahtar sözcüğü ile tanımlayın.

Row nesnesi ve yield kullanarak satır yayımlama.

İsteğe bağlı olarak, boş bir yineleyici döndürebilirsiniz. update çıkış modunu kullanır ve boş bir yineleyici yayımlarsanız, bu işlem gruplandırma anahtarının değerlerini null olarak günceller.

Scala programlama dili

Scala'da bir nesne kullanarak Iterator satırları yayarsınız. Şema, yayılan satırların şemasından otomatik olarak türetilir.

İsteğe bağlı olarak, boş bir Iteratordöndürebilirsiniz. update çıkış modunu kullanıyor ve boş bir Iterator yayıyorsanız, bu gruplandırma anahtarının değerlerini null olarak günceller.

Başlangıç durumunu işle

İsteğe bağlı olarak, başlangıç durumunu ilk micro-batch'e aktarabilirsiniz.

Örneğin, aşağıdakileri yapmak için bunu kullanabilirsiniz:

  • Mevcut bir iş akışını yeni bir özel uygulamaya geçirme.
  • Şemanızı veya mantığınızı değiştirmek için durum bilgisi olan bir işleci yükseltin.
  • Otomatik olarak onarılamayan ve manuel müdahale gerektiren bir hatayı giderin.

Not

Mevcut bir denetim noktasından durum bilgilerini sorgulamak için durum deposu okuyucusunu kullanın. Bkz. Yapılandırılmış Akış durumu bilgilerini okuyun.

Var olan bir Delta tablosunu durum bilgisi olan bir uygulamaya dönüştürüyorsanız, spark.read.table("table_name") kullanarak tabloyu okuyun ve sonuçta elde edilen DataFrame'i geçirin. Alanları, yeni durum bilgisi olan uygulamanızla tam uyumlu olacak şekilde isteğe bağlı olarak seçebilir veya değiştirebilirsiniz.

Giriş satırlarıyla aynı gruplandırma anahtarı şemasına sahip bir DataFrame kullanarak bir başlangıç durumu sağlarsınız.

Not

Python, handleInitialState tanımlarken ilk durumu belirtmek için StatefulProcessor kullanır. Scala, StatefulProcessorWithInitialStateayrı sınıfını kullanır.

Aşağıdaki örnek, anahtar başına bir sayacı mevcut bir Delta tablosundaki verilerle başlatır:

Python (satır tabanlı)

from pyspark.sql import Row
from pyspark.sql.streaming import StatefulProcessor, StatefulProcessorHandle
from pyspark.sql.types import StructType, StructField, IntegerType, StringType
from typing import Iterator

class CounterWithInitialState(StatefulProcessor):
  def init(self, handle: StatefulProcessorHandle) -> None:
    state_schema = StructType([StructField("count", IntegerType(), True)])
    self.count_state = handle.getValueState("countState", state_schema)

  def handleInitialState(self, key, initialState: Row, timerValues) -> None:
    self.count_state.update((initialState["count"],))

  def handleInputRows(self, key, rows: Iterator[Row], timerValues) -> Iterator[Row]:
    count = self.count_state.get()[0] if self.count_state.exists() else 0
    for _ in rows:
      count += 1
    self.count_state.update((count,))
    yield Row(id=key[0], count=count)

  def close(self) -> None:
    pass

output_schema = StructType([
  StructField("id", StringType(), True),
  StructField("count", IntegerType(), True),
])

import uuid

# Create a dedicated schema for the example tables
spark.sql("CREATE SCHEMA IF NOT EXISTS main.stateful_examples")

# Seed existing per-key counts to load as the initial state
spark.sql("DROP TABLE IF EXISTS main.stateful_examples.existing_counts")
spark.createDataFrame(
  [("x", 10)],
  "id string, count int",
).write.saveAsTable("main.stateful_examples.existing_counts")

# Seed a small Delta table to use as the streaming source
spark.sql("DROP TABLE IF EXISTS main.stateful_examples.tws_initial_source")
spark.createDataFrame(
  [("x", "a"), ("x", "b")],
  "id string, value string",
).write.saveAsTable("main.stateful_examples.tws_initial_source")

df = spark.readStream.table("main.stateful_examples.tws_initial_source")

# Load existing counts as initial state — must use the same grouping key as the input
initial_state = spark.read.table("main.stateful_examples.existing_counts").groupBy("id")

q = (
  df.groupBy("id")
    .transformWithState(
      statefulProcessor=CounterWithInitialState(),
      outputStructType=output_schema,
      outputMode="Update",
      timeMode="None",
      initialState=initial_state,
    )
    .writeStream.format("memory")
    .queryName("initial_state_output")
    .option("checkpointLocation", f"/tmp/checkpoint_{uuid.uuid4()}")
    .trigger(availableNow=True)
    .start()
)

q.awaitTermination()

# The initial state seeds "x" with 10, and the source adds two rows, so the count is 12
display(spark.sql("SELECT id, count FROM initial_state_output ORDER BY id"))

Scala programlama dili

import org.apache.spark.sql.streaming._
import org.apache.spark.sql.Encoders

class CounterWithInitialState
    extends StatefulProcessorWithInitialState[String, (String, String), (String, String), (String, Int)] {

  @transient private var countState: ValueState[Int] = _

  override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
    countState = getHandle.getValueState[Int]("countState", Encoders.scalaInt, TTLConfig.NONE)
  }

  override def handleInitialState(
      key: String, initialState: (String, Int), timerValues: TimerValues): Unit = {
    countState.update(initialState._2)
  }

  override def handleInputRows(
      key: String,
      rows: Iterator[(String, String)],
      timerValues: TimerValues): Iterator[(String, String)] = {
    val count = if (countState.exists()) countState.get() else 0
    val newCount = count + rows.size
    countState.update(newCount)
    Iterator((key, newCount.toString))
  }
}

// Load existing counts as initial state — must use the same grouping key as the input
val initialState = spark.read.table("existing_counts")
  .as[(String, Int)]
  .groupByKey(_._1)

val q = spark
  .readStream
  .format("delta")
  .load(srcDeltaTableDir)
  .as[(String, String)]
  .groupByKey(_._1)
  .transformWithState(
    new CounterWithInitialState(),
    TimeMode.None(),
    OutputMode.Update(),
    initialState,
  )
  .writeStream...

Asenkron işlem (Beta)

PythontransformWithState, durum işlemlerini ve kullanıcı mantığını eşzamanlı çalıştırmak için asenkron işlem yapmayı destekler. Asenkron işleme, senkron işlemeye göre daha yüksek veri verimliliğine sahiptir ve üçüncü taraf asenkron kütüphaneler olmadan sadece küçük kod değişiklikleri gerektirir. Asenkron işlem kullanmak için, senkron AsyncStatefulProcessoryerine bir StatefulProcessor uygulayın. Bkz. (Beta) ile Asenkron İşlemetransformWithState.

transformWithState ifadesini Lakeflow işlem hatlarında kullanın

Python kullanarak akış işlem hatlarınızda istenen durum bilgisi kullanan mantığı uygulamak için Lakeflow işlem hatlarında transformWithState işlecini kullanın.

Bunu yapmak için aşağıdaki adımları tamamlayın:

  1. Rastgele durum bilgisi olan dönüşümleriniz için çıkış şemasını ve durum bilgisi olan işlemci mantığını tanımlayın. Örnekler için bkz. Durum bilgisi olan uygulamalar örneği.
  2. DataFrame üzerinde transformWithState işleçini çağıran bir Lakeflow işlem hattı akışı oluşturun. Bkz. Öğretici: Lakeflow Pipelines Düzenleyicisi'ni kullanarak ilk işlem hattınızı oluşturma.
  3. İşlem hattınızı çalıştırın ve sonuçları hedef tablo veya havuz üzerinde doğrulayın.

Algılayıcı sinyallerini izlemek için transformWithState kullanan bir örnek için bkz Örnek: Algılayıcı sinyallerini izlemek için transformWithState kullanma.