For AI agents: the complete documentation index is available at https://docs.ovhcloud.com/fr/llms.txt, the full documentation bundle is available at https://docs.ovhcloud.com/fr/llms-full.txt, and this page is available as Markdown at https://docs.ovhcloud.com/fr/guides/public-cloud/data-platform/tutorials-kafka.md.

Streamer des données depuis Apache Kafka vers la Platform

Voir en Markdown

Ce tutoriel montre comment utiliser les données d'un Broker Apache Kafka dans la platform

Objectif

Ce tutoriel montre comment utiliser les données d'un Broker Apache Kafka dans la platform.
La première section est dédiée à la création de données de test sur votre serveur Kafka. Si vous avez déjà des messages sur votre Broker, vous pouvez ignorer cette étape.

Introduction

Prérequis

Pour suivre ce tutoriel, vous devez disposer d'un Broker Kafka opérationnel. L'exemple de code fourni a été écrit pour un serveur ne nécessitant aucune authentification particulière (c'est-à-dire que toute personne disposant de l'adresse IP peut lire les messages). Si votre serveur Kafka a son authentification configurée différemment, vous devez adapter le code utilisé ici en conséquence.

De plus, nous recommandons d'avoir au moins suivi le premier tutoriel de démarrage avant de faire celui-ci. Ici, nous supposons que vous êtes à l'aise avec Data Platform et familier avec les principaux composants de la platform.

Vue d'ensemble des concepts

Data Platform s'intègre à Apache Kafka via le connecteur Kafka dans les Connectors. Ce connecteur permet de récupérer des données depuis un ou plusieurs topics du même serveur vers Data Platform.

En général, les données sont ensuite ingérées dans les tables du Lakehouse Manager. Une table stockera les données d'un topic et chaque message d'un topic devient une ligne de données dans la table correspondante.

Info

À propos des champs imbriqués... actuellement, Data Platform ne prend en charge que les messages au format JSON sans imbrication. Par conséquent, seuls les champs situés à la racine de la représentation JSON sont pris en compte.

Une fois que vous avez configuré votre connexion aux topics Kafka dans Connectors et configuré vos tables du Lakehouse Manager, vous devrez charger les données des topics vers les tables en lançant une action Load à l'aide du Data Processing Engine.

Après le chargement des données, vos messages seront automatiquement chargés dans les tables du Lakehouse Manager tant que votre action est en cours d'exécution. Notez que vos actions seront exécutées en mode d'exécution Serverless par défaut, qui a un timeout. C'est pourquoi nous vous recommandons d'utiliser le mode d'exécution Always-up si vous utilisez le connecteur Kafka.

Voyons maintenant comment tout cela fonctionne en pratique !

Configurer des données de test (optionnel)

Pour envoyer des messages à votre broker Kafka à des fins de test, vous pouvez configurer un Producer dans Data Platform à l'aide d'une Custom Action DPE. Créez une Custom Action dans le Data Processing Engine, sélectionnez start with a boilerplate et remplacez le code du boilerplate par celui ci-dessous (un simple jeu de devinette de prénom) :

from forepaas.dwh.connect import connect
import logging, time, json, random
from kafka import KafkaProducer # kafka-python

logger = logging.getLogger(__name__)

TOPIC = "sample"
#   Les deux lignes suivantes DOIVENT être remplacées par votre propre adresse et port Kafka
KAFKA_PORT="9092"
KAFKA_ADDRESS=["10.152.1.186","10.152.1.187","10.152.7.65"]

def generate_bootstrap():
    servers = [f"{x}:{KAFKA_PORT}" for x in KAFKA_ADDRESS]
    return ",".join(servers)

def customfunc(event):
    logger.info("Begin function customfunc")
    i = 0

    #  Jeu de devinette de prénom : 
    #   - Gagnez 5 à 15 points en devinant le bon prénom
    #   - Perdez 5 à 15 points en devinant le mauvais prénom
    #   - Aucun point pour deviner les autres prénoms

    # génère la liste des prénoms
    nameslist = ["helene", "francoise", "lea", "lorene", "claire", "lise", "karen", "elise", "elia", "annabele"]

    # sélectionne un "mauvais" et un "bon" prénom
    correct_name = nameslist[0]
    bad_name = nameslist[1]
    logger.info("Correct name:" + correct_name)
    logger.info("Bad name:" + bad_name)

    try:
        producer = KafkaProducer(bootstrap_servers=generate_bootstrap())
        while True:
            i +=1

            name = random.choice(nameslist)
            # distingue les bons et mauvais prénoms, par rapport aux autres
            if name == correct_name:
                points = random.randint(5,15)
            elif name == bad_name:
                points = random.randint(-15,-5)
            else:
                points = 0

            value = {
                "index":i,
                "points":points,
                "name":name,
                }

            # Connexion Kafka, et publication vers le broker

            producer.send(TOPIC, json.dumps(value).encode("utf-8"))
            if i % 1000 == 0:
                logger.info("SENT 1000 records")

                # Pause de 1 seconde tous les 1000 messages
                time.sleep(1)
            if i ==10000:
                return
        logger.info("END function customfunc")
    except Exception as err:
        logger.critical(err)

Le code ci-dessus représente un jeu de devinette de prénom, il enverra simplement des messages représentant des tentatives. Chaque message contient un index, un prénom (la tentative) et les points gagnés pour cette tentative. Une fois 10 000 messages envoyés, l'action s'arrêtera et vous devriez avoir quelques messages dans votre Broker.

Warning

N'oubliez pas d'ajouter le module kafka aux dépendances Python de votre Custom action.

Maintenant, exécutez l'action pour peupler votre topic avec les données de test, puis arrêtez l'action une fois que quelques milliers d'enregistrements ont été envoyés.

actions log

Connecter votre serveur Kafka à Data Platform

Configurer votre connexion Kafka

La première chose à faire est de configurer votre connexion à un serveur Kafka et de choisir un topic depuis lequel lire les données. Si vous avez besoin d'aide, vous pouvez consulter notre article dédié Connecteur Apache Kafka.

Configurer votre schéma

Maintenant que votre connexion et votre topic sont correctement configurés, vous êtes prêt à accéder à vos messages. Pour cela, vous devez aller dans l'onglet Analyzer et extraire les métadonnées du topic de votre connexion.

La connexion apparaîtra dans la barre latérale gauche et le topic s'affichera en cliquant sur la connexion. Sélectionnez le topic et cliquez sur le bouton Extract metadata.

analyzer screen with metadata extracted

Après l'extraction des métadonnées, les messages apparaîtront dans le panneau d'aperçu où chaque ligne correspond à un message. Cochez les cases du panneau de métadonnées pour configurer les champs du message qui seront inclus lorsque vous utiliserez votre message dans Data Platform.

Info

Vous remarquerez peut-être qu'il y a des champs supplémentaires dans votre message. Le timestamp, la date ainsi que l'offset sont fournis par le Broker et correspondent respectivement au timestamp d'arrivée de vos messages, à la date d'arrivée et à l'offset du topic. Ils peuvent être utiles pour certains cas d'usage mais, si vous ne souhaitez pas les inclure dans votre Projet de données, décochez-les simplement dans le panneau de métadonnées et ils seront ignorés par le reste de la platform.

Créer et construire votre table

Avant de charger vos données dans Data Platform, vous devez créer et construire la table qui les stockera. Si vous n'êtes pas familier avec ces concepts, vous pouvez consulter notre article Tables. Vous voudrez probablement lire les sections Create a new table et Build all tables.

Charger vos données dans Data Platform

Configurer l'action Load

Par rapport aux autres connecteurs, il existe quelques différences lors de la création d'une action Load avec une source de streaming telle qu'Apache Kafka.

Pour commencer, sélectionnez la table liée à votre topic comme Source lors de la configuration de l'action (si vous avez déjà généré l'action lors de la création de la table, vous n'aurez pas besoin de sélectionner la table).

load

Modes d'exécution

Lors de l'exécution d'une action Load connectée à une source Kafka, nous vous recommandons vivement de sélectionner le mode d'exécution Always-up.

exec

Quel que soit le mode d'exécution que vous utilisez, votre action s'exécutera jusqu'à ce que des données arrivent. Une fois que c'est le cas, elles seront chargées dans Data Platform.

Si vous utilisez le mode d'exécution Serverless, l'action s'arrêtera une fois le timeout atteint (2 heures par défaut). Étant donné que de nouvelles données peuvent arriver à tout moment dans votre Broker Kafka, cela signifie que vous devez relancer cette action après son arrêt si vous souhaitez continuer à alimenter Data Platform depuis votre Kafka. C'est pourquoi nous recommandons l'utilisation du mode d'exécution Always-up, notamment pour les environnements de production.

Info

Si vous ne souhaitez pas utiliser le mode d'exécution Always-up, vous pouvez également configurer des déclencheurs basés sur le temps pour exécuter automatiquement votre action selon un intervalle de temps prédéfini.

Segmentation automatique

Lorsque vous utilisez un connecteur Kafka, vous pouvez bénéficier d'un temps d'exécution plus rapide en utilisant la fonctionnalité Automatic Segmentation. Cette option est disponible dans les Preferences de votre action, utilisez-la pour les charges de travail lourdes !

Offset personnalisé

Une autre option qui s'offre à vous est de commencer à lire vos messages à partir d'un offset personnalisé au lieu de configurer vos topics pour lire depuis le message le plus ancien ou le plus récent.

Pour remplacer la politique d'offset Latest ou Earliest que vous avez configurée sur votre topic et commencer à lire les messages à partir d'un offset défini, vous devez utiliser le mode Advanced des Actions. Ajoutez simplement un champ à l'intérieur du champ paras.load_from comme dans l'exemple suivant (démarre la lecture à partir de l'offset 7) :

"params": {
    "load_from": [
      {
        "offset_number": 7,
        ...
      }
    ],
...
offset
Info

Notez que cela remplacera la politique d'offset Latest ou Earliest que vous avez configurée sur votre topic.

Considérations techniques

Une autre considération à garder à l'esprit est que, si vous êtes en mode Earliest et que vous changez votre table de destination après avoir lu les messages les plus anciens, lors d'une nouvelle exécution de l'action Load, tous les messages du topic seront relus.

Cela se produit car l'offset du dernier message consommé par Data Platform est stocké dans les métadonnées de la table de destination de la base de données Data Platform. Si vous utilisez une nouvelle table de destination, l'offset repartira depuis le début.

Aller plus loin

Si vous avez besoin d'une formation ou d'une assistance technique pour la mise en oeuvre de nos solutions, contactez votre commercial ou cliquez sur ce lien pour obtenir un devis et demander une analyse personnalisée de votre projet à nos experts de l’équipe Professional Services.

Posez vos questions, faites-nous part de vos commentaires et interagissez directement avec l’équipe qui développe la Data Platform sur le canal Discord dédié.

Si vous avez besoin d'une assistance concernant vos services OVHcloud, créez une demande depuis notre centre d'aide.

Rejoignez notre communauté d'utilisateurs.

Cette page vous a-t-elle aidé ?