使用
你可以用 transformWithState 來建立有狀態串流應用程式,並實作低延遲且近即時的解決方案。 透過自訂的有狀態運算子,你可以建立任意的有狀態邏輯,讓你建立傳統結構化串流處理無法實現的新操作用例。
注意
對於有狀態操作,如聚合、重複去重和串流連接,Databricks 建議使用內建的結構化串流運算子,而非自訂邏輯。 請參閱 什麼是具狀態串流?。
Databricks 建議在進行任意狀態轉換時,使用 transformWithState 來取代舊版運算子,例如 flatMapGroupsWithState 和 mapGroupsWithState。 參見 遺留的任意有狀態運算子。
要求
transformWithState與transformWithStateInPandas運算子有以下要求:
- 適用於 Databricks Runtime 16.2 和更新版本。
- 即時模式則可使用 Databricks Runtime 17.3 LTS 或以上版本。 詳見 即時模式概念。
- 標準存取模式中,Python 支援 Databricks Runtime 16.3 及以上版本,Scala 則支援 Databricks 執行環境 17.3 及以上版本。
- RocksDB 是 Databricks Runtime 17.3 及以上版本的預設狀態儲存提供者。
對於 Databricks Runtime 17.2 及以下版本,您必須設定 RocksDB 狀態儲存提供者。 Databricks 建議在 Spark 配置中啟用 RocksDB。
spark.conf.set("spark.sql.streaming.stateStore.providerClass", "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")
什麼是 transformWithState?
transformWithState 運算符將自定義的具狀態處理器應用於結構化串流查詢。 您必須實作自訂具狀態處理器,才能使用 transformWithState。 結構化串流提供用於使用 Python、Scala 或 Java 建構有狀態處理器的 API。
用 transformWithState 來套用自訂邏輯到分組鍵。 下列描述高階設計:
- 定義一或多個狀態變數。
- 每個分組鍵的狀態資訊會持續存在。 你可以用使用者自訂的程式碼存取每個狀態變數。
- 對於每個已處理的微批次,該鍵的所有資料列都會以迭代器的形式提供。
- 使用
StatefulProcessorHandle、計時器和使用者自訂條件來控制如何輸出資料列。 - 為了管理狀態到期與狀態大小,狀態值支援個別的存活時間(TTL)定義。
因為 transformWithState 支援狀態儲存中的結構演化,你可以在不遺失歷史狀態資訊的情況下迭代和更新生產應用程式。 更新狀態結構後,你不需要重新處理資料列,這簡化了程式碼部署和維護。 請參閱 狀態存放區中的架構演進。
重要
Azure Databricks 文件用transformWithState來描述 Python 和 Scala 的實作:
- PySpark 支援基於
transformWithState資料列的 API 與基於 Pandas 的transformWithStateInPandas算子。-
transformWithStateInPandas在即時模式下不支援。 請改用transformWithState。 詳情請參見transformWithState即時模式。 - 以
transformWithState為基礎的資料列式 API 支援透過asyncio進行非同步處理,以提高吞吐量。 無伺服器運算不支援非同步處理。 參見非同步處理(Beta)。
-
- Scala 僅支援基於
transformWithState列的 API。
Scala 和 Python 的實作transformWithState功能相同,但語法上有些差異。
定義一個 StatefulProcessor
你透過擴充 StatefulProcessor 類別並實作其方法來定義有狀態處理器。
Spark 會將 StatefulProcessorHandle 傳遞給您的 init 的 StatefulProcessor 方法。 利用 handle 來建立狀態變數並與狀態儲存互動。
transformWithState支援三種狀態類型:ValueState、、 ListStateMapState和 。 每種類型都使用不同的底層資料結構,為每個群組鍵儲存狀態。
實作以下方法來定義你的自訂邏輯:
- 實作
handleInputRows,以控制應用程式如何處理資料、更新狀態,以及為每個微批次輸出資料列。 請參閱 處理輸入資料列。 - 實作
handleExpiredTimer,以便無論群組鍵是否在微批次中接收新的資料列,都能執行時間型邏輯。 請參見 處理過期計時器。 - 您可以選擇實作
handleInitialState,以便在應用程式處理任何輸入列之前預先填入狀態。 請參見 處理初始狀態。
下表比較了這些方法的功能行為:
| 行為 | handleInputRows |
handleExpiredTimer |
|---|---|---|
| 取得、放置、更新或清除狀態值 | 是的 | 是的 |
| 建立或刪除定時器 | 是的 | 是的 |
| 輸出資料列 | 是的 | 是的 |
| 在目前微批次中反覆迭代各列 | 是的 | 不 |
| 根據時間推移觸發邏輯 | 不 | 是的 |
你可以結合兩者handleInputRowshandleExpiredTimer,根據需要實作複雜的邏輯。
例如,您可以實作一個應用程序,使用 handleInputRows 來更新每個微批次的狀態值,並設置一個 10 秒後觸發的計時器。 如果沒有處理任何額外的資料列,你可以使用 handleExpiredTimer 輸出狀態存放區中的目前值。 如果有新的資料列針對該分組鍵進行處理,你可以清除現有的計時器,並設定新的計時器。
StatefulProcessorHandle
在 PySpark 中,這個 StatefulProcessorHandle 類別允許你存取控制程式碼如何使用狀態資訊的函式。
初始化 StatefulProcessor 時,您都必須匯入並將 StatefulProcessorHandle 傳遞給 handle 變數。
handle 變數會把你 Python 類別的本地變數綁定到狀態變數。
注意
Scala 會使用 getHandle 方法。
自定義狀態類型
您可以在單一具狀態運算符中實作多個狀態物件。
根據你完整的應用程式邏輯選擇狀態類型。 例如,你可以使用 ValueState 追蹤工作階段,並依 user_id 和 session_id 分組。 或者,若要評估跨多個工作階段的條件,可使用依 MapState 分組,並以 user_id 作為對應表鍵的 session_id。
如果你的狀態物件使用 StructType,你必須為綱要中結構的每個欄位定義唯一的名稱。 讀取狀態儲存區時,可以看到這些名稱。 請參閱 查看結構化串流狀態資訊。
以下章節說明由 transformWithState以下支援的狀態類型:
ValueState
ValueState 為每個分組鍵儲存一個值。
值狀態可以包含複雜類型,例如結構體或元組。 對於 ValueState,你必須實作邏輯來替換整個值。
值狀態的存活時間會在值更新時重置。 如果您處理 ValueState 的來源金鑰時未更新已儲存的 ValueState,存留時間不會重設。
ListState
ListState 儲存每個分組鍵的清單。
清單狀態是值的集合,每個值都可以包含複雜類型。 清單中的每個值都有自己的存活時間。
您可以附加個別專案、附加專案清單,或使用 put覆寫整個清單,以將專案新增至清單。 若要重設存留時間,您必須執行 put 作業。
MapState
MapState 會為每個分組索引鍵儲存一個對應表。 Maps 是 Apache Spark 上的 Python 字典(dict)。
映射狀態是一組不同鍵,每個鍵對應一個值,每個鍵值中可能包含複數型態。 地圖中的每個鍵值對都有自己的存活時間。
你可以更新特定鍵的值,或者移除鍵及其值。 您可以使用鍵傳回個別值、列出所有鍵、列出所有值,或傳回迭代器,以便操作映射中完整的鍵值對。
重要
群組索引鍵描述結構化串流查詢 GROUP BY 子句中指定的欄位。 Map 狀態可針對某個分組鍵包含任意數量的鍵值對。
例如,如果您的查詢使用 GROUP BY user_id,且您想為每個 session_id 定義一個對應,則分組鍵是 user_id,而 MapState 鍵是 session_id:
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
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())
在 StatefulProcessor 中建立自訂狀態變數
當您初始化 StatefulProcessor時,您會為每個狀態物件建立局部變數,讓您與自定義邏輯中的狀態對象互動。 透過覆寫 init 類別內建 StatefulProcessor 的方法來定義並初始化狀態變數。
你可以在你的 getValueState 中使用 getListState、getMapState 和 StatefulProcessor 方法定義任意數量的狀態物件。
每個狀態對象都必須具有以下項目:
- 唯一名稱
- 一個結構描述
- 在 Python 裡,你必須指定 schema。
- 在 Scala 裡,你可以傳一個
Encoder指定狀態架構。
你也可以選擇性地提供以毫秒為單位的生存時間(TTL)持續時間。 如果要實作映射狀態,您必須為映射的鍵與值提供個別的架構定義。
注意
StatefulProcessor 分別處理查詢、更新及發出狀態資訊的邏輯。 請參閱在具有自訂邏輯的方法中使用您的狀態變數。
在有自訂邏輯的方法裡使用你的狀態變數
狀態物件有取得狀態、更新現有狀態資訊及清除當前狀態的方法。
每個分組金鑰都有專屬的狀態資訊。
-
StatefulProcessor會根據你的自訂邏輯和指定的輸出結構描述輸出資料列。 請參見 發射列。 - 使用讀取器
statestore存取狀態儲存中的數值。 此讀卡器設計用於批次工作負載,並非低延遲工作負載。 請參閱 查看結構化串流狀態資訊。 - 使用
handleInputRows指定的邏輯,只有在微批次中存在該索引鍵對應的資料列時才會執行。 請參閱 處理輸入資料列。 - 用
handleExpiredTimer來實作基於時間的邏輯,不需要依賴觀察射擊列。 請參見 處理過期計時器。
注意
狀態物件會藉由將索引鍵分組來隔離,其含意如下:
- 狀態值不會被與不同群組鍵相關的列影響。
- 您 無法 實作取決於比較值或更新群組索引鍵狀態的邏輯。
您可以比較群組索引鍵內的值。 使用 MapState 透過第二個供自訂邏輯使用的索引鍵來實作邏輯。 例如,依 user_id 分組,並將 ip_address 用於您的 MapState 金鑰,即可追蹤同時進行的使用者工作階段。
操作狀態的進階考慮事項
狀態更新具備容錯能力。 若任務在微批次處理完成前當機,重試將使用最後一次成功微批次的值。
為了優化效能,Databricks 建議你在迭代器中處理同一鍵的所有值,並在一次寫入中提交更新。 當你寫入狀態變數時,會觸發寫入 RocksDB。
狀態值沒有預設值。 如果你的邏輯需要讀取現有狀態資訊,就用這個 exists 方法。
為了實作空狀態的邏輯, MapState 變數允許你檢查單一鍵值或列出所有鍵。
處理輸入數據列
使用此 handleInputRows 方法定義應用程式如何處理資料列並更新狀態值。 這個方法每次結構化串流查詢處理分組鍵的列時都會執行。
對於使用 transformWithState實作的大部分具狀態應用程式,核心邏輯是使用 handleInputRows定義。
每次處理一次微批次更新,該微批次中特定分組鍵的所有列皆可透過迭代器取得。 使用者自訂邏輯可與目前微批次中的所有列及狀態儲存中的值互動。
處理過期計時器
使用這個 handleExpiredTimer 方法來實作基於經過時間的自訂邏輯。
在群組鍵值內,定時器會以其時間戳進行唯一識別。
定時器到期時,結果會由應用程式中實作的邏輯決定。 常見的模式包括:
- 發出儲存在狀態變數中的資訊。
- 清除儲存的狀態資訊。
- 建立新的定時器。
即使在微批次中沒有處理其關聯金鑰的任何資料列,已過期的計時器仍會觸發。
指定時間模式
當你將 傳送 StatefulProcessor 到 transformWithState時,必須用參數 timeMode 指定時間模式。
支援下列選項:
| 時間模式 | 描述 |
|---|---|
ProcessingTime |
支援計時器和 TTL,並會根據 Apache Spark 處理每個微批次時的壁鐘時間進行評估。 當你希望計時器相對於資料列處理的時間以固定間隔觸發,且不受資料中時間戳記影響時,請使用 ProcessingTime。 |
EventTime |
支援計時器,並依據事件時間浮水印進行判定。 浮水印會隨著 Apache Spark 在輸入資料中觀察到的時間戳而向前推進。 TTL 不支援 EventTime。 當您的資料包含時間戳記,且您希望根據這些時間戳記的進展來觸發計時器時,請使用 EventTime。 使用 EventTime時,也必須指定參數。eventTimeColumnName 參見 eventTimeColumnName。 |
NoTime 或 TimeMode.None() |
不支援計時器和 TTL。 當你的有狀態應用程式不需要基於時間的邏輯時使用 NoTime 。 |
eventTimeColumnName
使用 EventTime 時間模式時, eventTimeColumnName 參數會指定輸出結構中包含事件時間戳記的欄位名稱。 Apache Spark 利用此欄位將浮水印傳播至輸出串流,從而實現正確的下游時間基礎操作。
Python
eventTimeColumnName 是 transformWithState 或 transformWithStateInPandas 的額外引數:
q = (
df.groupBy("key")
.transformWithState(
statefulProcessor=MyProcessor(),
outputStructType=output_schema,
outputMode="Append",
timeMode="EventTime",
eventTimeColumnName="outputTimestamp",
)
.writeStream...
)
程式語言 Scala
transformWithState 接受 eventTimeColumnName 代替 timeMode。 此方法一律使用 EventTime 模式:
val q = spark
.readStream
.format("delta")
.load(srcDeltaTableDir)
.as[(String, String)]
.groupByKey(x => x._1)
.transformWithState(
new MyProcessor(),
"outputTimestamp",
OutputMode.Append(),
)
.writeStream...
內建定時器值
Databricks 建議不要在自定義具狀態應用程式中使用系統時鐘,因為這可能會導致工作失敗時的重試不可靠。 當您必須存取處理時間或浮水印時,請使用 TimerValues 類別中的 方法:
TimerValues |
描述 |
|---|---|
getCurrentProcessingTimeInMs |
傳回自 epoch 以來目前批次的處理時間時間戳,以毫秒為單位。 |
getCurrentWatermarkInMs |
傳回自紀元時刻起算,目前批次水印的時間戳記,以毫秒為單位。 |
注意
處理時間描述 Apache Spark 處理微批次的時間。 許多串流來源,例如 Kafka,也包含系統處理時間。
串流查詢中的浮水印通常是針對事件時間或串流來源的處理時間來定義的。 請參閱 套用浮水印來控制資料處理臨界值。
水印和視窗都可以與 transformWithState搭配使用。 您可以利用TTL、定時器和 MapState 或 ListState 功能,在自定義具狀態應用程式中實作類似的功能。
狀態類型的存活時間(TTL)
為防止記憶體外錯誤並移除過時狀態類型值, transformWithState 支援每個狀態類型值的可選性存活時間(TTL)值。 過期後,TTL 會靜默地移除狀態型態值。 TTL 不會執行 handleExpiredTimer ,也不會有任何自訂邏輯。 如果要在狀態到期時執行程式碼,可以用計時器。
重要
如果你沒有實作 TTL,就必須清除狀態資料,以避免發生記憶體耗盡錯誤。
對於所有狀態類型,TTL 在更新狀態資訊時會重置。 TTL 對每種狀態類型值強制執行,且每種狀態類型有不同的規則:
- 狀態變數的範圍設定為群組鍵值。
- 對於
ValueState物件,每個群組索引鍵只會儲存單一值。 TTL 適用於此值。 - 對於
ListState物件,清單可以包含許多值。 TTL 會獨立套用至清單中的每個值。- 雖然 TTL 是針對 中的
ListState個別值設限,但更新個別值的唯一方法是使用 該put方法,該方法會覆蓋變數的ListState全部內容,並重置列表中所有值的 TTL。
- 雖然 TTL 是針對 中的
- 對於
MapState物件,每個對應索引鍵都有相關聯的狀態值。 TTL 會獨立套用至映射中的每個鍵值對。
注意
計時器可讓你定義狀態清除以外的自訂邏輯,包括發出資料列。 你可以選擇使用計時器來清除給定狀態值的狀態資訊,並發出值或觸發條件邏輯。 請參見 處理過期計時器。
具狀態應用程式的範例
以下範例定義了自訂的有狀態處理器, SimpleCounterProcessor包含範例狀態變數。
SimpleCounterProcessor 使用 ValueState、 ListState、 和 MapState 來計算每個分組鍵的行數。
Python(熊貓)
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(基於行)
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
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...
完整執行這個範例
注意
本頁可執行的範例是在專用 main.stateful_examples 架構中建立資料表,讓它們能在不影響現有資料的情況下執行。 如果你沒有權限在目錄中 main 建立結構,請將範例中的目錄和結構改成可以建立資料表的位置。
上述處理器定義了有狀態邏輯,但不會啟動查詢。 若要透過複製貼上來執行 SimpleCounterProcessor,請先建立一個小型的 Delta Lake 資料表作為串流來源,接著啟動將資料寫入記憶體內接收端的查詢。 這個範例使用 Trigger.AvailableNow,使查詢處理已植入的資料列後停止。 若要初始化來源並啟動查詢,請執行下列命令:
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()
查詢完成後,查看每個分組鍵的計數:
display(spark.sql("SELECT id, countAsString FROM counter_output ORDER BY id"))
a 鍵有三列,而 b 鍵有兩列,因此查詢回傳:
id countAsString
a 3
b 2
如需更多範例,請參閱 範例具狀態應用程式。
注意
在 Python 中,狀態值是元組。 將元組傳遞給 put 和 update,並期望從 get中獲得元組。
例如,若你的 ValueState 結構是一個整數:
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 中的項目或 MapState 中的值。
輸出資料列
你必須使用 handleInputRows 或 handleExpiredTimer 來定義 transformWithState 如何針對每個分組鍵發出資料列。 請參見 「處理輸入列 」和 「處理過期計時器」。
自訂有狀態應用程式對於如何使用狀態資訊不做任何假設。 對於某個特定條件,應用程式可能不輸出任何資料列、輸出一筆資料列,或輸出多筆資料列。
注意
你可以實作多個狀態值並定義多個發出列的條件,但所有列都必須使用相同的結構。
Python(熊貓)
使用 transformWithStateInPandas,以 outputStructType 關鍵字定義您的輸出結構描述。
使用 pandas DataFrame 物件和 yield 輸出資料列。
或者,你可以 yield 空的 DataFrame。 如果你使用 update 輸出模式並輸出空的 DataFrame,這會將分組鍵的值更新為 null。
Python(基於行)
使用 transformWithState,以 outputStructType 關鍵字定義您的輸出結構描述。
使用 Row 物件和 yield 發出資料列。
你可以選擇回傳一個空的迭代器。 如果你使用 update 輸出模式並輸出空迭代器,這會更新分組鍵的值為 null。
程式語言 Scala
在 Scala 中,你用 Iterator 物件來發射列。 該結構描述會根據輸出的資料列結構描述自動推導而成。
你可以選擇性地回傳一個空的 Iterator。 如果你使用 update 輸出模式並發出空值 Iterator,這會更新分組鍵的值為 null。
處理初始狀態
你也可以選擇將初始狀態傳給第一個微批次。
例如,你可以用這個來:
- 將現有工作流程遷移到新的自訂應用程式。
- 升級一個有狀態運算子來改變你的結構架構或邏輯。
- 修復無法自動修復且需要人工介入的故障。
注意
使用狀態存放區讀取器從現有的檢查點查詢狀態資訊。 請參閱 查看結構化串流狀態資訊。
如果您要將現有的 Delta 資料表轉換成具狀態應用程式,請使用 spark.read.table("table_name") 讀取數據表,並傳遞產生的 DataFrame。 您可以選擇性地選取或修改欄位,以符合您新的有狀態的應用程式。
您可以使用 DataFrame 與輸入資料列相同的群組索引鍵架構,提供初始狀態。
注意
Python 使用 handleInitialState 來指定初始狀態,同時定義 StatefulProcessor。 Scala 會使用相異類別 StatefulProcessorWithInitialState。
以下範例示範如何從現有的 Delta 資料表初始化每個鍵的計數器:
Python(基於行)
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
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...
非同步處理(Beta)
Python transformWithState 支援使用 asyncio 進行非同步處理,以同時執行狀態操作與使用者邏輯。 非同步處理的吞吐量高於同步處理,且只需少量程式碼修改,無需第三方非同步函式庫。 若要使用非同步處理,請實作 a AsyncStatefulProcessor 代替同步處理 StatefulProcessor。 請參見 非同步處理(Beta transformWithState )。
在湖流管線中的應用transformWithState
在 Lakeflow pipelines 中使用transformWithState運算子,使用 Python 在你的串流管道中實作任意有狀態邏輯。
若要這樣做,請完成下列步驟:
- 定義任意有狀態轉換的輸出結構和有狀態處理器邏輯。 如需範例,請參見 具狀態應用程式範例。
- 建立 Lakeflow 管線流程,在 DataFrame 上呼叫
transformWithState運算子。 請參考 教學:使用 Lakeflow Pipelines Editor 建立你的第一個管線。 - 執行你的管線並在目標表或匯表上驗證結果。
關於用於 transformWithState 監測感測器心跳的範例,請參見 範例:用於 transformWithState 監測感測器心跳。