非同步處理( transformWithState Beta)

Important

Python 列式 transformWithState API 的非同步處理目前仍處於測試階段。 請參見 Azure Databricks預覽版本。

非同步處理可在 Databricks Runtime 19 及以上版本中使用。

Python transformWithState 支援基於 asyncio的非同步處理。 透過同時在分組鍵與處理程序間通訊,非同步處理的吞吐量高於僅有少量程式碼變更的同步處理。 這種吞吐量提升不需要任何第三方非同步函式庫。 進階使用者可透過非同步程式設計模式及啟用非同步函式庫進一步優化應用程式。

若要使用非同步處理,請實作 a AsyncStatefulProcessor 代替同步處理 StatefulProcessor。 API 會 AsyncStatefulProcessor 鏡 StatefulProcessor 像同步 API,因此大多數應用程式只需小幅修改即可使用非同步 API。 參見 實作一個 AsyncStatefulProcessor。

關於同步 transformWithState API 與核心概念,請參見 「建構自訂有狀態應用程式與 transformWithState。

Note

非同步處理僅支援 Python 列式 transformWithState API。 它不支援 transformWithStateInPandas 或支援 Scala transformWithState API。 無伺服器運算不支援非同步處理。

實作一個 AsyncStatefulProcessor

要將同步轉換 StatefulProcessor 成 AsyncStatefulProcessor,請進行以下修改:

  • 用關鍵字定義 API 方法(init, handleExpiredTimerhandleInitialStateclosehandleInputRows和)。async def
  • 讀取並更新狀態與計時器值,await或使用 Python 的asyncio函式庫執行。 這適用於狀態操作,如 , valueState.get() 以及計時器操作,如 registerTimer。 建立狀態物件(如 handle.getValueState)則保持同步。

非同步處理需考慮以下事項:

  • 如果你的應用程式將資料儲存在成員變數或外部系統中,Databricks 建議你重寫邏輯以保障並行執行的安全。 由於 handleInputRows 和 handleExpiredTimer 可以在分組鍵間同時執行,交錯執行不得破壞共享資料。 大多數申請已經符合此要求。
  • Databricks 建議不要捕捉或抑制狀態操作的錯誤。 Apache Spark 會幫你處理這些錯誤。 若狀態操作失敗,Apache Spark 會失敗該任務並重新嘗試。
    • 在 AsyncStatefulProcessor狀態操作錯誤中,錯誤會被你管理,且不會顯示在你的程式碼中。
    • 在同步 StatefulProcessor的狀態操作中,程式碼會產生錯誤,但若抑制這些錯誤,可能會損害資料的正確性。

範例:每個分組鍵的計數列數

以下範例定義了 , AsyncCountProcessor 計算每個分組鍵的列數。 變 value_schema 數定義了儲存運行計數的結構 ValueState 。 與同步方式 StatefulProcessor相比,變更是 async def 每個方法及 await 狀態讀取與更新操作的關鍵字。 getValueState init呼叫進入仍保持同步。 以下程式碼定義處理器:

from pyspark.sql import Row
from pyspark.sql.streaming import AsyncStatefulProcessor, AsyncStatefulProcessorHandle
from pyspark.sql.types import StructType, StructField, LongType

value_schema = StructType([StructField("count", LongType(), True)])

class AsyncCountProcessor(AsyncStatefulProcessor):
  async def init(self, handle: AsyncStatefulProcessorHandle) -> None:
    self.count = handle.getValueState("count", value_schema)

  async def handleInputRows(self, key, rows, timerValues):
    total = (await self.count.get() or (0,))[0]
    for _ in rows:
      total += 1
    await self.count.update((total,))
    yield Row(action=key[0], count=total)

  async def close(self) -> None:
    pass

用非同步處理器執行查詢

要用非同步處理器執行查詢,請將你的 AsyncStatefulProcessor 傳給 transformWithState。 查詢使用與同步路徑相同的語法。 非同步和同步 API 共享相同的狀態格式,因此你可以在重複使用相同檢查點的情況下,在StatefulProcessor現有查詢間AsyncStatefulProcessor切換同步查詢。

範例:樣本資料集中的 events 計數事件

以下範例以範例資料集為例events。AsyncCountProcessor 每個記錄有一個 time 欄位(紀元秒)和 action 一個欄位,值為 Open 或 Close。 查詢會依 分 action 組並計數每個動作類型的事件。 更多範例資料集請參見 範例資料集。

變 input_schema 數定義了來源記錄的結構, output_schema 變數定義了處理器所發出的資料列結構。 若要將範例資料集視為串流,請定義兩個結構,然後如以下程式碼開始查詢:

from pyspark.sql.types import StructType, StructField, StringType, LongType

input_schema = StructType([
  StructField("time", LongType(), True),
  StructField("action", StringType(), True),
])

output_schema = StructType([
  StructField("action", StringType(), True),
  StructField("count", LongType(), True),
])

events = (
  spark.readStream.schema(input_schema)
    .option("maxFilesPerTrigger", 10)
    .json("/databricks-datasets/structured-streaming/events")
)

q = (
  events.groupBy("action")
    .transformWithState(
      statefulProcessor=AsyncCountProcessor(),
      outputStructType=output_schema,
      outputMode="Update",
      timeMode="None",
    )
    .writeStream.format("memory")
    .queryName("async_counts")
    .trigger(availableNow=True)
    .start()
)

q.awaitTermination()

查詢完成後,請查看以下程式碼中每種動作類型的執行次數:

display(spark.sql("SELECT action, MAX(count) AS count FROM async_counts GROUP BY action ORDER BY action"))

非同步狀態與計時器操作

在 中 AsyncStatefulProcessor,狀態變數與計時器操作讀取或寫入值是非同步的。 大多數這些運算會回傳一個結果,你用 await來取得。 回傳集合的操作則會回傳一個非同步迭代器,你用 async for來處理。 關於 Python 中非同步迭代器的介紹async/await,請參閱 Python asyncio 文件。

下表列出了回傳單一結果的操作,你可以用:await

Class 使用 await
AsyncValueState exists、get、update、clear
AsyncMapState exists、getValue、containsKey、updateValue、removeKey、clear
AsyncListState exists、、 put、 appendValue、 appendList、 clear
AsyncStatefulProcessorHandle registerTimer、deleteTimer

下表列出了回傳非同步迭代器的操作,你可以用:async for

Class 使用 async for
AsyncMapState iterator、keys、values
AsyncListState get
AsyncStatefulProcessorHandle listTimers

範例:async for

例如,要讀取 中的值AsyncListState,請依以下程式碼迭代:async for

total = 0
async for value in self.items.get():
  total += value[0]

建立狀態物件並刪除狀態變數的方法保持同步: getValueState、 getMapState、 getListState、 deleteIfExists。

關於每種狀態類型的描述,請參見 自訂狀態類型。

以非同步程式設計模式進行優化

當邏輯等待外部操作(如網路請求)時,非同步處理非常有用。 與其連續等待每個請求,不如同時 asyncio 執行這些請求,減少閒置時間。

範例:執行並行請求 asyncio.gather

以下範例使用 asyncio.gather 同時觸發所有每列 HTTP 請求並等待它們完成,然後將最高分數儲存在狀態中。 以下程式碼定義處理器:

import asyncio
import aiohttp
from pyspark.sql import Row
from pyspark.sql.streaming import AsyncStatefulProcessor

class HttpScoreRowGatherProcessor(AsyncStatefulProcessor):
  async def init(self, handle):
    self._score_state = handle.getValueState("last_score", "score double")
    self._session = aiohttp.ClientSession()

  async def _fetch_score(self, row) -> float:
    async with self._session.get(
      f"https://api.example.com/score/{row.event_id}"
    ) as resp:
      return (await resp.json())["score"]

  async def handleInputRows(self, key, rows, timerValues):
    user_id = key[0]
    scores = await asyncio.gather(*[self._fetch_score(row) for row in rows])

    max_score = max(scores)
    await self._score_state.update((max_score,))
    yield Row(user_id=user_id, score=max_score)

  async def close(self):
    await self._session.close()