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
-
Nota
O Python 2.7 e 3.x são compatíveis. Este tópico utiliza o Python 3.9.
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
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.
-
Acesse aliware-kafka-demos, clique em
e selecione Download ZIP na lista suspensa para baixar e descompactar o projeto de demonstração.NotaO projeto de demonstração baixado inclui o certificado raiz SSL. Para usar o certificado raiz SSL separadamente, baixe o certificado raiz SSL.
-
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).
NotaSe 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.
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
-
Execute o comando a seguir para acessar o subdiretório /home/kafka-confluent-python-demo/vpc:
cd /home/kafka-confluent-python-demo/vpc -
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:
SSL endpoint
-
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 -
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:
Inscrever-se em mensagens
Inscreva-se nas mensagens conforme o endpoint utilizado.
Default endpoint
-
Execute o comando a seguir para acessar o subdiretório /home/kafka-confluent-python-demo/vpc:
cd /home/kafka-confluent-python-demo/vpc -
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:
SSL endpoint
-
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 -
Execute o comando a seguir para se inscrever nas mensagens:
python kafka_consumer.py
Confira abaixo um exemplo do arquivo kafka_consumer.py: