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/databases/kafka-dev-python-basics.md.

Créer des applications 'publisher' et 'consumer' avec Analytics avec Kafka

Voir en Markdown

Développez vos premières applications Python en utilisant Kafka

Objectif

Public Cloud Databases pour Kafka vous permet de vous concentrer sur la création et le déploiement de vos applications cloud, pendant qu'OVHcloud prend en charge l'infrastructure Kafka et son maintien en conditions opérationnelles.

Kafka est une plateforme utilisée pour le traitement de flux. Il s'agit fondamentalement d'une file de messages pub/sub massivement scalable.

L'objectif de ce tutoriel est de vous montrer les étapes pour disposer de vos premières applications Python utilisant Kafka.
Une application pourra s'abonner à un topic et consommer des messages, l'autre pourra produire et publier des messages dans un topic.
Vous disposerez ainsi de toutes les bases pour développer votre propre solution utilisant Kafka.

Prérequis

  • Un projet Public Cloud dans votre compte OVHcloud.
  • Un service Public Cloud Databases pour Kafka exécuté et configuré. Ce guide peut vous aider à répondre à ce prérequis.
  • En suivant le guide de démarrage, enregistrez tous les certificats dans un dossier dédié :
    • le certificat du serveur sous le nom ca.pem
    • le certificat utilisateur sous le nom service.cert
    • la clé d'accès utilisateur sous le nom service.key
  • Un environnement Python avec une version stable et une connectivité réseau publique (Internet). Ce guide a été réalisé avec Python 3.12.2.

En pratique

Info

L'ensemble du code source est disponible sur le dépôt GitHub public-cloud-examples.

Étape 1 - Consommer des messages

L'une des applications s'abonnera à un topic de votre service Kafka et attendra de consommer tout message entrant.

import os

from confluent_kafka import Consumer, KafkaException
from rich.console import Console

# retrieve URI to connect to your Kafka service from your OS environment variables
KAFKA_SERVICE_URI = os.getenv("KAFKA_SERVICE_URI")
if KAFKA_SERVICE_URI is None:
    raise ValueError("KAFKA_SERVICE_URI is not set")

# set up the configuration to access your Kafka service in a secure way
conf = {
    "bootstrap.servers": KAFKA_SERVICE_URI,
    "client.id": "customer",
    "group.id": "readers",
    "security.protocol": "SSL",
    "ssl.ca.location": "./sslcerts/ca.pem",
    "ssl.certificate.location": "./sslcerts/service.cert",
    "ssl.key.location": "./sslcerts/service.key",
}

# create a Kafka consumer instance
consumer = Consumer(conf)

Comme vous pouvez le voir, les premières lignes de code définissent la configuration à utiliser pour vous abonner à votre service Kafka.
N'oubliez pas de définir une variable d'environnement appelée KAFKA_SERVICE_URI qui pointera vers votre service.
La bibliothèque confluent-kafka fournit une classe appelée Consumer qui représentera votre connexion à votre service Kafka.

# create a console instance to display messages in a TUI
console = Console()

finished = False
local_count = 0

# subscribe to the "heroes" topic and wait for a message
consumer.subscribe(["heroes"])

Vous préparez ensuite l'outil qui vous aidera à afficher les messages de façon lisible.
Il ne vous reste plus qu'à utiliser l'objet Consumer pour vous abonner à votre service Kafka.

with console.status("Waiting for messages..."):
    while not finished:
        if (msg := consumer.poll(timeout=1.0)) is None:
            continue
        elif msg.error():
            raise KafkaException(msg.error())
        else:
            console.print(f"{msg.offset()}: {msg.key()}:{msg.value().decode()}\n\n")
            local_count += 1
            finished = local_count == 2

Le dernier morceau de code attendra les messages entrants via la fonction poll.
Vous pourrez utiliser l'objet Console pour afficher le contenu des messages.

Étape 2 - Publier des messages

Maintenant que vous avez une application qui attend des messages, créons-en une pour les produire et les publier.

import os
from confluent_kafka import Producer

# retrieve URI to connect to your Kafka service from your OS environment variables
KAFKA_SERVICE_URI = os.getenv("KAFKA_SERVICE_URI")
if KAFKA_SERVICE_URI is None:
    raise ValueError("KAFKA_SERVICE_URI is not set")

# set up the configuration to access your Kafka service in a secure way
conf = {
    "bootstrap.servers": KAFKA_SERVICE_URI,
    "client.id": "producer",
    "security.protocol": "SSL",
    "ssl.ca.location": "./sslcerts/ca.pem",
    "ssl.certificate.location": "./sslcerts/service.cert",
    "ssl.key.location": "./sslcerts/service.key",
}

# create a Kafka producer instance
producer = Producer(conf)

Cela se fait de manière très similaire à votre Consumer, mais cette fois vous utiliserez un objet Producer.

# when the message is published, this callback will be triggered
def delivery_callback(err, msg):
    if err:
        print(f"Message failed delivery: {err}")
    else:
        print(f"Published event to topic {msg.topic()} ")

# example data to send as a message
jsonValue = """{
        'id': 1,
        'name': 'Spider-Man (Peter Parker)',
        'description': 'Bitten by a radioactive spider, high school student Peter Parker gained the speed, strength and powers of a spider. Adopting the name Spider-Man, Peter hoped to start a career using his new abilities. Taught that with great power comes great responsibility, Spidey has vowed to use his powers to help people.',
    }"""

# create and publish the message to the "heroes" topic
producer.produce(
    "heroes", key="1", value=jsonValue.encode(), callback=delivery_callback
)
producer.flush()

Il est maintenant temps de préparer les éléments utilisés pour publier un message.
Le delivery_callback vous permet de garder le contrôle sur ce qu'il faut faire une fois que votre ```Producer`` a publié le message.
L'action de publication se déroule en fait en deux étapes :

  • Préparez d'abord le message dans le format requis par Kafka et définissez la fonction de callback.
  • Utilisez ensuite flush pour utiliser votre connexion à Kafka et publier le message.

Aller plus loin

Documentation officielle de Kafka

Bibliothèque Python Confluent Kafka

Rejoignez notre communauté d'utilisateurs.

Cette page vous a-t-elle aidé ?