Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Автономная таблица потоковой передачи — это таблица, зарегистрированная в каталоге Unity с дополнительной поддержкой потоковой или добавочной обработки данных, определенной за пределами конвейера Lakeflow. Конвейер создается автоматически для каждой потоковой таблицы. Таблицы потоковой передачи можно использовать для добавочной загрузки данных из Kafka и облачного хранилища объектов.
Вы можете создавать и обновлять автономные таблицы потоковой передачи из хранилища SQL Databricks или записной книжки, работающей на бессерверных общих вычислительных ресурсах. Дополнительные сведения о различиях между двумя вариантами вычислений см. в разделе "Требования" для автономных конвейеров.
Чтобы создавать и обновлять автономные потоковые таблицы с помощью Python в записной книжке, см. статью Использование Python с автономными конвейерами.
Замечание
Сведения о том, как использовать таблицы Delta Lake в качестве источников и приемников потоковых данных, см. в разделе Потоковое чтение и запись таблиц Delta Lake.
Требования
Сведения о параметрах вычислений, разрешениях и других требованиях для создания, обновления и запроса автономных таблиц потоковой передачи см. в разделе "Требования к автономным конвейерам".
Создание потоковых таблиц
Потоковая таблица определяется SQL-запросом в Databricks SQL. При создании потоковой таблицы данные из исходных таблиц используются для её формирования. После этого вы обновляете таблицу, как правило, по расписанию, чтобы извлечь все добавленные данные в исходных таблицах, чтобы добавить в потоковую таблицу.
Когда вы создаёте потоковую таблицу, вы считаетесь её владельцем.
Чтобы создать потоковую таблицу из существующей таблицы, используйте инструкциюCREATE STREAMING TABLE, как показано в следующем примере:
CREATE OR REFRESH STREAMING TABLE sales
SCHEDULE EVERY 1 hour
AS SELECT product, price FROM STREAM raw_data;
В этом случае потоковая таблица sales создается из определенных raw_data столбцов таблицы с расписанием обновления каждый час. Используемый запрос должен быть потоковым запросом. Используйте ключевое слово STREAM для применения семантики потоковой передачи при чтении из источника.
Вычисления, используемые для обновления
При создании потоковой таблицы с помощью CREATE OR REFRESH STREAMING TABLE инструкции обновление и начальное заполнение данных начинается немедленно. Эти операции не используют вычислительные ресурсы хранилища Databricks SQL. Вместо этого потоковые таблицы используют бессерверные конвейеры для создания и обновления. Выделенный бессерверный конвейер автоматически создается и управляется системой для каждой потоковой таблицы.
Загрузка файлов с помощью автозагрузчика
Чтобы создать потоковую таблицу из файлов в томе, используйте автозагрузчик. Используйте автозагрузчик для большинства задач приема данных из облачного хранилища объектов. Автозагрузчик и конвейеры предназначены для добавочной и идемпотентной загрузки постоянно растущих данных по мере поступления в облачное хранилище.
Чтобы использовать автозагрузчик в Databricks SQL, используйте функцию read_files . В следующих примерах показано использование автозагрузчика для чтения тома JSON-файлов в потоковую таблицу:
CREATE OR REFRESH STREAMING TABLE sales
SCHEDULE EVERY 1 hour
AS SELECT * FROM STREAM read_files(
"/Volumes/my_catalog/my_schema/my_volume/path/to/data",
format => "json"
);
Для чтения данных из облачного хранилища можно также использовать автозагрузчик:
CREATE OR REFRESH STREAMING TABLE sales
SCHEDULE EVERY 1 hour
AS SELECT *
FROM STREAM read_files(
'abfss://myContainer@myStorageAccount.dfs.core.windows.net/analysis/*/*/*.json',
format => "json"
);
Дополнительные сведения о автозагрузчике см. в разделе "Что такое автозагрузчик?". Дополнительные сведения об использовании автозагрузчика в SQL см. в статье "Загрузка данных из хранилища объектов".
Прием потоковой передачи из других источников
Пример приема из других источников, включая Kafka, см. в разделе "Загрузка данных в конвейерах".
Применение отслеживания измененных данных (CDC) с автоматическими потоками CDC
Используйте оператор FLOW AUTO CDC для обработки записей регистрации изменений данных (CDC) из источника в потоковую таблицу. Ранее инструкция MERGE INTO часто использовалась для обработки записей CDC в Azure Databricks.
MERGE INTO Однако может привести к неправильным результатам из-за неупорядоченных записей или требуется сложная логика для повторного упорядочивания записей. См. сведения об отслеживании и моментальных снимках измененных данных.
AUTO CDC упрощает CDC путем автоматической обработки несортированных записей. Вы указываете ключи для идентификации записей, столбец последовательности для упорядочивания, и нужно выбрать, хранить результаты как SCD типа 1 (прямые обновления) или SCD типа 2 (историческое отслеживание).
В следующем примере создается потоковая таблица, которая применяет изменения CDC с помощью SCD типа 1:
CREATE OR REFRESH STREAMING TABLE target
FLOW AUTO CDC
FROM stream(cdc_data.users)
KEYS (userId)
SEQUENCE BY sequenceNum
STORED AS SCD TYPE 1;
В следующем примере для сохранения журнала изменений используется тип SCD 2:
CREATE OR REFRESH STREAMING TABLE target
FLOW AUTO CDC
FROM stream(cdc_data.users)
KEYS (userId)
APPLY AS DELETE WHEN operation = "DELETE"
SEQUENCE BY sequenceNum
COLUMNS * EXCEPT (operation, sequenceNum)
STORED AS SCD TYPE 2;
Полные сведения о параметрах и поведении auto CDC см. в api-интерфейсах AUTO CDC: упрощение отслеживания измененных данных с помощью конвейеров. Полный справочник по синтаксису см. в разделе CREATE STREAMING TABLE.
Примените выборочную пакетную замену в потоках REPLACE WHERE
Используйте предложение FLOW REPLACE WHERE, чтобы пересчитать и перезаписать определённое подмножество потоковой таблицы без повторной обработки всей истории таблицы.
REPLACE WHERE потоки хорошо подходят для инкрементной пакетной обработки операций соединения и агрегаций, поздно поступающих данных, повторной обработки вышестоящих данных, эволюции схемы и обратного заполнения данных.
Подробные сведения о процессах REPLACE WHERE, включая требования, переопределения предикатов и инкрементальное обновление, см. в статье REPLACE WHERE потоки для автономных потоковых таблиц.
Примените частичную замену снимков с помощью REPLACE USING flows
Это важно
Процессы REPLACE USING находятся в бета-версии.
Используйте предложение FLOW REPLACE USING, чтобы поддерживать синхронизацию потоковой таблицы с потоком частичных моментальных снимков. При каждом обновлении поток REPLACE USING заменяет все строки, соответствующие указанным ключевым столбцам, и оставляет все остальные строки неизменными. Столбец SEQUENCE BY упорядочивает обновления так, что обновление с наибольшим номером последовательности для ключа всегда имеет приоритет, даже если обновления поступают не по порядку. Рассмотрим пример.
CREATE OR REFRESH STREAMING TABLE payments_current
FLOW REPLACE USING (payment_id) SEQUENCE BY payment_date BY NAME
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);
BY NAME является обязательным. Сопоставляет столбцы по именам, а не по их порядку.
REPLACE USING ведёт себя с отдельными потоковыми таблицами так же, как и в конвейерах Lakeflow. О том, как это работает, о последовательности, ожиданиях, ограничениях и примерах см. Частичная замена снимка с помощью потоков REPLACE USING. Следующие различия применимы к отдельным стриминговым таблицам:
- Определите поток в SQL. Создайте поток REPLACE USING со встроенной инструкцией SQL
FLOW REPLACE USINGвCREATE OR REFRESH STREAMING TABLE. Отдельная инструкцияCREATE FLOW— это конструкция конвейера Lakeflow и не используется для автономных потоковых таблиц. - Вычисления управляются за вас. Изолированные потоковые таблицы выполняются в бессерверных конвейерах, управляемых системой, и требуют Databricks Runtime 18.2 и выше. Вы не выбираете между классическими и бессерверными вычислениями.
Загружайте только новые данные
По умолчанию read_files функция считывает все существующие данные в исходной папке во время создания таблицы, а затем обрабатывает только что поступающие записи при каждом обновлении.
Чтобы избежать приема данных, которые уже существуют в исходной папке на момент создания таблицы, установите параметр includeExistingFiles на значение false. Это означает, что только данные, поступающие в папку после создания таблицы, обрабатываются. Рассмотрим пример.
CREATE OR REFRESH STREAMING TABLE sales
SCHEDULE EVERY 1 hour
AS SELECT *
FROM STREAM read_files(
'/path/to/files',
includeExistingFiles => false
);
Версия для выполнения
Потоковые таблицы всегда работают на последней версии Databricks SQL runtime. Свойство pipelines.channel таблицы, ранее использовавшееся для выбора preview канала current во время выполнения, больше не поддерживается и не действует. Если существующее определение включает это свойство, оно безопасно игнорируется, и вам не нужно его удалять.
Скрытие конфиденциальных данных
Таблицы стриминга можно использовать для скрытия конфиденциальных данных от пользователей, обращающихся к таблице. Одним из способов является определение запроса, чтобы он полностью исключил конфиденциальные столбцы или строки. Кроме того, можно применить маски столбцов или фильтры строк на основе разрешений пользователя запроса. Например, можно скрыть tax_id столбец для пользователей, которые не находятся в группе HumanResourcesDept. Для этого используйте синтаксис ROW FILTER и MASK во время создания потоковой таблицы. Дополнительные сведения см. в статьях "Фильтры строк" и маски столбцов.
Обновление потоковой таблицы
Потоковые таблицы автоматически создают и используют бессерверные конвейеры для обработки обновления. Обновление управляется конвейером и обновление отслеживается хранилищем Databricks SQL, используемым для создания потоковой таблицы. Таблицы потоковой обработки могут быть обновлены с помощью конвейера, выполняющегося по расписанию.
Даже если у вас есть запланированное обновление, вы можете в любое время вызвать обновление вручную. Обновления обрабатываются тем же конвейером, который был автоматически создан вместе с потоковой таблицей.
Чтобы обновить таблицу потоковой передачи, выполните следующее:
REFRESH STREAMING TABLE sales;
Вы можете проверить состояние последнего обновления с помощью DESCRIBE TABLE EXTENDED.
Замечание
Вам может потребоваться обновить таблицу потоковой передачи перед использованием запросов перемещения по времени.
Сведения о планировании обновления см. в разделе "Расписание обновлений". Запланированные обновления могут содержать уведомления об обновлении, и вы можете задать режим производительности для обновления.
Как работает обновление
Обновление потоковой таблицы оценивает только новые строки, поступающие после последнего обновления, и добавляет только новые данные.
Каждое обновление использует текущее определение таблицы потоковой передачи для обработки этих новых данных. Изменение описания таблицы потоковой передачи не приводит к автоматическому пересчету существующих данных. Если изменение несовместимо с существующими данными (например, изменение типа данных), следующее обновление завершается ошибкой.
В следующих примерах объясняется, как изменения в определении таблицы потоковой передачи влияют на поведение обновления:
- Удаление фильтра не обрабатывает ранее отфильтрованные строки.
- Изменение проекций столбцов не влияет на способ обработки существующих данных.
- Соединения со статическими моментальными снимками используют состояние моментального снимка во время начальной обработки. Поздно поступающие данные, которые бы совпали с обновленным снимком, игнорируются. Это может привести к упущению фактов, если измерения задерживаются.
- Изменение CAST существующего столбца приводит к ошибке.
Если данные изменяются таким образом, что их нельзя поддерживать в существующей потоковой таблице, вы можете выполнить полное обновление таблицы.
Полное обновление стриминговой таблицы
Полные обновления повторно обрабатывают все данные, доступные в источнике с помощью последнего определения. Не рекомендуется вызывать полные обновления в таких источниках, как Kafka, которые не хранят всю историю данных или имеют короткие периоды хранения, так как полное обновление обрезает существующие данные. Возможно, вы не сможете восстановить старые данные, если данные больше не доступны в источнике.
Рассмотрим пример.
REFRESH STREAMING TABLE sales FULL;
Планирование и мониторинг обновлений
Вы можете автоматически обновить таблицу потоковой передачи по расписанию или при изменении вышестоящих данных, а также настроить время ожидания обновления, уведомления и режимы производительности. См. статью "Расписание обновлений".
Управление доступом к потоковым таблицам
Потоковые таблицы поддерживают расширенные элементы управления доступом для совместного использования данных, не подвергая частные данные риску. Владелец таблицы потоковой передачи или пользователь с MANAGE привилегией может предоставить SELECT права другим пользователям. Пользователям с SELECT доступом к стриминговой таблице не требуется SELECT доступ к таблицам, на которые ссылается эта таблица. Этот контроль доступа обеспечивает общий доступ к данным при управлении доступом к базовым данным.
Вы также можете изменить владельца таблицы стриминга.
Предоставьте привилегии потоковой таблице
Чтобы предоставить доступ к таблице потоковой передачи, используйте инструкциюGRANT:
GRANT <privilege_type> ON <st_name> TO <principal>;
privilege_type может быть:
-
SELECT— пользователь может управлятьSELECTпотоковой таблицей. -
REFRESH— пользователь может управлятьREFRESHпотоковой таблицей. Обновления выполняются с помощью разрешений владельца.
В следующем примере создается потоковая таблица и предоставляются права выбора и обновления пользователям:
CREATE OR REFRESH STREAMING TABLE st_name AS SELECT * FROM source_table;
-- Grant read-only access:
GRANT SELECT ON st_name TO read_only_user;
-- Grant read and refresh access:
GRANT SELECT ON st_name TO refresh_user;
GRANT REFRESH ON st_name TO refresh_user;
Дополнительные сведения о предоставлении привилегий для защищаемых объектов каталога Unity см. в справочнике по привилегиям каталога Unity.
Отозвать привилегии из потоковой таблицы
Чтобы отозвать доступ к потоковой таблице, используйте инструкциюREVOKE:
REVOKE privilege_type ON <st_name> FROM principal;
Если SELECT права на исходную таблицу отзываются у владельца потоковой таблицы или любого другого пользователя, которому предоставлены привилегии MANAGE или SELECT на потоковую таблицу, или если исходная таблица удаляется, владелец потоковой таблицы или пользователь, имеющий доступ, всё равно может выполнять запросы к потоковой таблице. Однако происходит следующее поведение:
- Владелец таблицы потоковой передачи или другие пользователи, которые потеряли доступ к таблице потоковой передачи, больше
REFRESHне смогут использовать потоковую таблицу, и с течением времени потоковая таблица становится устаревшей. - Если автоматизировано с расписанием, следующий запланированный
REFRESHлибо не выполняется, либо завершится сбоем.
В следующем примере привилегия SELECT отзывается у read_only_user.
REVOKE SELECT ON st_name FROM read_only_user;
Изменение владельца потоковой таблицы
Пользователь с правами MANAGE на отдельную потоковую таблицу может назначить нового владельца через Catalog Explorer. Новый владелец может быть самим владельцем или сервисным принципалом, которому назначена роль пользователя сервисного принципала.
В рабочей области Azure Databricks щелкните
Каталог , чтобы открыть обозреватель каталогов.
Выберите таблицу потоковой передачи, которую требуется обновить.
В правой боковой панели в разделе "Сведения об этой потоковой таблице" найдите Владелец и щелкните
, чтобы изменить.
Замечание
Если вы получите сообщение о необходимости обновить владельца, изменив пользователя Run as в настройках конвейера, это означает, что потоковая таблица определена в конвейере Lakeflow, а не как отдельная таблица. Сообщение содержит ссылку на параметры конвейера, где можно изменить пользователя "Запуск от имени ".
Выберите нового владельца для потоковой таблицы.
Владельцы автоматически имеют
MANAGEиSELECTпривилегии на потоковые таблицы, которыми они владеют. Если вы назначаете служебный принципал в качестве владельца для собственной таблицы потоковой передачи, и у вас нет явных привилегийSELECTилиMANAGEна таблице потоковой передачи, это изменение приведет к полной потере вашего доступа к таблице потоковой передачи. В этом случае вам будет предложено явно предоставить эти привилегии.Выберите права "Предоставить управление " и "Предоставить SELECT " для предоставления прав на сохранение.
Нажмите кнопку "Сохранить", чтобы изменить владельца.
Владелец таблицы потокового вещания обновляется. Все будущие обновления выполняются, используя идентификацию нового владельца.
Когда владелец теряет права доступа к исходным таблицам
Если изменить владельца, а новый владелец не имеет доступа к исходным таблицам (или SELECT привилегии отзываются в базовых исходных таблицах), пользователи по-прежнему могут запрашивать потоковую таблицу. Тем не менее
- Они не могут
REFRESHвыполнять потоковую таблицу. - Следующее запланированное обновление таблицы потоковой передачи завершается сбоем.
Потеря доступа к исходным данным предотвращает обновление, но не немедленно запрещает чтение существующей потоковой таблицы.
окончательное удаление записей из потоковой таблицы
Это важно
Поддержка инструкции REORG с потоковыми таблицами доступна в общедоступной предварительной версии.
Замечание
- Использование инструкции
REORGс потоковой таблицей требует Databricks Runtime 15.4 и более поздних версий. - Хотя инструкцию
REORGможно использовать с любой таблицей потоковой передачи, она требуется только при удалении записей из таблицы потоковой передачи с включенными векторами удаления . Команда не действует при использовании со стриминговой таблицей, если векторы удаления не включены.
Чтобы физически удалить записи из базового хранилища потоковой таблицы с включенными векторами удаления, например в целях соблюдения требований GDPR, необходимо принять дополнительные меры, чтобы обеспечить выполнение операции VACUUM на данных потоковой таблицы.
Чтобы физически удалить записи из базового хранилища, выполните приведенные ниже действия.
- Обновите записи или удалите записи из таблицы потоковой передачи.
- Выполните выражение
REORGнад потоковой таблицей, указав параметрAPPLY (PURGE). Например,REORG TABLE <streaming-table-name> APPLY (PURGE);. - Дождитесь окончания срока хранения данных в потоковой таблице. Срок хранения данных по умолчанию составляет семь дней, но его можно настроить с помощью свойства таблицы
delta.deletedFileRetentionDuration. См. Настройка сохранения данных для запросов по временному перемещению. -
REFRESHпотоковая таблица. См . статью "Обновить потоковую таблицу". В течение 24 часов после операции задачи обслуживания конвейера, включаяREFRESHоперацию, необходимую для окончательногоVACUUMудаления записей, выполняются автоматически.
Мониторинг запусков с помощью журнала запросов
Вы можете использовать страницу журнала запросов для доступа к сведениям о запросах и профилям запросов, которые помогут определить плохое выполнение запросов и узких мест в конвейере, используемом для запуска обновлений потоковой таблицы. Общие сведения о типе информации, доступной в журналах запросов и профилях запросов, см. в разделе "Журнал запросов" и "Профиль запросов".
Это важно
Эта функция доступна в общедоступной предварительной версии. Администраторы рабочей области могут управлять доступом к этой функции на странице "Предварительные версии ". См. статью "Управление предварительными версиями Azure Databricks".
Все записи, связанные со стриминговыми таблицами, отображаются в журнале запросов. Вы можете использовать раскрывающийся список фильтра Statement, чтобы выбрать любую команду и проверить связанные запросы. За всеми CREATE операторами следует REFRESH инструкция, которая выполняется асинхронно в конвейере. Инструкции REFRESH обычно включают подробные планы запросов, которые предоставляют аналитические сведения о оптимизации производительности.
Чтобы получить доступ к REFRESH запросам в журнале запросов, выполните следующие действия.
- Щелкните
в левой боковой панели, чтобы открыть интерфейс истории запросов.
- Выберите флажок REFRESH, используя фильтр раскрывающегося списка заявления.
- Щелкните имя инструкции запроса, чтобы просмотреть сводные сведения, такие как длительность запроса и агрегированные метрики.
- Щелкните Просмотреть профиль запроса, чтобы открыть профиль запроса. Для получения информации о навигации в профиле запроса см. Профиль запроса.
- При необходимости можно использовать ссылки в разделе "Источник запросов", чтобы открыть связанный запрос или конвейер.
Вы также можете получить доступ к сведениям о запросах с помощью ссылок в редакторе SQL или из записной книжки, подключенной к хранилищу SQL.
Доступ к потоковым таблицам из внешних клиентов
Чтобы получить доступ к таблицам потоковой передачи из внешних клиентов Delta Lake или Iceberg, которые не поддерживают открытые API, можно использовать режим совместимости. Режим совместимости создает версию вашей таблицы для потоковой передачи в режиме только для чтения, к которой может получить доступ любой клиент Delta Lake или Iceberg.
Дополнительные ресурсы
- Декларативные конвейеры Spark
-
read_filesФункция с табличным значением -
read_kafkaФункция с табличным значением - CREATE STREAMING TABLE
- ALTER STREAMING TABLE
- Использование автономных материализованных представлений
- API-интерфейсы AUTO CDC: упрощение отслеживания измененных данных с помощью конвейеров