Todos os produtos
Search
Central de documentação

ApsaraMQ for Kafka:Use the SDK for PHP to send and receive messages

Última atualização: Jun 27, 2026

Conecte uma aplicação PHP ao ApsaraMQ for Kafka para enviar e receber mensagens. Os exemplos neste guia usam a extensão php-rdkafka, um wrapper PHP para a biblioteca C librdkafka.

Pré-requisitos

Antes de começar, verifique se os seguintes softwares estão instalados no servidor Linux:

  • GNU Compiler Collection (GCC): necessário para compilar extensões nativas. Para instruções de instalação, consulte Installing GCC.

  • PHP: para instruções de download e instalação, consulte a página de downloads do PHP.

  • PHP Extension Community Library (PECL): usada para instalar extensões PHP. Para mais detalhes, consulte Downloading PECL extensions.

Você também precisa de:

  • Uma instância do ApsaraMQ for Kafka com um tópico e um grupo de consumidores criados. Obtenha o endpoint, o nome do tópico e o ID do grupo no console do ApsaraMQ for Kafka.

  • Se a conexão usar o endpoint SSL: o nome de usuário e a senha do Simple Authentication and Security Layer (SASL) da instância.

Etapa 1: Instale a librdkafka

A extensão php-rdkafka depende da biblioteca C librdkafka. Instale-a pelo repositório de pacotes da Confluent.

  1. Acesse o diretório de repositórios do yum:

       cd /etc/yum.repos.d/
  2. Crie um arquivo chamado confluent.repo com o seguinte conteúdo:

       [Confluent.dist]
       name=Confluent repository (dist)
       baseurl=https://packages.confluent.io/rpm/5.1/7
       gpgcheck=1
       gpgkey=https://packages.confluent.io/rpm/5.1/archive.key
       enabled=1
    
       [Confluent]
       name=Confluent repository
       baseurl=https://packages.confluent.io/rpm/5.1
       gpgcheck=1
       gpgkey=https://packages.confluent.io/rpm/5.1/archive.key
       enabled=1
  3. Instale o pacote de desenvolvimento da librdkafka:

       yum install librdkafka-devel

Etapa 2: Instale a extensão php-rdkafka

  1. Instale a extensão via PECL:

       pecl install rdkafka
  2. Adicione a linha abaixo ao arquivo php.ini para ativar a extensão:

       extension=rdkafka.so

Etapa 3: Configure os parâmetros de conexão

  1. (Opcional) Se você usar o endpoint SSL, baixe o certificado raiz SSL.

  2. Baixe o projeto de demonstração em aliware-kafka-demos e extraia o arquivo compactado.

  3. No pacote extraído, acesse a pasta kafka-php-demo. Abra a subpasta correspondente ao tipo de endpoint (SSL ou padrão) e edite o arquivo setting.php:

       <?php
    
       return [
           'sasl_plain_username' => '<YOUR_SASL_USERNAME>',
           'sasl_plain_password' => '<YOUR_SASL_PASSWORD>',
           'bootstrap_servers'   => '<HOST1>:<PORT1>,<HOST2>:<PORT2>',
           'topic_name'          => '<YOUR_TOPIC_NAME>',
           'consumer_id'         => '<YOUR_CONSUMER_GROUP_ID>',
       ];

    Substitua os marcadores pelos valores reais:

    Marcador

    Descrição

    Onde encontrar

    <YOUR_SASL_USERNAME>

    Nome de usuário SASL. Não necessário para o endpoint padrão.

    Instance Details > Configuration Information > Username no console do ApsaraMQ for Kafka. Se o recurso ACL estiver ativado, garanta que o usuário SASL tenha autorização para enviar e receber mensagens. Para mais informações, consulte Grant permissions to SASL users.

    <YOUR_SASL_PASSWORD>

    Senha SASL. Não necessária para o endpoint padrão.

    Mesmo local do nome de usuário.

    <HOST1>:<PORT1>,<HOST2>:<PORT2>

    Servidores bootstrap (endpoint da instância do ApsaraMQ for Kafka).

    Instance Details > Endpoint Information no console do ApsaraMQ for Kafka.

    <YOUR_TOPIC_NAME>

    Nome do tópico.

    Página Topics no console do ApsaraMQ for Kafka.

    <YOUR_CONSUMER_GROUP_ID>

    ID do grupo de consumidores.

    Página Groups no console do ApsaraMQ for Kafka.

  4. Faça upload de todos os arquivos da pasta para o diretório de instalação do PHP no servidor. Se usar o endpoint SSL, certifique-se de incluir o arquivo do certificado raiz SSL (ca-cert.pem).

Etapa 4: Enviar mensagens

Execute kafka-producer.php para enviar uma mensagem:

php kafka-producer.php

Endpoint SSL

Use este código ao conectar-se pelo endpoint SSL com autenticação SASL:

<?php

$setting = require __DIR__ . '/setting.php';

$conf = new RdKafka\Conf();

// SASL authentication
$conf->set('sasl.mechanisms', 'PLAIN');
$conf->set('sasl.username', $setting['sasl_plain_username']);
$conf->set('sasl.password', $setting['sasl_plain_password']);

// SSL/TLS encryption
$conf->set('security.protocol', 'SASL_SSL');
$conf->set('ssl.ca.location', __DIR__ . '/ca-cert.pem');
$conf->set('ssl.endpoint.identification.algorithm', 'none');

// Producer settings
$conf->set('api.version.request', 'true');
$conf->set('message.send.max.retries', 5);

$producer = new RdKafka\Producer($conf);
$producer->setLogLevel(LOG_INFO);  // Set to LOG_DEBUG for troubleshooting
$producer->addBrokers($setting['bootstrap_servers']);

$topic = $producer->newTopic($setting['topic_name']);

// Send a message to an automatically assigned partition
$topic->produce(RD_KAFKA_PARTITION_UA, 0, "Message hello kafka");
$producer->poll(0);

// Wait for all outstanding messages to be delivered
while ($producer->getOutQLen() > 0) {
    $producer->poll(50);
}

echo "send succ" . PHP_EOL;

Endpoint padrão

Use este código ao conectar-se pelo endpoint padrão sem autenticação SASL:

<?php

$setting = require __DIR__ . '/setting.php';

$conf = new RdKafka\Conf();

// Producer settings
$conf->set('api.version.request', 'true');
$conf->set('message.send.max.retries', 5);

$producer = new RdKafka\Producer($conf);
$producer->setLogLevel(LOG_INFO);  // Set to LOG_DEBUG for troubleshooting
$producer->addBrokers($setting['bootstrap_servers']);

$topic = $producer->newTopic($setting['topic_name']);

// Send a message to an automatically assigned partition
$topic->produce(RD_KAFKA_PARTITION_UA, 0, "Message hello kafka");
$producer->poll(0);

// Wait for all outstanding messages to be delivered
while ($producer->getOutQLen() > 0) {
    $producer->poll(50);
}

echo "send succ" . PHP_EOL;

Para mais informações sobre a API de produtor do php-rdkafka, consulte php-rdkafka.

Etapa 5: Inscrever-se em mensagens

Execute kafka-consumer.php para iniciar o consumo de mensagens:

php kafka-consumer.php

Endpoint SSL

Use este código ao conectar-se pelo endpoint SSL com autenticação SASL:

<?php

$setting = require __DIR__ . '/setting.php';

$conf = new RdKafka\Conf();

// SASL authentication
$conf->set('sasl.mechanisms', 'PLAIN');
$conf->set('sasl.username', $setting['sasl_plain_username']);
$conf->set('sasl.password', $setting['sasl_plain_password']);

// SSL/TLS encryption
$conf->set('security.protocol', 'SASL_SSL');
$conf->set('ssl.ca.location', __DIR__ . '/ca-cert.pem');
$conf->set('ssl.endpoint.identification.algorithm', 'none');

// Consumer settings
$conf->set('api.version.request', 'true');
$conf->set('group.id', $setting['consumer_id']);
$conf->set('session.timeout.ms', 10000);
$conf->set('request.timeout.ms', 305000);
$conf->set('metadata.broker.list', $setting['bootstrap_servers']);

$topicConf = new RdKafka\TopicConf();
$conf->setDefaultTopicConf($topicConf);

$consumer = new RdKafka\KafkaConsumer($conf);
$consumer->subscribe([$setting['topic_name']]);

echo "Waiting for partition assignment... (may take some time when" . PHP_EOL;
echo "quickly re-joining the group after leaving it.)" . PHP_EOL;

while (true) {
    $message = $consumer->consume(30 * 1000);  // 30-second poll timeout
    switch ($message->err) {
        case RD_KAFKA_RESP_ERR_NO_ERROR:
            // Successfully consumed a message
            var_dump($message);
            break;
        case RD_KAFKA_RESP_ERR__PARTITION_EOF:
            echo "No more messages; will wait for more" . PHP_EOL;
            break;
        case RD_KAFKA_RESP_ERR__TIMED_OUT:
            echo "Timed out" . PHP_EOL;
            break;
        default:
            throw new \Exception($message->errstr(), $message->err);
            break;
    }
}

?>

Endpoint padrão

Use este código ao conectar-se pelo endpoint padrão sem autenticação SASL:

<?php

$setting = require __DIR__ . '/setting.php';

$conf = new RdKafka\Conf();

// Consumer settings
$conf->set('api.version.request', 'true');
$conf->set('group.id', $setting['consumer_id']);
$conf->set('session.timeout.ms', 10000);
$conf->set('request.timeout.ms', 305000);
$conf->set('metadata.broker.list', $setting['bootstrap_servers']);

$topicConf = new RdKafka\TopicConf();
$conf->setDefaultTopicConf($topicConf);

$consumer = new RdKafka\KafkaConsumer($conf);
$consumer->subscribe([$setting['topic_name']]);

echo "Waiting for partition assignment... (may take some time when" . PHP_EOL;
echo "quickly re-joining the group after leaving it.)" . PHP_EOL;

while (true) {
    $message = $consumer->consume(30 * 1000);  // 30-second poll timeout
    switch ($message->err) {
        case RD_KAFKA_RESP_ERR_NO_ERROR:
            // Successfully consumed a message
            var_dump($message);
            break;
        case RD_KAFKA_RESP_ERR__PARTITION_EOF:
            echo "No more messages; will wait for more" . PHP_EOL;
            break;
        case RD_KAFKA_RESP_ERR__TIMED_OUT:
            echo "Timed out" . PHP_EOL;
            break;
        default:
            throw new \Exception($message->errstr(), $message->err);
            break;
    }
}

?>

Para mais informações sobre a API de consumidor do php-rdkafka, consulte php-rdkafka.

Próximos passos