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.
-
Acesse o diretório de repositórios do yum:
cd /etc/yum.repos.d/ -
Crie um arquivo chamado
confluent.repocom 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 -
Instale o pacote de desenvolvimento da librdkafka:
yum install librdkafka-devel
Etapa 2: Instale a extensão php-rdkafka
-
Instale a extensão via PECL:
pecl install rdkafka -
Adicione a linha abaixo ao arquivo
php.inipara ativar a extensão:extension=rdkafka.so
Etapa 3: Configure os parâmetros de conexão
(Opcional) Se você usar o endpoint SSL, baixe o certificado raiz SSL.
Baixe o projeto de demonstração em aliware-kafka-demos e extraia o arquivo compactado.
-
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 arquivosetting.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.
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
Explore as opções de configuração e o uso da API do php-rdkafka no repositório GitHub do php-rdkafka.
Envie e receba mensagens com SDKs para outras linguagens. Consulte a Referência do Desenvolvedor do ApsaraMQ for Kafka.