Aller au contenu principal

Concepts de Zerobus Ingest

Cette page décrit les concepts fondamentaux de Zerobus Ingest dans Lakeflow Connect : le fonctionnement du service, ses streams, son serveur et ses clients, ainsi que les types de données qu'il prend en charge.

Accéder à un concept :

Comment fonctionne Zerobus Ingest

Un producteur de données ouvre d'abord un Stream vers l'API Zerobus Ingest et spécifie une table Delta cible, construit un message correspondant à son schéma, puis envoie le message via le Stream ouvert. Le service rend les données durables et accuse réception du message du client. Il matérialise ensuite les données dans la table Delta, de manière optimisée, en tant qu'étape distincte. L'accusé de réception confirme la durabilité, et non la possibilité d'interrogation. Voir Communication asynchrone pour savoir comment cela fonctionne et ce que cela implique pour votre client.

Zerobus Ingest est un service Serverless qui monte en charge de manière élastique avec votre charge de travail. Pour savoir comment il monte en charge, voir Comment Zerobus Ingest monte en charge ci-dessous.

Comment fonctionne Zerobus Ingest

Cette section couvre également les manières de vous connecter à Zerobus Ingest et les formes que peuvent prendre vos données :

  • Protocoles d'API: les protocoles d'API, gRPC avec SDK, REST et OpenTelemetry, et quand utiliser chacun d'eux.
  • Types de messages: les formats d'enregistrement, JSON, Protocol Buffers (protobuf) et Apache Arrow, ainsi que les cas d'utilisation de chacun.

Serveur

Le service Zerobus Ingest ne crée ni ne manipule automatiquement les tables. Les utilisateurs doivent créer la table eux-mêmes. Les tables et leurs schémas sont les sources faisant autorité pour les attentes concernant les données entrantes.

Le serveur Zerobus Ingest accepte les données envoyées par les clients et vérifie qu'elles correspondent au schéma de la table cible. Si l'enregistrement est conforme, le serveur le rend durable et en confirme la réception au client. La matérialisation de l'enregistrement dans la table Delta, afin qu'il puisse être interrogé, s'effectue lors d'une étape distincte peu de temps après.

Les responsabilités du service incluent :

  • Validation du schéma du message par rapport à la table.
  • Rendre l'enregistrement durable et le confirmer au client. La confirmation atteste de la durabilité, et non du fait que l'enregistrement est déjà interrogeable.
  • Matérialisation des données dans la table cible en temps opportun, moment à partir duquel elles deviennent interrogeables. Pour les chiffres de latence, consultez Latence.

Client

Un client se connecte à Zerobus Ingest, envoie des enregistrements et confirme qu'ils sont durables. Lorsque vous utilisez un SDK Zerobus Ingest, le SDK gère la plupart de ces aspects pour vous ; il est donc utile de séparer ce que vous configurez de ce que le SDK fait automatiquement.

Vous configurez ou implémentez :

  • Sélection d'une table cible.
  • Ouverture d'un Stream vers le service Zerobus Ingest.
  • Construction d’un message compatible avec le schéma et envoi de celui-ci.

Voici les opérations dont le SDK se charge automatiquement :

  • Confirmations de message. Le SDK exécute la boucle de confirmation pour vous et fait remonter les confirmations de durabilité via des offsets ou un callback de confirmation. Vous ne bloquez sur un enregistrement spécifique que lorsque votre application en a besoin. Voir Communication asynchrone.
  • Récupération. Par default, le SDK se reconnecte et rejoue les enregistrements non acquittés en cas d’erreurs temporaires.
    • Vous pouvez désactiver la récupération intégrée et implémenter votre propre mécanisme de récupération à la place. Pour savoir ce qui Trigger la récupération, les options de configuration et les modèles de récupération personnalisés, consultez Modèles de récupération et de nouvelle tentative.

Vous n'avez pas à écrire manuellement la logique d'accusé de réception ou de récupération lorsque vous utilisez un SDK. Pour les intégrations personnalisées qui n'utilisent pas de SDK, le repository du SDK Zerobus sert de référence pour la structure d'intégration et la gestion de la récupération.

Stream

Un Stream est une connexion directe entre votre client et le serveur Zerobus Ingest, établie via une connexion gRPC persistante et bidirectionnelle. Les SDK utilisent des Stream pour faciliter les connexions longue durée à throughput élevé.

  • Les flux (Streams) ne sont utilisés que dans l'API gRPC avec les SDK.
  • Un Stream ingère des données vers une table cible unique.
  • Ouvrez des streams supplémentaires pour écrire dans différentes tables, ou pour monter en charge le throughput d'un seul client aussi haut que vos charges de travail l'exigent.

Les Stream sont également l'unité de classement (voir Garanties de classement) et l'unité par laquelle Zerobus Ingest monte en charge (voir Comment Zerobus Ingest monte en charge).

Garanties d'ordre

L'ordre est garanti par Stream. Les enregistrements sont validés dans la table cible dans l'ordre où ils sont mis en file d'attente sur un seul Stream. Il n'y a pas d'ordre global entre les différents Streams. Plusieurs points de conception en découlent :

  • Si vous répartissez les enregistrements sur plusieurs streams (par exemple, round-robin), il n'y a aucune garantie de classement entre ces streams.
  • Si votre cas d’utilisation nécessite un ordre total unique parmi de nombreux producteurs ou Stream, imposez cet ordre dans votre application (par exemple, avec un timestamp ou un numéro de séquence sur lequel vous effectuez une query) plutôt que de vous fier à l’ordre d’ingestion.

Pourquoi le streaming gRPC

Comme la connexion gRPC d'un Stream reste ouverte, le client évite le coût de configuration par requête d'un protocole sans état et peut pousser un flux continu et volumineux d'enregistrements via un seul Canal de distribution. C'est ce qui fait des SDK le moyen d'ingestion offrant le plus haut throughput. Pour les autres interfaces (REST et OpenTelemetry) et savoir quand choisir chacune d'elles, consultez API protocols.

Comment Zerobus Ingest monte en charge

Zerobus Ingest est conçu pour une haute évolutivité, et il atteint cette échelle sans vous demander de planifier la capacité. Deux choix de conception rendent cela possible :

  • C'est du Serverless. Le service ajoute et supprime automatiquement de la capacité en fonction de l'évolution de la charge, vous n'avez donc pas à dimensionner les brokers ni à effectuer le provisionnement des partitions. Vous pouvez ouvrir autant de Streams simultanés et écrire dans autant de tables que votre charge de travail l'exige.
  • Les streams sont des unités de partitionnement dynamiques. Plutôt qu’un ensemble fixe de partitions devant être repartitionnées et rééquilibrées pour monter en charge, les streams peuvent être ouverts, fermés et alternés. L’alternance des streams permet au service de rééquilibrer la capacité et les ressources à mesure que la demande évolue ; vous pouvez ainsi monter en charge en ouvrant davantage de streams et en exécutant plus de producteurs, tandis que le service absorbe le reste.

Le résultat pratique est qu’un client « hello world » et une charge de travail à l’échelle du pétaoctet exécutent essentiellement le même code. La différence réside dans le nombre de producteurs et de streams que vous exécutez. Cette conception a permis une ingestion soutenue de plus de 1 000 milliards d’enregistrements dans une table Delta unique. Pour le contexte technique, consultez le billet de blog Ingesting the Milky Way: Petabyte-Scale with Zerobus Ingest.

Exigences relatives aux tables

Zerobus Ingest écrit dans une table Delta que vous créez et possédez. La table cible et le Workspace doivent répondre aux exigences suivantes :

  • Zerobus Ingest écrit uniquement dans des tables Delta gérées. L'écriture dans le stockage « default » n'est pas prise en charge.
  • Zerobus Ingest n’écrit pas dans un stockage sécurisé via un endpoint privé.
  • Zerobus Ingest ne prend pas en charge la recréation d'une table cible.
  • Les noms de table ne prennent en charge que les lettres ASCII, les chiffres et les traits de soulignement.
  • Le Workspace et la table cible doivent tous deux se trouver dans l'une des régions prises en charge.

Pour savoir comment les enregistrements sont validés par rapport au schéma de la table, consultez Gestion des schémas. Pour les fonctionnalités de table telles que le partitionnement et le clustering liquide, consultez Fonctionnalités des tables Delta.

Types de données pris en charge

Le tableau suivant présente les types Delta pris en charge et leurs types Protobuf correspondants pour l'ingestion.

Types Delta

Types Protobuf

INTEGER

int32

STRING

string

FLOAT

float

LONG

int64

SHORT

int32

DOUBLE

double

DECIMAL(p, s)

Texte décimal, par ex. « 123,45 », « 1e2 », etc.

string

BOOLEAN

bool

BINARY

bytes

BYTE (TINYINT)

int32

DATE

Doit être converti en int32 (nombre de jours depuis l'époque).

int32

TIMESTAMP

Doit être converti en int64 (temps epoch en microsecondes).

int64

TIMESTAMPNTZ

Doit être converti en int64 (temps epoch en microsecondes).

int64

ARRAY<TYPE>

repeated TYPE

MAP<K,V>

map<K,V>

Le sucre syntaxique Protobuf map n'est disponible que pour les compilateurs Protobuf en version 3 et ultérieure.

STRUCT<FIELDS>

message Nested { FIELDS }

VARIANT

Via les SDK gRPC et REST, ingérez une valeur Variant sous forme de chaîne encodée en JSON avec des clés de type STRING, et Zerobus Ingest écrit les données sans les fragmenter dans la colonne. Pour Apache Arrow Flight, le client construit plutôt les champs metadata et value qui soutiennent la colonne Variant. Consultez Ingestion de colonnes VARIANT.

Les formats pris en charge sont les suivants :

  • Objets : "{\"id\":0,\"example\":\"this is variant example\"}"
  • Primitives : "5", "3.14", "\"string\""
  • Tableaux : "[1,2,3]"

string

Types Delta

Types Protobuf

INTEGER

int32

STRING

string

FLOAT

float

LONG

int64

SHORT

int32

DOUBLE

double

DECIMAL(p, s)

Texte décimal, par ex. « 123,45 », « 1e2 », etc.

string

BOOLEAN

bool

BINARY

bytes

BYTE (TINYINT)

int32

DATE

Doit être converti en int32 (nombre de jours depuis l'époque).

int32

TIMESTAMP

Doit être converti en int64 (temps epoch en microsecondes).

int64

TIMESTAMPNTZ

Doit être converti en int64 (temps epoch en microsecondes).

int64

ARRAY<TYPE>

repeated TYPE

MAP<K,V>

map<K,V>

Le sucre syntaxique Protobuf map n'est disponible que pour les compilateurs Protobuf en version 3 et ultérieure.

STRUCT<FIELDS>

message Nested { FIELDS }

VARIANT

Via les SDK gRPC et REST, ingérez une valeur Variant sous forme de chaîne encodée en JSON avec des clés de type STRING, et Zerobus Ingest écrit les données sans les fragmenter dans la colonne. Pour Apache Arrow Flight, le client construit plutôt les champs metadata et value qui soutiennent la colonne Variant. Consultez Ingestion de colonnes VARIANT.

Les formats pris en charge sont les suivants :

  • Objets : "{\"id\":0,\"example\":\"this is variant example\"}"
  • Primitives : "5", "3.14", "\"string\""
  • Tableaux : "[1,2,3]"

string