Remarque
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de modifier des répertoires.
Cette page décrit les concepts fondamentaux de l’ingestion de Zerobus dans Lakeflow Connect : le fonctionnement du service, ses flux, son serveur et ses clients, ainsi que les types de données qu’il prend en charge.
Passons à un concept :
- Comment fonctionne Zerobus Ingest : le processus de bout en bout, du producteur à la table interrogeable.
- Serveur : ce dont est responsable le service d’ingestion Zerobus.
- Client : ce dont votre producteur est responsable.
- Flux : la connexion que votre client ouvre pour ingérer les données.
- Comment Zerobus Ingest évolue : pourquoi il évolue sans planification de capacité.
- Exigences de table : ce que votre table et votre espace de travail cibles doivent satisfaire.
- Types de données pris en charge : le mappage de type Delta-à-Protobuf.
Fonctionnement de Zerobus Ingest
Un producteur de données ouvre d’abord un flux vers l’API Zerobus Ingest, spécifie une table Delta cible, construit un message correspondant à son schéma, puis pousse le message à travers le flux ouvert. Le service rend les données durables et reconnaît le 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é, pas la possibilité de faire des questions. Consultez la communication asynchrone pour savoir comment cela fonctionne et ce que cela signifie pour votre client.
Zerobus Ingest est un service serverless qui évolue élastiquement avec votre charge de travail. Pour voir comment cela évolue, voir Comment Zerobus Ingest évolue ci-dessous.
Cette section explique également les façons dont vous vous connectez à Zerobus Ingest et les formes que vos données peuvent prendre :
- Protocoles API : Les protocoles API, gRPC avec SDK, REST et OpenTelemetry, et quand utiliser chacun.
- Types de messages : Les formats d’enregistrement, JSON, tampons de protocole (protobuf) et Apache Arrow, et quand utiliser chacun.
Server
Le service d’ingestion Zerobus ne crée pas ou manipule automatiquement des tables. Les utilisateurs doivent créer eux-mêmes la table. Les tables et leurs schémas sont les sources faisant autorité pour les attentes des données entrantes.
Le serveur Zerobus Ingest accepte les données envoyées par les clients et valide qu’elles correspondent au schéma de la table cible. Si l’enregistrement convient, le serveur le rend persistant et en accuse réception auprès du client. La matérialisation de l’enregistrement dans la table Delta, pour qu’il devienne consultable, se fait comme une étape séparée peu après.
Les responsabilités de service sont les suivantes :
- Validation du schéma du message par rapport à la table.
- Rendre le dossier durable et en accuser réception auprès du client. L’accusé de réception confirme la durabilité, mais pas encore que l’enregistrement puisse faire l’objet de requêtes.
- Matérialiser les données dans la table cible en temps voulu, c’est à ce moment-là qu’elles deviennent interrogables. Pour les chiffres de latence, voir Latence.
Client
Un client se connecte à Zerobus Gest, envoie des enregistrements et confirme qu’ils sont durables. Quand vous utilisez un SDK d’ingestion Zerobus, le SDK gère la plupart de ces choses pour vous, donc cela aide à séparer automatiquement ce que vous configurez de ce que fait le SDK.
Vous configurez ou implémentez :
- Sélection d’une table cible.
- Ouverture d’un flux vers le service Zerobus Ingest.
- Construire un message compatible avec le schéma et l’envoyer.
Le SDK gère automatiquement :
- Accusés de réception du message Le SDK exécute la boucle d’acquittement pour vous et rend disponibles les confirmations de durabilité via des décalages ou une fonction de rappel d’accusé de réception. Vous ne verrouillez un enregistrement spécifique que lorsque votre application en a besoin. Voir Communication asynchrone.
- Récupération. Par défaut, le SDK se reconnecte et renvoie les enregistrements non confirmés en cas de défaillances temporaires.
- Vous pouvez désactiver la récupération intégrée et mettre en place votre propre mécanisme de récupération. Pour ce qui déclenche la récupération, les options de configuration et les motifs de récupération personnalisés, voir Patrons de récupération et de réessayage.
Vous n’avez pas besoin de coder vous-même les mécanismes d’accusé de réception ou de reprise lorsque vous utilisez un SDK. Pour les intégrations personnalisées qui n’utilisent pas de SDK, le dépôt SDK Zerobus est une référence pour la structure d’intégration et la gestion de la récupération.
Flux de données
Un flux est une connexion directe entre votre client et le serveur Zerobus Ingest, établie via une connexion gRPC bidirectionnelle et persistante. Les kits SDK utilisent des flux pour faciliter les connexions à long terme et à débit élevé.
- Les « streams » sont utilisés uniquement dans l'API gRPC avec les SDKs.
- Un flux ingère des données dans une table cible unique.
- Ouvrez des flux supplémentaires pour écrire dans différentes tables, ou pour augmenter le débit d’un seul client autant que votre charge de travail l’exige.
Les flux sont également l’unité d’ordonnancement (voir Garanties d’ordonnancement) et l’unité selon laquelle Zerobus Ingest monte en charge (voir Comment Zerobus Ingest monte en charge).
Garanties de commande
La commande est garantie par flux. Les enregistrements sont écrits dans la table cible dans l’ordre dans lequel ils sont placés en file d’attente dans un flux unique. Il n’y a pas d’ordre global entre les courants. Plusieurs points de conception découlent de cela :
- Si vous répartissez les enregistrements sur plusieurs flux (par exemple, round-robin), il n’existe aucune garantie d’ordre entre ces flux.
- Si votre cas d’usage nécessite une commande totale unique sur plusieurs producteurs ou streams, faites respecter cette commande dans votre application (par exemple, avec un horodatage ou un numéro de séquence que vous interrogez) plutôt que de vous fier à une commande d’ingestion.
Pourquoi le streaming gRPC
Parce que la connexion gRPC d’un flux reste ouverte, le client évite le coût d’installation par requête d’un protocole sans état et peut faire passer un flux continu et à fort volume d’enregistrements sur un seul canal. C’est ce qui fait des SDK le moyen d’ingérer offrant le débit le plus élevé. Pour les autres interfaces (REST et OpenTelemetry) et le moment de choisir chacune, voir protocoles API.
Comment Zerobus Ingest évolue
Zerobus Ingest est conçu pour une grande scalabilité, et il atteint cette échelle sans que vous ayez à planifier la capacité. Deux choix de conception rendent cela possible :
- C’est sans serveur. 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 à provisionner les partitions. Vous pouvez ouvrir autant de flux simultanés et écrire sur autant de tables que votre charge de travail en a besoin.
- Les flux sont des unités de partitionnement dynamique. Plutôt qu’un ensemble fixe de partitions qu’il faut repartitionner et rééquilibrer pour passer à l’échelle, les flux peuvent être ouverts, fermés et renouvelés. Faire tourner les flux permet au service de rééquilibrer la capacité et les ressources selon les variations de la demande, donc vous pouvez évoluer en ouvrant plus de flux et en faisant tourner plus de producteurs pendant que le service absorbe le reste.
Le résultat pratique est qu’un client « hello world » et une charge de travail à l’échelle pétaoctet exécutent essentiellement le même code. La différence réside dans le nombre de producteurs et de streams que vous gérez. Ce design a permis d’ingérer plus d’un trillion de disques dans une seule table Delta. Pour le contexte technique, voir l’article de blog Ingesting the Milky Way: Petabyte-Scale with Zerobus Ingest.
Exigences du tableau
Zerobus Ingest écrit dans une table Delta que vous créez et dont vous êtes propriétaire. La table et l’espace de travail ciblés doivent répondre à ces exigences :
- Zerobus Ingest écrit uniquement sur des tables Delta gérées. L’écriture dans le stockage par défaut n’est pas prise en charge.
- Zerobus Ingest n’écrit pas sur un stockage sécurisé via un point de terminaison privé.
- Zerobus Ingest ne supporte pas la recréation d’une table cible.
- Les noms des tableaux ne prennent en charge que les lettres, chiffres et soulignements ASCII.
- L’espace de travail 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 table, voir Gestion du schéma. Pour les caractéristiques de table telles que le partitionnement et le clustering liquide, voir les caractéristiques de la table 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 exemple « 123.45 », « 1e2 », etc. |
string |
BOOLEAN |
bool |
BINARY |
bytes |
BYTE (TINYINT) |
int32 |
DATEDoit être converti en int32 (nombre de jours depuis l'époque). |
int32 |
TIMESTAMPDoit être converti en int64 (heure epoch en microsecondes). |
int64 |
TIMESTAMPNTZDoit être converti en int64 (heure epoch en microsecondes). |
int64 |
ARRAY<TYPE> |
repeated TYPE |
MAP<K,V> |
map<K,V>Le map sucre syntactique Protobuf est disponible uniquement pour les compilateurs Protobuf version 3 et ultérieure. |
STRUCT<FIELDS> |
message Nested { FIELDS } |
VARIANTSur les SDK gRPC et REST, on intègre 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 non déchiquetées dans la colonne. Pour Apache Arrow Flight, le client construit plutôt les champs sous-jacents metadata et value de la colonne Variant. Voir l’ingestion des colonnes VARIANT.Les formats pris en charge sont les suivants :
|
string |