Используйте ForEachBatch для записи данных в произвольные приемники в конвейерах

Приемник ForEachBatch обрабатывает поток как последовательность микропакетов. Каждую серию можно обрабатывать в Python с пользовательской логикой, похожей на структурированную потоковую передачу в Apache Spark foreachBatch. При использовании приемника ForEachBatch в конвейерах Lakeflow можно преобразовывать потоковые данные, выполнять их слияние или записывать их в одну или несколько целевых систем, которые изначально не имеют встроенной поддержки потоковой записи.

Приемник ForEachBatch предоставляет следующие функции:

  • Настраиваемая логика для каждого микропакета: ForEachBatch — это гибкий приемник потоковой передачи. Можно применять произвольные действия (например, слияние во внешнюю таблицу, запись в несколько мест назначения или выполнение вставок/обновлений) с помощью кода Python.
  • Полная поддержка обновления: Конвейеры управляют контрольными точками на основе каждого потока данных, поэтому контрольные точки сбрасываются автоматически при полном обновлении конвейера. При использовании объекта ForEachBatch вы отвечаете за управление сбросом данных в нижестоящем потоке, когда это происходит.
  • Поддержка каталога Unity: приемник ForEachBatch поддерживает все функции каталога Unity, такие как чтение или запись в томах или таблицах каталога Unity.
  • Ограниченное обслуживание: конвейер не отслеживает данные, записанные в приемник ForEachBatch, поэтому не может очистить эти данные. Вы несете ответственность за управление данными на последующих этапах.
  • Записи журнала событий: Журнал событий конвейера фиксирует создание и использование каждого приемника ForEachBatch. Если функция Python не сериализуема, в журнале событий появится предупреждение с дополнительными предложениями.

Замечание

  • Приемник ForEachBatch предназначен для потоковых запросов, таких как append_flow. Он не предназначен для конвейеров, работающих только с пакетами, или для AutoCDC семантики.
  • Приемник ForEachBatch, описанный на этой странице, предназначен для конвейеров. Структурированная потоковая передача Apache Spark также поддерживает foreachBatch. Сведения о структурированной потоковой передаче foreachBatchсм. в разделе Use foreachBatch для записи в произвольные приемники данных.

Когда использовать приемник ForEachBatch

Используйте приемник ForEachBatch всякий раз, когда конвейер требует функциональных возможностей, которые недоступны через встроенный формат приемника, например deltaили kafka. Типичные варианты использования:

  • Слияние или вставка в таблицу Delta Lake: запуск пользовательской логики слияния для каждого микро-пакета (например, обработка обновленных записей).
  • Запись в несколько или неподдерживаемых назначений: записывайте результаты каждого пакета в несколько таблиц или в внешние системы хранения, которые не поддерживают потоковую запись (например, определенные приемники JDBC).
  • Применение настраиваемой логики или преобразований: Манипулирование данными в Python напрямую (например, с использованием специализированных библиотек или сложных преобразований).

Сведения о встроенных приемниках или создании пользовательских приемников с помощью Python см. в разделе "Приемники" в конвейерах Lakeflow.

@dp.foreach_batch_sink() См. справочник по Python API: foreach_batch_sink.

Полное обновление

Так как ForEachBatch использует потоковый запрос, конвейер отслеживает каталог контрольных точек для каждого потока. При полном обновлении:

  • Каталог контрольной точки сбрасывается.
  • Функция получения данных (foreach_batch_sink UDF) распознает новый batch_id цикл, начиная с 0.
  • Данные в целевой системе не очищаются конвейером автоматически (так как конвейер не знает, где записываются данные). Если требуется сценарий чистого листа, необходимо вручную удалить или обрезать внешние таблицы или места хранения, заполняемые функцией ForEachBatch.

Использование функций каталога Unity

Все существующие возможности каталога Unity в Структурированной потоковой передаче foreach_batch_sink Spark остаются доступными.

Это включает запись в управляемые или внешние таблицы каталога Unity. Микропакеты можно записывать в управляемые или внешние таблицы каталога Unity точно так же, как и в любом задании структурированной потоковой передачи Apache Spark.

Записи журнала событий

При создании приемника ForEachBatch, событие SinkDefinition, связанное с "format": "foreachBatch", добавляется в журнал событий конвейера.

Это позволяет отслеживать использование приемников ForEachBatch и просматривать предупреждения о приемнике.

Использование с Databricks Connect

Если указанная функция не сериализуема (важное требование для Databricks Connect), журнал событий содержит WARN запись, которая рекомендует упростить или рефакторинг кода, если требуется поддержка Databricks Connect.

Например, если вы используете dbutils для получения параметров в UDF ForEachBatch, можно вместо этого получить аргумент перед его использованием в UDF:

# Instead of accessing parameters within the UDF...
def foreach_batch(df, batchId):
  value = dbutils.widgets.get ("X") + str (i)

# ...get the parameters first, and use them within the UDF:
argX = dbutils.widgets.get ("X")

def foreach_batch(df, batchId):
  value = argX + str (i)

Лучшие практики

  1. Делайте функцию ForEachBatch краткой: избегайте многопоточности, тяжелых библиотечных зависимостей, или больших операций с данными в памяти. Сложная или с состоянием логика может привести к ошибкам сериализации или проблемам с производительностью.
  2. Следите за папкой контрольных точек: Для потоковых запросов конвейер управляет контрольными точками на уровне потока, а не на уровне приемника. При наличии нескольких потоков в конвейере каждый поток имеет собственный каталог контрольных точек.
  3. Проверьте внешние зависимости: если вы используете внешние системы или библиотеки, убедитесь, что они установлены на всех узлах кластера или в контейнере.
  4. Помните, что Databricks Connect: если ваша среда может перейти в Databricks Connect в будущем, убедитесь, что ваш код можно сериализовать и не зависит от dbutils внутри foreach_batch_sink UDF.

Ограничения

  • No для ForEachBatch: поскольку пользовательский код Python может записывать данные в любом месте, конвейер не может очистить или отслеживать эти данные. Необходимо обеспечивать соблюдение собственных политик управления данными и хранения данных для мест назначения, в которые вы записываете данные.
  • Микропакетные метрики: конвейеры собирают метрики потоковой обработки, но некоторые сценарии могут вызывать неполные или необычные метрики при использовании оператора ForEachBatch. Это обусловлено базовой гибкостью ForEachBatch, что затрудняет отслеживание потока данных и строк для системы.
  • Поддержка записи в несколько целевых расположений без многократного чтения: Некоторые клиенты могут использовать ForEachBatch, чтобы один раз считать данные из источника, а затем записать их в несколько целевых расположений. Для этого необходимо включить df.persist и df.cache в функцию ForEachBatch. Используя эти параметры, Azure Databricks пытается считывать данные только один раз. Без этих параметров запрос приводит к нескольким считываниям. Это не включается в приведенные ниже примеры кода.
  • Использование с Databricks Connect: если конвейер выполняется в Databricks Connect, foreachBatch определяемые пользователем функции (UDF) должны быть сериализуемыми и не могут использоваться dbutils. Конвейер вызывает предупреждения, если он обнаруживает несериализируемый UDF, но не завершает конвейер.
  • Несериализируемая логика: код, ссылающийся на локальные объекты, классы или неуправляемые ресурсы, может нарушить контексты Databricks Connect. Используйте чистые модули Python и убедитесь, что ссылки (например, dbutils) не используются, если Databricks Connect является обязательным требованием.

Примеры

Пример простого синтаксиса

from pyspark import pipelines as dp

# Create a ForEachBatch sink
@dp.foreach_batch_sink(name = "my_foreachbatch_sink")
def feb_sink(df, batch_id):
  # Custom logic here. You can perform merges,
  # write to multiple destinations, etc.
  return

# Create source data for example:
@dp.table()
def example_source_data():
  return spark.range(5)

# Add sink to an append flow:
@dp.append_flow(
    target="my_foreachbatch_sink",
)
def my_flow():
  return spark.readStream.format("delta").table("example_source_data")

Использование примеров данных для простого конвейера

В этом примере используется пример такси Нью-Йорка. Предполагается, что администратор рабочей области включил каталог общедоступных наборов данных Databricks. Для приемника измените my_catalog.my_schema на каталог и схему, к которым у вас есть доступ.

from pyspark import pipelines as dp
from pyspark.sql.functions import current_timestamp

# Create foreachBatch sink
@dp.foreach_batch_sink(name = "my_foreach_sink")
def my_foreach_sink(df, batch_id):
    # Custom logic here. You can perform merges,
    # write to multiple destinations, etc.
    # For this example, we are adding a timestamp column.
    enriched = df.withColumn("processed_timestamp", current_timestamp())
    # Write to a Delta location
    enriched.write \
      .format("delta") \
      .mode("append") \
      .saveAsTable("my_catalog.my_schema.trips_sink_delta")
    # Return is optional here, but generally not used for the sink
    return

# Create an append flow that reads sample data,
# and sends it to the ForEachBatch sink
@dp.append_flow(
    target="my_foreach_sink",
)
def taxi_source():
  df = spark.readStream.table("samples.nyctaxi.trips")
  return df

Запись в несколько направлений

В этом примере выполняется запись в несколько мест назначения. Он демонстрирует, как использование txnVersion и txnAppId делает операции записи в таблицы Delta Lake идемпотентными. Дополнительные сведения см. в разделе "Использование foreachBatch для записи идемпотентной таблицы".

Предположим, что мы записываем в две таблицы, table_a и table_b, и что в пакете запись в table_a выполняется успешно, в то время как запись в table_b завершается ошибкой. При повторном запуске пакета пара (txnVersion, txnAppId) позволит Delta игнорировать дубликат записи table_a и записывать пакет только в table_b.

from pyspark import pipelines as dp

app_id = "my-app-name" # different applications that write to the same table should have unique txnAppId

# Create the ForEachBatch sink
@dp.foreach_batch_sink(name="user_events_feb")
def user_events_handler(df, batch_id):
    # Optionally do transformations, logging, or merging logic
    # ...

    # Write to a Delta table
    df.write \
     .format("delta") \
     .mode("append") \
     .option("txnVersion", batch_id) \
     .option("txnAppId", app_id) \
     .saveAsTable("my_catalog.my_schema.example_table_1")

    # Also write to a JSON file location
    df.write \
      .format("json") \
      .mode("append") \
      .option("txnVersion", batch_id) \
      .option("txnAppId", app_id) \
      .save("/tmp/json_target")
    return

# Create source data for example
@dp.table()
def example_source():
  return spark.range(5)


# Create the append flow, and target the ForEachBatch sink
@dp.append_flow(target="user_events_feb", name="user_events_flow")
def read_user_events():
    return spark.readStream.format("delta").table("example_source")

С использованием spark.sql()

Вы можете использовать spark.sql() в приемнике ForEachBatch, как показано в следующем примере.

from pyspark import pipelines as dp
from pyspark.sql import Row

@dp.foreach_batch_sink(name = "example_sink")
def feb_sink(df, batch_id):
  df.createOrReplaceTempView("df_view")
  df.sparkSession.sql("MERGE INTO target_table AS tgt " +
            "USING df_view AS src ON tgt.id = src.id " +
            "WHEN MATCHED THEN UPDATE SET tgt.id = src.id * 10 " +
            "WHEN NOT MATCHED THEN INSERT (id) VALUES (id)"
          )
  return

# Create target delta table
spark.range(5).write.format("delta").mode("overwrite").saveAsTable("target_table")

# Create source table
@dp.table()
def src_table():
  return spark.range(5)

@dp.append_flow(
    target="example_sink",
)
def example_flow():
  return spark.readStream.format("delta").table("source_table")

Слияние с внешней таблицей Delta Lake

from pyspark import pipelines as dp
from pyspark.sql.functions import col
from delta.tables import DeltaTable

@dp.foreach_batch_sink(name = "external_merge_feb")
def foreachBatchFunc(df, batchId):
  out = DeltaTable.forName(df.sparkSession, $table)
  out.alias("target") \
    .merge(df.alias("source"), "source.value = target.value") \
    .whenMatchedUpdateAll() \
    .whenNotMatchedInsertAll() \
    .whenNotMatchedBySourceDelete() \
    .execute()

@dp.update_flow(
    target="external_merge_feb",
    name="merge_flow"
)
def read_data():
    return (
        spark.readStream.format("delta")
        .load("/tmp/source_delta_table")
        .filter(col("value").isNotNull())
    )

Часто задаваемые вопросы (FAQ)

Можно ли использовать dbutils в приемнике ForEachBatch?

Если вы планируете запустить конвейер в среде, отличной от Databricks Connect, dbutils может работать. Однако, если вы используете Databricks Connect, dbutils недоступен в функции foreachBatch. Конвейер может порождать предупреждения, если он обнаруживает использование dbutils, чтобы помочь избежать сбоев.

Можно ли использовать несколько потоков данных с одним приемником ForEachBatch?

Да. Можно определить несколько потоков (с @dp.append_flow), которые предназначены для одного и того же имени приемника, но каждый из них поддерживает свои собственные контрольные точки.

Обрабатывает ли конвейер хранение или очистку данных для моего целевого объекта?

Нет. Так как приемник ForEachBatch может записывать данные в любое произвольное расположение или систему, конвейер не может автоматически управлять или удалять данные в этом целевом объекте. Эти операции необходимо обрабатывать как часть пользовательского кода или внешних процессов.

Как устранить ошибки сериализации или сбои в функции ForEachBatch?

Просмотрите журналы драйверов кластера или журналы событий конвейера. Для проблем сериализации, связанных с Spark Connect, убедитесь, что функция зависит только от сериализуемых Python объектов и не ссылается на запрещенные объекты (например, открытые дескрипторы файлов или dbutils).