Подключение к Lakebase

Important

Эта функция доступна в общедоступной предварительной версии.

Используйте структурированную потоковую передачу для записи в Lakebase со встроенной пакетной обработкой, автоматическими повторными попытками и управляемой рабочей областью проверкой подлинности.

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

Используйте приемник Lakebase для потоковой записи с низкой задержкой в Lakebase. Этот приемник не требует, чтобы вы реализовывали пользовательские foreachBatch функции для пакетной обработки, управления подключениями и обработки ошибок.

Распространенные варианты использования:

  • Обновите базы данных приложений в режиме реального времени для операционных панелей мониторинга или функций, доступных для клиентов.
  • Синхронизация непрерывно изменяющихся данных, таких как агрегированные или отфильтрованные результаты потоковой передачи, в базу данных транзакций.
  • Записывайте результаты запроса Structured Streaming в таблицу Lakebase с задержкой менее секунды с помощью режима реального времени.

Чтобы синхронизировать данные из Lakebase с таблицами Delta Lake в Lakehouse, о синхронизации в обратном направлении см. канал передачи данных об изменениях Lakebase.

Требования

Соединение с базой данных

Приемник Lakebase поддерживает следующие методы подключения:

Таблицы Lakebase, зарегистрированные в каталоге Unity

Для таблиц Lakebase, зарегистрированных в Unity Catalog, коннектор автоматически управляет учетными данными и использует идентификационные данные пользователя или сервисного субъекта, выполняющего запрос. Если таблица не существует, соединитель создает таблицу.

Сведения о регистрации базы данных Lakebase в каталоге Unity см. в разделе "Регистрация базы данных Lakebase" в каталоге Unity.

Чтобы записать данные в таблицу Lakebase, используйте метод .toTable() с полным именем таблицы catalog.schema.table. В следующем примере показаны необходимые параметры, а также необязательный upsertkey параметр:

Python

(df.writeStream
  .outputMode("update")
  .option("upsertkey", "<primary-key-column>")  # Optional. Inferred from the table's primary key if omitted.
  .option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
  .toTable("<catalog>.<schema>.<table>")
)

Scala

df.writeStream
  .outputMode("update")
  .option("upsertkey", "<primary-key-column>")  // Optional. Inferred from the table's primary key if omitted.
  .option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
  .toTable("<catalog>.<schema>.<table>")

Замените следующие заполнители:

  • <catalog>.<schema>.<table>: полное имя целевой таблицы. Это catalog каталог каталога Unity, созданный при регистрации базы данных Lakebase, см. статью "Регистрация базы данных Lakebase в каталоге Unity". Если таблица не существует, соединитель создает ее.
  • <primary-key-column>: необязательный параметр. Разделенный запятыми список столбцов, которые образуют ключ upsert, например id или user_id,event_type. Если не указано upsertkey, приемник выводит ключ из первичного ключа целевой таблицы. См. поведение Upsert.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Путь тома каталога Unity, в котором запрос сохраняет контрольную точку. Вы также можете использовать универсальный код ресурса (URI) хранилища облачных объектов. Расположение должно быть хранилищем, которое можно записать, а не локальный диск, и должно быть уникальным для каждого потокового запроса. Это не зависит от целевой таблицы. Смотрите контрольные точки структурированной потоковой передачи.

Сведения о необязательных конфигурациях, таких как batchsize и batchinterval, см. в разделе Параметры конфигурации.

Таблицы Lakebase, не зарегистрированные в каталоге Unity

Для таблиц Lakebase, не зарегистрированных в каталоге Unity, соединитель автоматически управляет учетными данными и использует удостоверение пользователя или субъекта-службы, выполняющего запрос. Если таблица не существует, соединитель создает таблицу.

Чтобы записать данные в таблицу Lakebase, используйте параметры endpoint и dbtable. В следующем примере также содержатся необязательные database и upsertkey параметры:

Python

(df.writeStream
  .format("postgresql")
  .outputMode("update")
  .option("endpoint", "<project-id>.<branch-id>.<endpoint-id>")
  .option("database", "<database>")  # Optional. Defaults to databricks_postgres.
  .option("dbtable", "<schema>.<table>")
  .option("upsertkey", "<primary-key-column>")  # Optional. Inferred from the table's primary key if omitted.
  .option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
  .start()
)

Scala

df.writeStream
  .format("postgresql")
  .outputMode("update")
  .option("endpoint", "<project-id>.<branch-id>.<endpoint-id>")
  .option("database", "<database>")  // Optional. Defaults to databricks_postgres.
  .option("dbtable", "<schema>.<table>")
  .option("upsertkey", "<primary-key-column>")  // Optional. Inferred from the table's primary key if omitted.
  .option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
  .start()

Замените следующие заполнители:

  • <project-id>.<branch-id>.<endpoint-id>: конечная точка Lakebase. Найдите все три значения в имени ресурса в меню "Получить идентификатор " вкладки "Вычисления ", которая имеет формат projects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id>. Ознакомьтесь с идентификаторами вычислений.
  • <database>: необязательный параметр. Имя целевой базы данных Postgres. По умолчанию — databricks_postgres. См. раздел "Управление базами данных".
  • <schema>.<table>: целевая таблица в schema.table формате. Если вы не укажете схему, приемник использует схему public. Используйте простые идентификаторы, которые начинаются с буквы или знака подчеркивания и содержат только буквы, цифры и знаки подчеркивания; идентификаторы в кавычках и специальные символы, например дефисы, не поддерживаются.
  • <primary-key-column>: необязательный параметр. Разделенный запятыми список столбцов, которые образуют ключ upsert, например id или user_id,event_type. Если не указано upsertkey, приемник выводит ключ из первичного ключа целевой таблицы. См. поведение Upsert.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Путь тома каталога Unity, в котором запрос сохраняет контрольную точку. Вы также можете использовать универсальный код ресурса (URI) хранилища облачных объектов. Расположение должно быть хранилищем, которое можно записать, а не локальный диск, и должно быть уникальным для каждого потокового запроса. Это не зависит от целевой таблицы. Смотрите контрольные точки структурированной потоковой передачи.

Сведения о необязательных конфигурациях, таких как batchsize и batchinterval, см. в разделе Параметры конфигурации.

Параметры конфигурации

Приёмник сообщает об ошибке при использовании нераспознанных параметров, JDBC_STREAMING_SINK_INVALID_OPTIONS.

Следующие параметры применяются ко всем методам подключения:

Ключ Default Description
batchinterval 100 milliseconds Optional. Максимальное время хранения строк в буфере перед очисткой. Например: "50 milliseconds".
batchsize 1000 Optional. Максимальное количество строк для каждой транзакции базы данных.
checkpointLocation None Required. Путь к каталогу контрольных точек, например том каталога Unity (/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>). Должен быть уникальным для каждого запроса. Смотрите контрольные точки структурированной потоковой передачи.
upsertkey None Optional. Список имен столбцов, разделенных запятыми, которые образуют ключ upsert. Например, "id" или "user_id,event_type". При указании upsertkeyстолбцы должны соответствовать первичному ключу таблицы или запрос завершается ошибкой. Если его опустить, приёмник автоматически использует первичный ключ. Дополнительные сведения см. в разделе "Поведение Upsert".

Таблицы Lakebase, не зарегистрированные в каталоге Unity

Следующие параметры применяются при подключении к таблице Lakebase, не зарегистрированной в каталоге Unity:

Ключ Default Description
database databricks_postgres Optional. Имя целевой базы данных PostgreSQL.
dbtable None Required. Имя целевой таблицы в формате schema.table. Если схема не указана, значение схемы по умолчанию равно public. Используйте простые идентификаторы, которые начинаются с буквы или символа подчеркивания и содержат только буквы, цифры и символы подчеркивания. Не заключайте имена таблиц и схем в кавычки; идентификаторы в кавычках и имена со специальными символами, например дефисами, не поддерживаются.
endpoint None Required. Конечная точка Lakebase в формате project_id.branch_id или project_id.branch_id.endpoint_id. Элемент endpoint_id необязателен; если его опустить и в ветви есть только одна конечная точка с доступом на чтение и запись, синк по умолчанию выбирает эту конечную точку.

Поведение операции Upsert

Если ключи upsert существуют — либо указаны с помощью upsertkey, либо определены приёмником на основе первичных ключей таблицы, — приёмник выполняет upsert в таблицу, используя синтаксис PostgreSQL INSERT INTO ... ON CONFLICT (<upsert_key>) DO UPDATE SET ....

Если ключи upsert отсутствуют, приёмник выполняет операции вставки. Режим вывода запроса не влияет на поведение операций upsert или insert.

Столбцы upsertkey должны:

  • Быть непустым подмножеством столбцов DataFrame.
  • Точно соответствует целевой таблице PRIMARY KEY . Если указанные столбцы не соответствуют первичному ключу, запрос завершается ошибкой.
  • Быть типами, которые можно сравнивать, например числовыми или строковыми. Чтобы предотвратить взаимоблокировку базы данных во время параллельной записи, приемник сортирует строки по ключу upsert в каждом пакете. Ключи Upsert не поддерживают сложные или типы структур.

Имена столбцов автоматически заключаются в стандартные для PostgreSQL двойные кавычки ", что позволяет обрабатывать зарезервированные слова и имена со смешанным регистром.

Имена таблиц и схем должны использовать простые идентификаторы, которые начинаются с буквы или подчеркивания и содержат только буквы, цифры и символы подчеркивания. Приёмник не поддерживает идентификаторы, заключённые в кавычки, а также специальные символы, такие как дефисы, в именах таблиц или схем.

Настройка производительности

Пакетная обработка и обратное давление

При выполнении любого условия выполняется очистка:

  • Размер буфера достигает batchsize строк, значение по умолчанию — 1000.
  • Возраст буфера превышает batchinterval, которое по умолчанию равно 100 milliseconds.

Если база данных не справляется со скоростью поступления данных, приёмник передаёт обратное давление вверх по потоку к источнику.

Руководство по задержке и пропускной способности:

  • Для рабочих нагрузок с низкой задержкой, использующих режим реального времени, уменьшите batchinterval, чтобы гарантировать меньшее максимальное время до сброса. См. концепции режима реального времени для концепций и примеры режима реального времени для примера кода.
  • Для рабочих нагрузок с высокой пропускной способностью увеличьте batchsize, чтобы уменьшить накладные расходы на каждую транзакцию.

Поведение подключения

Приемник использует пул соединений для исполнителей. По умолчанию каждая задача использует одно подключение к базе данных.

Databricks рекомендует использовать значение 1 задачи по умолчанию для каждого подключения. Если увеличить количество задач для каждого подключения, это может привести к конфликтам на соединении и увеличить задержки для соединений с высокой пропускной способностью.

Чтобы настроить отношение задач к подключениям, задайте конфигурацию spark.databricks.sql.streaming.jdbc.tasksPerConnection Spark. Если для целевой базы данных установлен низкий лимит подключений, уменьшите количество разделов shuffle или увеличьте spark.databricks.sql.streaming.jdbc.tasksPerConnection.

Приемник автоматически повторяет попытки при возникновении временных ошибок JDBC, включая сбои подключения, взаимоблокировки и ограничение частоты запросов. Если приемник исчерпает все повторные попытки, запрос завершается ошибкой.

Поддерживаемые триггеры и режимы вывода

Triggers

В этой таблице показана поддержка типов триггеров структурированной потоковой передачи:

Триггер Supported
realTime Yes
ProcessingTime Yes
AvailableNow Yes
Once Yes

Режимы вывода

В этой таблице показана поддержка режимов вывода структурированной потоковой передачи:

Режим вывода Supported
update Yes
append Yes. Поведение идентично update. Запрос выполняет операцию upsert, если целевая таблица имеет первичный ключ; в противном случае выполняется вставка. См. поведение Upsert.
complete No

Ограничения

  • Бессерверные вычислительные ресурсы и конвейеры Lakeflow не поддерживаются.
  • Только Lakebase поддерживается как целевой объект записи. Внешние базы данных, совместимые с PostgreSQL, не поддерживаются.