Tous les produits
Search
Centre de documentation

ApsaraMQ for Kafka:Utiliser le SDK pour Python pour envoyer et recevoir des messages

Dernière mise à jour :Aug 11, 2026

Cette rubrique explique comment utiliser le SDK pour Python pour se connecter à ApsaraMQ for Kafka afin d'envoyer et de recevoir des messages sur un serveur Linux.

Avant de commencer

Installer la bibliothèque de dépendances Python

Exécutez la commande suivante pour installer la bibliothèque de dépendances Python :

pip install confluent-kafka==1.9.2
Important

Nous vous recommandons d'installer confluent-kafka version 1.9.2 ou antérieure. Dans le cas contraire, l'erreur SSL_HANDSHAKE s'affiche lors de l'envoi de messages via Internet.

Préparer un fichier de configuration

Téléchargez le projet de démonstration, modifiez les configurations correspondantes en fonction de l'endpoint que vous utilisez, puis téléversez le projet de démonstration sur le serveur Linux.

  1. Accédez à aliware-kafka-demos, cliquez sur l'icône image, puis sélectionnez Download ZIP dans la liste déroulante pour télécharger et décompresser le projet de démonstration.

    Remarque

    Le projet de démonstration téléchargé inclut le certificat racine SSL. Pour utiliser le certificat racine SSL séparément, téléchargez le certificat racine SSL.

  2. Dans le projet de démonstration décompressé, localisez le dossier kafka-confluent-python-demo et modifiez le fichier de configuration setting.py en fonction de l'endpoint que vous utilisez.

    Default endpoint

    Dans le répertoire vpc, modifiez le fichier de configuration setting.py.

    kafka_setting = {
        'bootstrap_servers': 'XXX:xxx,XXX:xxx',
        'topic_name': 'XXX',
        'group_name': 'XXX'
    }
    

    Parameter

    Description

    bootstrap_servers

    L'endpoint par défaut de l'instance ApsaraMQ for Kafka. Vous pouvez obtenir l'endpoint dans la section Endpoint Information de la page Instance Details dans la console ApsaraMQ for Kafka.

    topic_name

    Le nom du topic. Vous pouvez obtenir le nom du topic sur la page Topics dans la console ApsaraMQ for Kafka.

    group_name

    Le nom du group. Vous pouvez obtenir le nom du group sur la page Groups dans la console ApsaraMQ for Kafka.

    SSL endpoint

    Dans le répertoire vpc-ssl, modifiez le fichier de configuration setting.py.

    kafka_setting = {
        'sasl_plain_username': 'XXX',
        'sasl_plain_password': 'XXX',
        'ca_location': '/XXX/mix-4096-ca-cert',
        'bootstrap_servers': 'XXX:xxx,XXX:xxx',
        'topic_name': 'XXX',
        'group_name': 'XXX'
    }
    

    Parameter

    Description

    sasl_plain_username

    Le nom d'utilisateur de l'utilisateur Simple Authentication and Security Layer (SASL).

    Remarque
    • Si la fonctionnalité ACL n'est pas activée pour l'instance ApsaraMQ for Kafka, vous pouvez obtenir le nom d'utilisateur et le mot de passe de l'utilisateur SASL à partir des paramètres Username et Password dans la section Configuration Information de la page Instance Details dans la console ApsaraMQ for Kafka.

    • Si la fonctionnalité ACL est activée pour l'instance ApsaraMQ for Kafka, assurez-vous que l'utilisateur SASL est autorisé à envoyer et recevoir des messages via l'instance. Pour plus d'informations, consultez la rubrique Accorder des autorisations aux utilisateurs SASL.

    sasl_plain_password

    Le mot de passe de l'utilisateur SASL.

    ca_location

    Le chemin d'accès où le certificat racine SSL est enregistré. Remplacez XXX dans l'exemple de code par le chemin d'accès local. Exemple : /home/kafka-confluent-python-demo/vpc-ssl/mix-4096-ca-cert.

    bootstrap_servers

    L'endpoint SSL de l'instance ApsaraMQ for Kafka. Vous pouvez obtenir l'endpoint dans la section Endpoint Information de la page Instance Details dans la console ApsaraMQ for Kafka.

    topic_name

    Le nom du topic. Vous pouvez obtenir le nom du topic sur la page Topics dans la console ApsaraMQ for Kafka.

    group_name

    Le nom du group. Vous pouvez obtenir le nom du group sur la page Groups dans la console ApsaraMQ for Kafka.

  3. Téléversez le dossier kafka-confluent-python-demo vers le répertoire /home sur le serveur Linux.

Envoyer des messages

Envoyez des messages en fonction de l'endpoint que vous utilisez.

Default endpoint

  1. Exécutez la commande suivante pour accéder au sous-répertoire /home/kafka-confluent-python-demo/vpc :

    cd /home/kafka-confluent-python-demo/vpc
  2. Exécutez la commande suivante pour envoyer des messages :

    python kafka_producer.py

L'exemple de code suivant illustre le fichier kafka_producer.py :

kafka_producer.py

from confluent_kafka import Producer
import setting

conf = setting.kafka_setting
# Initialize a producer. 
p = Producer({'bootstrap.servers': conf['bootstrap_servers']})

def delivery_report(err, msg):
    """ Called once for each message produced to indicate delivery result.
        Triggered by poll() or flush(). """
    if err is not None:
        print('Message delivery failed: {}'.format(err))
    else:
        print('Message delivered to {} [{}]'.format(msg.topic(), msg.partition()))

# Send messages in asynchronous transmission mode. 
p.produce(conf['topic_name'], "Hello".encode('utf-8'), callback=delivery_report)
p.poll(0)

# When the program is ended, call the flush() method. 
p.flush()

SSL endpoint

  1. Exécutez la commande suivante pour accéder au sous-répertoire /home/kafka-confluent-python-demo/vpc-ssl :

    cd /home/kafka-confluent-python-demo/vpc-ssl
  2. Exécutez la commande suivante pour envoyer des messages :

    python kafka_producer.py

L'exemple de code suivant illustre le fichier kafka_producer.py :

kafka_producer.py

from confluent_kafka import Producer
import setting

conf = setting.kafka_setting

p = Producer({'bootstrap.servers':conf['bootstrap_servers'],
   'ssl.endpoint.identification.algorithm': 'none',
   'sasl.mechanisms':'PLAIN',
   'ssl.ca.location':conf['ca_location'],
   'security.protocol':'SASL_SSL',
   'sasl.username':conf['sasl_plain_username'],
   'sasl.password':conf['sasl_plain_password']})

def delivery_report(err, msg):
    if err is not None:
        print('Message delivery failed: {}'.format(err))
    else:
        print('Message delivered to {} [{}]'.format(msg.topic(), msg.partition()))

p.produce(conf['topic_name'], "Hello".encode('utf-8'), callback=delivery_report)
p.poll(0)

p.flush()

S'abonner aux messages

Abonnez-vous aux messages en fonction de l'endpoint que vous utilisez.

Default endpoint

  1. Exécutez la commande suivante pour accéder au sous-répertoire /home/kafka-confluent-python-demo/vpc :

    cd /home/kafka-confluent-python-demo/vpc
  2. Exécutez la commande suivante pour vous abonner aux messages :

    python kafka_consumer.py

L'exemple de code suivant illustre le fichier kafka_consumer.py :

kafka_consumer.py

from confluent_kafka import Consumer, KafkaError

import setting

conf = setting.kafka_setting

c = Consumer({
    'bootstrap.servers': conf['bootstrap_servers'],
    'group.id': conf['group_name'],
    'auto.offset.reset': 'latest'
})

c.subscribe([conf['topic_name']])

while True:
    msg = c.poll(1.0)

    if msg is None:
        continue
    if msg.error():
        if msg.error().code() == KafkaError._PARTITION_EOF:
            continue
        else:
            print("Consumer error: {}".format(msg.error()))
            continue

    print('Received message: {}'.format(msg.value().decode('utf-8')))

c.close()

SSL endpoint

  1. Exécutez la commande suivante pour accéder au sous-répertoire /home/kafka-confluent-python-demo/vpc-ssl :

    cd /home/kafka-confluent-python-demo/vpc-ssl
  2. Exécutez la commande suivante pour vous abonner aux messages :

    python kafka_consumer.py

L'exemple de code suivant illustre le fichier kafka_consumer.py :

kafka_consumer.py

from confluent_kafka import Consumer, KafkaError

import setting

conf = setting.kafka_setting

c = Consumer({
    'bootstrap.servers': conf['bootstrap_servers'],
    'ssl.endpoint.identification.algorithm': 'none',
    'sasl.mechanisms':'PLAIN',
    'ssl.ca.location':conf['ca_location'],
    'security.protocol':'SASL_SSL',
    'sasl.username':conf['sasl_plain_username'],
    'sasl.password':conf['sasl_plain_password'],
    'group.id': conf['group_name'],
    'auto.offset.reset': 'latest',
    'fetch.message.max.bytes':'1024*512'
})

c.subscribe([conf['topic_name']])

while True:
    msg = c.poll(1.0)

    if msg is None:
        continue
    if msg.error():
       if msg.error().code() == KafkaError._PARTITION_EOF:
          continue
       else:
           print("Consumer error: {}".format(msg.error()))
           continue

    print('Received message: {}'.format(msg.value().decode('utf-8')))

c.close()