Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Important
Эта функция доступна в общедоступной предварительной версии.
Используйте структурированную потоковую передачу для записи в Lakebase со встроенной пакетной обработкой, автоматическими повторными попытками и управляемой рабочей областью проверкой подлинности.
Когда следует использовать приемник Lakebase
Используйте приемник Lakebase для потоковой записи с низкой задержкой в Lakebase. Этот приемник не требует, чтобы вы реализовывали пользовательские foreachBatch функции для пакетной обработки, управления подключениями и обработки ошибок.
Распространенные варианты использования:
- Обновите базы данных приложений в режиме реального времени для операционных панелей мониторинга или функций, доступных для клиентов.
- Синхронизация непрерывно изменяющихся данных, таких как агрегированные или отфильтрованные результаты потоковой передачи, в базу данных транзакций.
- Записывайте результаты запроса Structured Streaming в таблицу Lakebase с задержкой менее секунды с помощью режима реального времени.
Чтобы синхронизировать данные из Lakebase с таблицами Delta Lake в Lakehouse, о синхронизации в обратном направлении см. канал передачи данных об изменениях Lakebase.
Требования
- Databricks Runtime 18 и более поздние версии
- Классические вычисления с выделенными или стандартными режимами доступа.
- База данных 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, не поддерживаются.