Обзор Zerobus Ingest

Zerobus Ingest — это push-ориентированный потоковый API, который записывает данные напрямую в таблицы Unity Catalog Delta в большом масштабе, без запуска шины сообщений. Это убирает средний слой, который многие команды ставят между своими продюсерами и домом на озере. Рабочий процесс состоит из двух этапов: создание таблицы, затем отправка данных в неё. Клиент «hello world» и рабочая нагрузка в масштабе петабайта выполняют по сути один и тот же код без инфраструктуры для управления.

Прием данных через шину сообщений направляет производителей через брокер и задание приема данных, прежде чем данные попадут в таблицы Delta, тогда как Zerobus Ingest напрямую подключает производителей к lakehouse.

Zerobus Ingest бессерверный, добавляя и удаляя ёмкость по мере изменения нагрузки. Он поглотил более 1 триллиона записей в одну таблицу менее чем за 24 часа (см. блог Ingesting the Milky Way: Petabyte-Scale with Zerobus Ingest ) и устанавливает рекорды за считанные секунды.

Advantages

Zerobus Ingest упрощает приём данных и при этом масштабируется для крупнейших рабочих нагрузок:

  • Простой по дизайну. Создайте таблицу, затем отправьте в неё данные — нет ни брокеров, ни разделов, ни конвейеров для управления. Вместо того чтобы направлять данные через шину сообщений и задачу загрузки, прежде чем они попадут в таблицу, источники данных записывают их напрямую в таблицу, поэтому промежуточных этапов меньше и компонентов, которые нужно поддерживать, тоже меньше.
  • Без сервера и эластичный. Zerobus Ingest включен по умолчанию и добавляет или убирает ёмкость при изменении нагрузки. Вы масштабируете систему, запуская больше продюсеров, а не переписывая приложение. Чтобы узнать, как это происходит, посмотрите, как масштабируется Zerobus Ingest.
  • Высокая пропускная нагрузка. Zerobus Ingest разработан для крупномасштабного ввода, поддерживая высокие скорости записи в одну таблицу.
  • Почти в режиме реального времени свежесть. Записи попадают в Delta за считанные секунды и готовы к запросу почти сразу после поступления.
  • Высокая степень параллелизма. Zerobus Ingest обрабатывает одновременные записи от тысяч клиентов в одну и ту же таблицу.

Если ваша цель — дом на озере, Zerobus Ingest — самый прямой путь. Другие инструменты Azure Databricks подходят для соседних потребностей и хорошо работают вместе с ними:

  • Для случаев, когда вы используете Kafka для поддержки потребителей, не входящих в Lakehouse, вам также может понадобиться копия данных, созданных в Lakehouse. Используйте управляемые стриминговые коннекторы , чтобы воспроизвести это.
  • Для данных, уже появляющихся в виде файлов в облачном хранилище, используйте Auto Loader.
  • Когда нужна операционная задержка ниже секунды на пути обработки, используйте режим реального времени.

Создайте таблицу, затем отправьте данные

Использование Zerobus Ingest — это так же просто, как создать таблицу и затем отправить в неё данные. Схема таблицы определяет, что должна содержать каждая запись. Сначала создайте целевую таблицу:

CREATE TABLE main.default.air_quality (
  device_name STRING,
  temp INT,
  humidity INT
);

Затем для загрузки записи достаточно всего нескольких строк кода:

from zerobus.sdk.sync import ZerobusSdk
from zerobus.sdk.shared import TableProperties

sdk = ZerobusSdk(SERVER_ENDPOINT, DATABRICKS_WORKSPACE_URL)

table_properties = TableProperties("main.default.air_quality")
stream = sdk.create_stream(CLIENT_ID, CLIENT_SECRET, table_properties)

stream.ingest_record_offset({"device_name": "sensor-1", "temp": 22, "humidity": 55})
# ingest more records...
stream.close()

Тот же код, который вы используете в среде разработки, можно масштабировать до промышленных нагрузок. Полный обзор смотрите в разделе «Использовать Zerobus Ingest».

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

  • IoT и телеметрия устройств: передавайте потоки данных с датчиков, транспортных средств и умных устройств из крупных распределённых парков напрямую в управляемые таблицы Delta.
  • Из локальной среды в облако: свяжите локальные и гибридные системы с lakehouse без развертывания промежуточной брокерской инфраструктуры. Для приватного подключения и настройки межсетевого экрана см. раздел «Сетевые аспекты».
  • События приложений и clickstream: отправляйте события из облачных и периферийных приложений для аналитики практически в реальном времени.
  • Захват изменённых данных (CDC): загрузка изменений строк из операционных систем в Delta.
  • Данные наблюдаемости: отправляйте следы, логи и метрики OpenTelemetry в принадлежащие вам Delta-таблицы. См. Ввод данных OpenTelemetry с помощью Zerobus Ingest.

Принцип работы

Продюсер открывает поток к Zerobus Ingest и передаёт записи в целевую таблицу Delta. Сервис проверяет каждую запись по схеме таблицы и делает её надёжной. Как только запись надежно сохранена, Zerobus Ingest быстро отправляет подтверждение, поэтому ваш производитель может продолжать отправлять записи, не дожидаясь подтверждения для каждой из них. Данные материализуются в таблице вскоре после этого, обычно в течение нескольких секунд. Динамичная архитектура Zerobus Ingest без разделов делает приём данных эластичным, благодаря чему её бессерверные вычислительные ресурсы масштабируются в соответствии с вашими рабочими нагрузками.

Как работает Zerobus Ingest: продюсеры отправляют записи на конечную точку Zerobus Ingest, которая проверяет, делает их надёжными, подтверждает и материализует их в таблицы Unity Catalog Delta

Для более глубокого объяснения потоков и масштабирования Zerobus Ingest см. концепции Zerobus Ingest. Для модели асинхронной клиентско-серверной коммуникации см. Асинхронная коммуникация.

Способы передачи данных

Zerobus Ingest — это одна конечная точка, поддерживающая несколько интерфейсов, так что вы можете выбрать лучший вариант для каждого производителя:

  • SDK над gRPC: высокопроизводительные потоковые клиенты на Python, Java, Rust, Go, TypeScript и (в бета-версии) C++ и C# / .NET. Лучше всего подходит для упорядоченного приёма данных в больших объёмах. См. Напишите клиенту.
  • REST API: интерфейс без сохранения состояния для легковесных или «разговорчивых» клиентов, таких как крупные парки периферийных устройств. См. Напишите клиенту.
  • OpenTelemetry (OTLP): настройте существующие коллекторы OpenTelemetry на отправку в Zerobus Ingest, чтобы передавать трассировки, журналы и метрики без какой-либо специальной интеграции. См. Ввод данных OpenTelemetry с помощью Zerobus Ingest.
  • Kafka-совместимые API (Beta): указать существующего производителя Apache Kafka на Zerobus Ingest, без Azure Databricks SDK. См. Как использовать Kafka-совместимые API с Zerobus Ingest.

Архитектура масштабирования Zerobus Ingest: источники отправляют записи Protocol Buffers (protobuf), JSON и Arrow через API gRPC, REST, OpenTelemetry и API, совместимые с Kafka; затем эти данные проходят через механизмы автомасштабирования и балансировки нагрузки и поступают в горизонтально масштабируемый пул узлов Zerobus без состояния, каждый из которых имеет журнал упреждающей записи и модуль записи в Lakehouse, который пакетно записывает данные в Delta-таблицу под управлением Unity Catalog

Все они записывают напрямую в таблицы Delta. Для полного сравнения и выбора смотрите протоколы API. Чтобы создать свой первый клиент, см. «Использование Zerobus Ingest».

Себестоимость

Плата за Zerobus Ingest тарифицируется по SKU «Automated Serverless». Цены доступны на странице цен Lakeflow Connect.

Мониторинг использования

Вы можете отслеживать свои расходы через таблицу системы оплачиваемого использования. См. справочную таблицу по использованию системы выставления счетов. Фильтрация использования Zerobus Ingest:

  • billing_origin_product = 'LAKEFLOW_CONNECT'
  • product_features.lakeflow_connect.zerobus_request_type определяет, как данные были получены: 'GRPC' (SDK streaming), 'HTTP' (REST) 'OTEL_GRPC' и 'OTEL_HTTP' (OpenTelemetry/OTLP), или 'KAFKA' (совместимые с Kafka API).

Дополнительные ресурсы