Используйте Zerobus Ingest

На этой странице описано, как принимать данные с помощью Zerobus Ingest в Lakeflow Connect.

Начните с Zerobus Ingest

Перед началом убедитесь, что Zerobus Ingest доступен в регионе вашего рабочего пространства. См. доступность приема данных.

  1. Получите URL-адрес инжеста Zerobus.
  2. Создайте или определите таблицу, в которую требуется принять данные.
  3. Создайте основной объект службы и предоставьте привилегии для таблицы.
  4. Подключите клиента или экспортера, чтобы начать отправку данных.

Выберите руководство для вашего варианта использования:

  • Импортирование собственных данных: используйте SDK или REST API Zerobus с определенной вами схемой. Следуйте инструкциям на этой странице.

  • Сбор данных OpenTelemetry: используйте стандартные SDK OpenTelemetry или сборщики для отправки трассировок, журналов и метрик в заранее определенные схемы таблиц. Полные инструкции см. в разделе Ingest OpenTelemetry data with Zerobus Ingest.

Выбор интерфейса

Zerobus Ingest поддерживает несколько интерфейсов, все они записываются непосредственно в таблицы Unity Catalog Delta. Иными словами:

  • SDK через gRPC: наивысшая стабильная пропускная способность, лучший вариант для производителей потоковых данных с большим объёмом трафика.
  • REST: без сохранения состояния, лучше всего подходит для большого парка лёгких или часто обменивающихся данными периферийных устройств.
  • OpenTelemetry (OTLP): для систем, уже излучающих следы, логи и метрики OpenTelemetry. См. Ввод данных OpenTelemetry с помощью Zerobus Ingest.

Для полного сравнения и выбора смотрите протоколы API. Помимо SDK, вы также можете выбрать формат записи (JSON, Protocol Buffers (protobuf) или Apache Arrow). См. Типы сообщений. Остальная часть этой страницы использует SDK и REST API.

Получите URL-адрес вашей рабочей области и Zerobus точку входа для приёма данных

URL-адрес рабочей области отображается в браузере при входе. Хотя полный URL-адрес соответствует формату https://<databricks-instance>.net/o=XXXXX, URL-адрес рабочей области состоит из всего, прежде чем /o=XXXXX. Например, учитывая следующий полный URL-адрес, можно определить URL-адрес рабочей области и идентификатор рабочей области.

  • Полный URL-адрес: https://abcd-teste2-test-spcse2.azuredatabricks.net/?o=2281745829657864#
  • URL-адрес рабочей области: https://abcd-teste2-test-spcse2.azuredatabricks.net
  • Идентификатор рабочей области: 2281745829657864

Конечная точка сервера зависит от рабочей области и региона:

  • Конечная точка сервера: <workspace-id>.zerobus.<region>.azuredatabricks.net

Чтобы найти регион рабочего пространства, откройте переключатель workspace в верхней панели навигации интерфейса Databricks. Регион отображается под именем каждого рабочего пространства (например, eastus). Его также можно найти в консоли учетной записи в разделе "Рабочие области".

Сведения о доступности по регионам см. в квотах Zerobus Ingest.

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

Определите целевую таблицу, в которую требуется принять данные. Чтобы создать целевую таблицу, выполните CREATE TABLE команду SQL. Например, создайте новую таблицу с именем unity.default.air_quality.

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

Zerobus Ingest может записывать как в управляемые Delta-таблицы, так и в потоковые таблицы, которые работают одинаково, с одинаковыми ограничениями и квотами.

Замечание

Для приема openTelemetry таблицы должны использовать предопределенные схемы для каждого типа сигнала (трассировки, журналы, метрики). См. статью "Создание целевых таблиц в каталоге Unity".

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

По умолчанию Zerobus Ingest отклоняет записи с полями, не соответствующими схеме целевой таблицы. Чтобы записать эти поля, а не потерять их, настройте спасательный столбец. См. колонку спасения Zerobus.

Создание субъекта-службы и предоставление разрешений

Специализированный сервисный идентификатор обеспечивает большую безопасность, чем персональные учетные записи. Для получения дополнительной информации о принципах сервиса и о том, как их использовать для аутентификации, см. раздел «Authorize service principal access to Azure Databricks with OAuth».

Вы можете создавать и управлять принципами сервисов программно с помощью API или SDK Azure Databricks REST, либо через интерфейс рабочего пространства, как описано ниже. Разрешения в конце этого раздела — это SQL, который можно запускать с любого клиента.

  1. Чтобы создать принципал сервиса, перейдите в раздел «Настройки>идентичности и доступа».

  2. В разделе "Субъекты-службы" выберите "Управление".

  3. Нажмите Добавить учетную запись службы.

  4. В окне "Добавление учетной записи службы" создайте новую учетную запись службы, нажав "Добавить новое".

  5. Создайте и сохраните идентификатор клиента и секрет клиента для субъекта-службы.

  6. Предоставьте необходимые разрешения для каталога, схемы и таблицы служебному принципалу.

    1. На странице главного сервиса перейдите на вкладку «Конфигурации ».
    2. Скопируйте идентификатор приложения (UUID).
    3. Используйте следующий SQL для предоставления разрешений, заменив пример UUID, каталог, имя схемы и имена таблиц, если требуется.
    GRANT USE CATALOG ON CATALOG <catalog> TO `<UUID>`;
    GRANT USE SCHEMA ON SCHEMA <catalog.schema> TO `<UUID>`;
    GRANT MODIFY, SELECT ON TABLE <catalog.schema.table_name> TO `<UUID>`;
    

Разработка клиента

Используйте пакет SDK Zerobus в предпочитаемом языке программирования или REST API для приема данных в целевую таблицу. SDK — это открытый код. Полная библиотека, специфическая языковая документация и дополнительные примеры см. репозиторий Zerobus SDK.

В приведенных ниже примерах используется ingest_record_offset, который сохраняет порядок отправки записей.

пакет SDK Python

требуется Python 3.9 или более поздней версии. SDK обеспечивает высокую пропускную способность и эффективный сетевой ввод-вывод благодаря асинхронной среде выполнения. Он поддерживает JSON (самый простой) и Protocol Buffers (рекомендуется для использования в продуктивной среде). SDK также поддерживает как синхронизирующую, так и асинхронную реализацию, а также методы запуска на основе смещения и будущего.

pip install databricks-zerobus-ingest-sdk

Пример JSON:

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

# See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
SERVER_ENDPOINT="https://1234567890123456.zerobus.eastus.azuredatabricks.net"
DATABRICKS_WORKSPACE_URL="https://adb-1234567890123456.12.azuredatabricks.net"
TABLE_NAME="main.default.air_quality"
CLIENT_ID="your-client-id"
CLIENT_SECRET="your-client-secret"

sdk = ZerobusSdk(
    SERVER_ENDPOINT,
    DATABRICKS_WORKSPACE_URL
)

table_properties = TableProperties(TABLE_NAME)
stream = sdk.create_stream(CLIENT_ID, CLIENT_SECRET, table_properties)

try:
    for i in range(1000):
        record_dict = {
            "device_name": f"sensor-{i}",
            "temp": 20 + i % 15,
            "humidity": 50 + i % 40
        }
        stream.ingest_record_offset(record_dict)
finally:
    stream.close()

В приведённых выше примерах используется метод на основе смещения ingest_record_offset без ожидания возвращаемого смещения. Чтобы узнать о доступных методах приёма данных, о том, когда следует дожидаться подтверждения устойчивости для смещения, и о том, как отслеживать ход выполнения с помощью функции обратного вызова для подтверждения, см. раздел «Блокировка сообщений и подтверждение».

Protocol Buffers: Для типобезопасной загрузки данных передайте protobuf-дескриптор в TableProperties (формат выбирается автоматически). Сгенерируйте схему из вашей таблицы с помощью generate_proto инструмента, скомпилируйте её с protoc, затем передайте скомпилированный дескриптор для создания потока.

Arrow Flight: Для колоночного или пакетного приёма данных Apache Arrow RecordBatch по тому же соединению gRPC см. раздел «Использование Arrow Flight с Zerobus Ingest». Требуется дополнение [arrow]: pip install "databricks-zerobus-ingest-sdk[arrow]" pyarrow.

Полные сведения о параметрах конфигурации, пакетном приеме и примерах буфера протокола см. в репозитории пакета SDK Python.

Rust SDK

Требуется Rust 1.70 или более поздней версии. SDK использует асинхронный ввод-вывод и gRPC для высокопроизводительного приёма данных. Он поддерживает JSON (самый простой) и Protocol Buffers (рекомендуется для использования в продуктивной среде).

Сначала импортируйте пакет.

cargo add databricks-zerobus-ingest-sdk

Или добавьте его в свой Cargo.toml.

[dependencies]
databricks-zerobus-ingest-sdk = "2.0.0" # Latest version at time of publication

Пример JSON:

  use databricks_zerobus_ingest_sdk::{JsonString, ZerobusSdk};
  use std::error::Error;

  // See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
  const DATABRICKS_WORKSPACE_URL: &str = "https://adb-1234567890123456.12.azuredatabricks.net";
  const SERVER_ENDPOINT: &str = "1234567890123456.zerobus.eastus.azuredatabricks.net";
  const TABLE_NAME: &str = "main.default.air_quality";
  const CLIENT_ID: &str = "your-client-id";
  const CLIENT_SECRET: &str = "your-client-secret";


  #[tokio::main]
  async fn main() -> Result<(), Box<dyn Error>> {
      let sdk_handle = ZerobusSdk::builder()
          .endpoint(SERVER_ENDPOINT)
          .unity_catalog_url(DATABRICKS_WORKSPACE_URL)
          .build()?;

      let mut stream = sdk_handle
          .stream_builder()
          .table(TABLE_NAME)
          .oauth(CLIENT_ID, CLIENT_SECRET)
          .json()
          .max_inflight_requests(100)
          .build()
          .await?;

      stream.ingest_record_offset(
        JsonString("{
          \"device_name\": \"sensor\",
          \"temp\": 22,
          \"humidity\": 55}".to_string())).await?;

      println!("Record ingested successfully");
      stream.close().await?;
      println!("Stream closed successfully");

      Ok(())
  }

Protocol Buffers: Для типобезопасного приема данных используйте Protocol Buffers через .compiled_proto(descriptor) в конструкторе потока вместо .json(), где descriptor — это prost_types::DescriptorProto. Сгенерируйте необходимые файлы с помощью инструмента generate_proto и импортируйте их в ваш проект. Arrow Flight: Для колоночного или пакетного приёма данных Apache Arrow RecordBatch по тому же соединению gRPC см. Использование Arrow Flight с Zerobus Ingest. Включите с помощью функции Cargo: cargo add databricks-zerobus-ingest-sdk --features arrow-flight

Полная документация, параметры конфигурации, пакетная приемка, generate_proto примеры инструментов и буфера протокола см. в репозитории пакета SDK Rust.

пакет SDK Java

требуется Java 8 или более поздней версии. SDK обеспечивает низкую задержку и эффективный сетевой ввод-вывод для высокоскоростного приёма данных. Он поддерживает JSON (самый простой) и Protocol Buffers (рекомендуется для использования в продуктивной среде).

Maven:

<dependency>
    <groupId>com.databricks</groupId>
    <artifactId>zerobus-ingest-sdk</artifactId>
    <version>0.2.0</version>
</dependency>

Пример JSON:

import com.databricks.zerobus.*;

public class ZerobusClient {

// See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
    private static final String SERVER_ENDPOINT =
        "https://1234567890123456.zerobus.eastus.azuredatabricks.net";
    private static final String DATABRICKS_WORKSPACE_URL =
        "https://adb-1234567890123456.12.azuredatabricks.net";
    private static final String TABLE_NAME = "main.default.air_quality";
    private static final String CLIENT_ID = "your-client-id";
    private static final String CLIENT_SECRET = "your-client-secret";

    public static void main(String[] args) throws Exception {
        ZerobusSdk sdk = new ZerobusSdk(
            SERVER_ENDPOINT,
            DATABRICKS_WORKSPACE_URL
        );

        ZerobusJsonStream stream = sdk.streamBuilder()
            .table(TABLE_NAME)
            .oauth(CLIENT_ID, CLIENT_SECRET)
            .json()
            .build()
            .join();

        try {
            for (int i = 0; i < 100; i++) {
                String record = String.format(
                    "{\"device_name\": \"sensor-%d\", \"temp\": 22, \"humidity\": 55}", i
                );
                stream.ingestRecordOffset(record);
            }
        } finally {
            stream.close();
        }
    }
}

Protocol Buffers: Для типобезопасной загрузки создайте ZerobusProtoStream с помощью streamBuilder() и .compiledProto(...). Сгенерируйте схему по вашей таблице, используя встроенную JAR-утилиту, затем скомпилируйте её с помощью protoc.

Arrow Flight: Для столбчатой или пакетной загрузки данных Apache Arrow RecordBatch через то же соединение gRPC см. раздел «Использование Arrow Flight с Zerobus Ingest».

Полная документация, параметры конфигурации, пакетная приемка и примеры буфера протокола см. в репозитории пакета SDK Java sdk.

Пакет SDK для GO

Требуется Go 1.21 или более поздняя версия. SDK обеспечивает высокую пропускную способность и производительность для потокового приема данных. Он поддерживает JSON (самый простой) и Protocol Buffers (рекомендуется для использования в продуктивной среде).

go get github.com/databricks/zerobus-sdk/go@latest

Пример JSON:

Для простоты ошибки игнорируются здесь. В рабочем коде всегда проверяйте ошибки.

package main

import (
	"fmt"

	zerobus "github.com/databricks/zerobus-sdk/go"
)

// See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
const (
	ServerEndpoint         = "https://1234567890123456.zerobus.eastus.azuredatabricks.net"
	DatabricksWorkspaceURL = "https://adb-1234567890123456.12.azuredatabricks.net"
	TableName              = "main.default.air_quality"
	ClientID               = "your-client-id"
	ClientSecret           = "your-client-secret"
)

func main() {
	sdk, _ := zerobus.NewZerobusSdk(
		ServerEndpoint,
		DatabricksWorkspaceURL,
	)
	defer sdk.Free()

	options := zerobus.DefaultStreamConfigurationOptions()
	options.RecordType = zerobus.RecordTypeJson

	stream, _ := sdk.CreateStream(
		zerobus.TableProperties{
			TableName: TableName,
		},
		ClientID,
		ClientSecret,
		options,
	)
	defer stream.Close()

	_, _ = stream.IngestRecordOffset(`{
		"device_name": "sensor-001",
		"temp": 20,
		"humidity": 60
	}`)

  fmt.Println("Record ingested successfully")

  _ = stream.Close()
  fmt.Println("Stream closed successfully")
}

Protocol Buffers: Для безопасного приема данных используйте Protocol Buffers с RecordTypeProto (по умолчанию) и укажите descriptorProto в свойствах таблицы. Создайте proto-файл, соответствующий схеме таблицы, и запустите generate_proto скрипт, чтобы импортировать файлы в проект.

Arrow Flight: Для колоночной или пакетной загрузки данных Apache Arrow RecordBatch по тому же соединению gRPC см. раздел «Использование Arrow Flight с Zerobus Ingest».

Полную документацию, параметры конфигурации, пакетную приемку, инструмент generate_proto и примеры Protocol Buffer см. в репозитории Go SDK.

C++ SDK

Important

Пакет SDK для C++ находится в бета-версии.

Требуется C++17 или более поздней версии. SDK обеспечивает нативную трансляцию gRPC, OAuth и автоматическое восстановление через интерфейс RAII на C++. Он поддерживает JSON для простых настроек и буферов протокола для рабочих нагрузок.

Пакет SDK поставляется в виде предварительно созданного пакета выпуска для каждой платформы, поэтому для его использования не требуется цепочка инструментов Rust. Скачайте пакет для вашей платформы (macOS, Linux (включая musl) или Windows) со страницы релизов, распакуйте его, а затем укажите в CMake путь к архиву FFI из этого пакета. Архив называется libzerobus_ffi.a в macOS и Linux и zerobus_ffi.lib на Windows:

# macOS and Linux
cmake -S cpp -B build \
  -DZEROBUS_FFI_LIBRARY="$PWD/lib/libzerobus_ffi.a" \
  -DZEROBUS_FFI_HEADER_DIR="$PWD/lib"
cmake --build build -j

В Windows (PowerShell) вместо этого укажите архив .lib:

cmake -S cpp -B build `
  -DZEROBUS_FFI_LIBRARY="$PWD/lib/zerobus_ffi.lib" `
  -DZEROBUS_FFI_HEADER_DIR="$PWD/lib"
cmake --build build -j

Чтобы создать пакет SDK из исходного заказа в собственном проекте CMake, добавьте его в качестве подкаталога и свяжите целевой объект. Вы также можете использовать FetchContent, чтобы получить его на этапе настройки. Это создает FFI из источника Rust, поэтому для него требуется цепочка инструментов Rust:

add_subdirectory(path/to/zerobus-sdk/cpp)
target_link_libraries(your_app PRIVATE zerobus::zerobus)

Чтобы вместо этого использовать готовую сборку через add_subdirectory, сначала укажите пути FFI, чтобы CMake компоновал архив из комплекта, а не пытался собрать его из отсутствующих исходников Rust. Использование zerobus_ffi.lib в Windows:

set(ZEROBUS_FFI_LIBRARY "/path/to/bundle/lib/libzerobus_ffi.a")
set(ZEROBUS_FFI_HEADER_DIR "/path/to/bundle/lib")
add_subdirectory(path/to/zerobus-sdk/cpp zerobus-cpp)
target_link_libraries(your_app PRIVATE zerobus::zerobus)

Пример JSON:

Прием является асинхронным и конвейерным. Методы ingest_* помещывают запись и возвращаются немедленно. Поставьте пакет в очередь и вызовите flush() один раз, вместо того чтобы ждать после каждой записи.

#include "zerobus/zerobus.hpp"
#include <string>
#include <vector>

int main() {
  // See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
  const std::string SERVER_ENDPOINT = "https://1234567890123456.zerobus.eastus.azuredatabricks.net";
  const std::string DATABRICKS_WORKSPACE_URL = "https://adb-1234567890123456.12.azuredatabricks.net";
  const std::string TABLE_NAME = "main.default.air_quality";
  const std::string CLIENT_ID = "your-client-id";
  const std::string CLIENT_SECRET = "your-client-secret";

  zerobus::Sdk sdk = zerobus::Sdk::builder()
                         .endpoint(SERVER_ENDPOINT)
                         .unity_catalog_url(DATABRICKS_WORKSPACE_URL)
                         .application_name("my-app")
                         .build();

  zerobus::TableProperties table;
  table.table_name = TABLE_NAME;   // empty descriptor => JSON stream

  zerobus::StreamOptions options;
  options.record_type = zerobus::RecordType::Json;

  zerobus::Stream stream =
      sdk.create_stream(table, CLIENT_ID, CLIENT_SECRET, options);

  std::vector<std::string> batch = {
      R"({"device_name": "sensor-001", "temp": 20, "humidity": 60})",
      R"({"device_name": "sensor-002", "temp": 22, "humidity": 55})",
  };
  stream.ingest_json_records(batch);   // queue the batch — no per-record wait
  stream.flush();                      // wait once for all acks
  stream.close();

  return 0;
}

Каждый сбой вызывает исключение zerobus::ZerobusException, которое содержит сообщение и флаг is_retryable(). Чтобы отслеживать надёжность в потоке непрерывных данных без блокировки, зарегистрируйте AckCallback через StreamOptions::ack_callback. Обратные вызовы выполняются последовательно в фоновом потоке и должны быть noexcept. См. документацию по C++ SDK с полным описанием многопоточности, drain-policy и контракта времени жизни.

Для типобезопасной загрузки данных можно использовать Protocol Buffers одним из двух способов:

  • Создайте схему из каталога Unity с помощью ProtoSchema::from_uc_json(). Это создаёт дескриптор и кодировщик JSON в proto непосредственно на основе метаданных таблицы, поэтому ему не нужен ни файл .proto, ни protoc.

    1. Получение JSON метаданных таблицы из API получения таблицы (GET /api/2.1/unity-catalog/tables/{full_name}). Субъект-служба нуждается SELECT в таблице.
    2. Передайте метаданные в ProtoSchema::from_uc_json(), чтобы создать дескриптор и кодировщик.
    3. Задайте TableProperties::descriptor_proto, затем загрузите с помощью ingest_proto_records().
  • Скомпилируйте флажок .protoprotoc для ввода во время компиляции.

Пошаговое руководство по выполнению см. в примерах буферов протокола.

Сведения о столбчатой или пакетно-ориентированной загрузке пакетов записей Apache Arrow по тому же соединению gRPC см. в разделе Использование Arrow Flight с Zerobus Ingest.

Полная документация, параметры конфигурации, пакетная приемка и примеры буфера протокола см. в репозитории пакета SDK для C++.

Пакет SDK для C#

Important

SDK для C# / .NET находится в бета-версии. Пакет Databricks.Zerobus находится в предварительном релизе.

Требуется версия .NET 8.0 и выше. SDK обеспечивает нативное gRPC-потоковое вещание, OAuth и автоматическое восстановление. Он поддерживает JSON для простых настроек и буферов протокола для рабочих нагрузок. Arrow Flight недоступен в CDK C#.

Добавьте пакет Databricks.Zerobus в проект:

dotnet add package Databricks.Zerobus

Пример JSON:

using Databricks.Zerobus;

// See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
const string SERVER_ENDPOINT = "https://1234567890123456.zerobus.eastus.azuredatabricks.net";
const string DATABRICKS_WORKSPACE_URL = "https://adb-1234567890123456.12.azuredatabricks.net";
const string TABLE_NAME = "main.default.air_quality";
const string CLIENT_ID = "your-client-id";
const string CLIENT_SECRET = "your-client-secret";

using var sdk = ZerobusSdk.CreateBuilder()
    .Endpoint(SERVER_ENDPOINT)
    .UnityCatalogUrl(DATABRICKS_WORKSPACE_URL)
    .Build();

using var stream = sdk.CreateJsonStream(TABLE_NAME, CLIENT_ID, CLIENT_SECRET);

long offset = stream.IngestRecord(
    """{"device_name": "sensor-1", "temp": 22, "humidity": 55}""");
stream.WaitForOffset(offset);
stream.Close();

IngestRecord возвращает смещение пластинки и WaitForOffset блокирует, пока запись не станет долговечной. Чтобы загрузить пакет, используйте IngestRecords, который принимает массив записей и возвращает последнее смещение. Блокировка по смещению необязательна. См. раздел «Блокировка сообщений и подтверждение».

Protocol Buffers: Для типобезопасной загрузки создайте поток данных с sdk.CreateProtoStream(TABLE_NAME, descriptorProto, CLIENT_ID, CLIENT_SECRET), где descriptorProto — сериализованные байты DescriptorProto для вашего скомпилированного сообщения, затем выполните загрузку с помощью stream.IngestRecord(protoBytes).

Полная документация, опции конфигурации и примеры буфера протоколов смотрите в репозитории C# SDK.

TypeScript SDK

требуется Node.js 16 или более поздней версии. SDK обеспечивает высокую производительность с асинхронной поддержкой через JavaScript Promises. Он поддерживает JSON (самый простой) и Protocol Buffers (рекомендуется для использования в продуктивной среде).

npm install @databricks/zerobus-ingest-sdk

Пример JSON:

import { ZerobusSdk, RecordType } from '@databricks/zerobus-ingest-sdk';

// See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
const SERVER_ENDPOINT = 'https://1234567890123456.zerobus.eastus.azuredatabricks.net';
const DATABRICKS_WORKSPACE_URL = 'https://adb-1234567890123456.12.azuredatabricks.net';
const TABLE_NAME = 'main.default.air_quality';
const CLIENT_ID = 'your-client-id';
const CLIENT_SECRET = 'your-client-secret';

const sdk = new ZerobusSdk(SERVER_ENDPOINT, DATABRICKS_WORKSPACE_URL);

const stream = await sdk.createStream({ tableName: TABLE_NAME }, CLIENT_ID, CLIENT_SECRET, {
  recordType: RecordType.Json,
});

try {
  for (let i = 0; i < 100; i++) {
    const record = { device_name: `sensor-${i}`, temp: 22, humidity: 55 };
    await stream.ingestRecordOffset(record);
  }
} finally {
  await stream.close();
}

Protocol Buffers: Для безопасного приема данных используйте Protocol Buffers с RecordType.Proto (по умолчанию) и укажите descriptorProto в свойствах таблицы.

Arrow Flight: Для колоночного или пакетного приёма данных Apache Arrow RecordBatch по тому же соединению gRPC см. раздел «Использование Arrow Flight с Zerobus Ingest».

Полная документация, параметры конфигурации, пакетная приемка и примеры буфера протокола см. в репозитории пакета SDK TypeScript.

REST API

REST API позволяет загрузить одну запись, отправив POST-запрос HTTP в конечную точку /zerobus/v1/tables/<table-name>/insert. Сама запись включается в текст запроса и должна быть в формате JSON.

В этом примере показано, как использовать CURL для отправки данных в Загрузчик Zerobus через REST API.

Заголовки

Для запроса требуется два определенных заголовка HTTP для правильной проверки подлинности и форматирования запроса.

  • Тип содержания: application/json
    • Обязательное поле для указания типа контента. В настоящее время JSON является единственным поддерживаемым форматом сообщения.
  • Авторизация: маркер носителя <>
    • Замените <токен>, который вы получили, используя команду curl, на токен OAuth.

Получение токена OAuth: Срок действия этих токенов истекает каждый час, поэтому они должны обновляться. Их можно обновить, повторно выбрав маркер OAuth.

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

  • $CATALOG, $SCHEMA, , $TABLE, $WORKSPACE_ID$WORKSPACE_URL
  • $DATABRICKS_CLIENT_ID и $DATABRICKS_CLIENT_SECRET.
    • Эти два параметра соответствуют созданному принципу службы.
authorization_details=$(cat <<EOF
[{
  "type": "unity_catalog_privileges",
  "privileges": ["USE CATALOG"],
  "object_type": "CATALOG",
  "object_full_path": "$CATALOG"
},
{
  "type": "unity_catalog_privileges",
  "privileges": ["USE SCHEMA"],
  "object_type": "SCHEMA",
  "object_full_path": "$CATALOG.$SCHEMA"
},
{
  "type": "unity_catalog_privileges",
  "privileges": ["SELECT", "MODIFY"],
  "object_type": "TABLE",
  "object_full_path": "$CATALOG.$SCHEMA.$TABLE"
}]
EOF
)

export OAUTH_TOKEN=$(curl -X POST \
  -u "$DATABRICKS_CLIENT_ID:$DATABRICKS_CLIENT_SECRET" \
  -d "grant_type=client_credentials" \
  -d "scope=all-apis" \
  -d "resource=api://databricks/workspaces/$WORKSPACE_ID/zerobusDirectWriteApi" \
  --data-urlencode "authorization_details=$authorization_details" \
  "$WORKSPACE_URL/oidc/v1/token" | jq -r '.access_token')

Прием записей:

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

  • $ZEROBUS_ENDPOINT
  • $CATALOG, $SCHEMA, , $TABLE, $WORKSPACE_ID$WORKSPACE_URL
  • $OAUTH_TOKEN
    • Это было создано на предыдущем шаге.

Текст запроса должен быть списком объектов JSON.

curl -X POST \
  "$ZEROBUS_ENDPOINT/zerobus/v1/tables/$CATALOG.$SCHEMA.$TABLE/insert" \
  -H "Content-Type: application/json" \
  -H "Authorization: Bearer $OAUTH_TOKEN" \
  -d '[{ "device_name": "device_num_1", "temp": 28, "humidity": 60 },
       { "device_name": "device_num_1", "temp": 28, "humidity": 60 }]'

Если все данные заполнены правильно, следует получить пустой ответ JSON с кодом состояния HTTP 200.

Управление ошибками

Приведённые выше примеры показывают счастливый путь. В рабочей среде добавьте обработку ошибок для процесса загрузки данных. SDK автоматически выполняет повторные попытки при временных ошибках, например при проблемах с сетью, благодаря встроенному механизму восстановления. Сбои, от которых не удаётся восстановиться, такие как недействительные учетные данные или отсутствующая таблица, появляются следующим ZerobusExceptionобразом:

from zerobus.sdk.shared import ZerobusException

try:
    stream.ingest_record_offset(record)
except ZerobusException as e:
    # Handle the failure: log it, fix the cause, recover on a new stream, or stop.
    ...

Пакеты SDK также автоматически восстанавливаются после временных сбоев и позволяют восстанавливать неподтверждённые записи при окончательном сбое потока. Информацию о шаблонах отказоустойчивого клиента и полный справочник по ошибкам см. в разделах Шаблоны восстановления и повторных попыток и Обработка ошибок Zerobus Ingest.

Дальнейшие действия