Todos os produtos
Search
Central de documentação

ApsaraMQ for Kafka:Usar o SDK for Python para enviar e receber mensagens

Última atualização: Jun 27, 2026

Este tópico descreve como usar o SDK for Python para se conectar ao ApsaraMQ for Kafka e enviar e receber mensagens em um servidor Linux.

Antes de começar

Instale a biblioteca de dependências do Python

Execute o comando a seguir para instalar a biblioteca de dependências do Python:

pip install confluent-kafka==1.9.2
Importante

Recomendamos instalar a versão 1.9.2 ou anterior do confluent-kafka. Caso contrário, o erro SSL_HANDSHAKE será retornado ao enviar mensagens pela Internet.

Preparar um arquivo de configuração

Baixe o projeto de demonstração, modifique as configurações correspondentes com base no endpoint utilizado e envie o projeto para o servidor Linux.

  1. Acesse aliware-kafka-demos, clique em image e selecione Download ZIP na lista suspensa para baixar e descompactar o projeto de demonstração.

    Nota

    O projeto de demonstração baixado inclui o certificado raiz SSL. Para usar o certificado raiz SSL separadamente, baixe o certificado raiz SSL.

  2. No projeto de demonstração descompactado, localize a pasta kafka-confluent-python-demo e modifique o arquivo de configuração setting.py conforme o endpoint utilizado.

    Default endpoint

    No diretório vpc, modifique o arquivo de configuração setting.py.

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

    Parâmetro

    Descrição

    bootstrap_servers

    Endpoint padrão da instância do ApsaraMQ for Kafka. Obtenha o endpoint na seção Endpoint Information da página Instance Details no console do ApsaraMQ for Kafka.

    topic_name

    Nome do tópico. Consulte o nome do tópico na página Topics no console do ApsaraMQ for Kafka.

    group_name

    Nome do grupo. Verifique o nome do grupo na página Groups no console do ApsaraMQ for Kafka.

    SSL endpoint

    No diretório vpc-ssl, modifique o arquivo de configuração 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'
    }
    

    Parâmetro

    Descrição

    sasl_plain_username

    Nome de usuário do SASL (Simple Authentication and Security Layer).

    Nota
    • Se o recurso ACL não estiver ativado para a instância do ApsaraMQ for Kafka, obtenha o nome de usuário e a senha do SASL nos parâmetros Username e Password, na seção Configuration Information da página Instance Details no console do ApsaraMQ for Kafka.

    • Caso o recurso ACL esteja ativado na instância do ApsaraMQ for Kafka, certifique-se de que o usuário SASL tenha autorização para enviar e receber mensagens por meio da instância. Para mais informações, consulte Conceder permissões a usuários SASL.

    sasl_plain_password

    Senha do usuário SASL.

    ca_location

    Caminho onde o certificado raiz SSL está salvo. Substitua XXX no código de exemplo pelo caminho local. Exemplo: /home/kafka-confluent-python-demo/vpc-ssl/mix-4096-ca-cert.

    bootstrap_servers

    Endpoint SSL da instância do ApsaraMQ for Kafka. Obtenha o endpoint na seção Endpoint Information da página Instance Details no console do ApsaraMQ for Kafka.

    topic_name

    Nome do tópico. Consulte o nome do tópico na página Topics no console do ApsaraMQ for Kafka.

    group_name

    Nome do grupo. Consulte o nome do grupo na página Groups no console do ApsaraMQ for Kafka.

  3. Envie a pasta kafka-confluent-python-demo para o diretório /home no servidor Linux.

Enviar mensagens

Envie mensagens de acordo com o endpoint utilizado.

Default endpoint

  1. Execute o comando a seguir para acessar o subdiretório /home/kafka-confluent-python-demo/vpc:

    cd /home/kafka-confluent-python-demo/vpc
  2. Execute o comando a seguir para enviar mensagens:

    python kafka_producer.py

O código de exemplo a seguir mostra uma implementação do arquivo 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. Execute o comando a seguir para acessar o subdiretório /home/kafka-confluent-python-demo/vpc-ssl:

    cd /home/kafka-confluent-python-demo/vpc-ssl
  2. Execute o comando a seguir para enviar mensagens:

    python kafka_producer.py

O código de exemplo a seguir ilustra o conteúdo do arquivo 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()

Inscrever-se em mensagens

Inscreva-se nas mensagens conforme o endpoint utilizado.

Default endpoint

  1. Execute o comando a seguir para acessar o subdiretório /home/kafka-confluent-python-demo/vpc:

    cd /home/kafka-confluent-python-demo/vpc
  2. Execute o comando a seguir para se inscrever nas mensagens:

    python kafka_consumer.py

Veja a seguir um exemplo de implementação do arquivo 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. Execute o comando a seguir para acessar o subdiretório /home/kafka-confluent-python-demo/vpc-ssl:

    cd /home/kafka-confluent-python-demo/vpc-ssl
  2. Execute o comando a seguir para se inscrever nas mensagens:

    python kafka_consumer.py

Confira abaixo um exemplo do arquivo 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()