//producer.php
$conf = new RdKafka\Conf();
$conf->set('bootstrap.servers', 'localhost:9092'); // Cambia según tu broker
$producer = new RdKafka\Producer($conf);
$topic = $producer->newTopic("mi-primer-topic");
echo "Producer PHP iniciado...\n";
for ($i = 1; $i <= 10; $i++) {
$mensaje = json_encode([
'id' => $i,
'mensaje' => "Hola desde PHP Kafka! Mensaje $i",
'timestamp' => time(),
'php_version' => phpversion()
]);
$topic->produce(RD_KAFKA_PARTITION_UA, 0, $mensaje);
echo "Mensaje enviado: $mensaje\n";
$producer->poll(0); // Necesario para procesar eventos
usleep(500000); // 0.5 segundos
}
$producer->flush(10000); // Espera máximo 10 segundos a que se envíen todos
echo "Todos los mensajes enviados.\n";
//consumer.php
$conf = new RdKafka\Conf();
$conf->set('bootstrap.servers', 'localhost:9092');
$conf->set('group.id', 'mi-grupo-php'); // Muy importante
$conf->set('auto.offset.reset', 'earliest'); // Lee desde el principio la primera vez
$consumer = new RdKafka\KafkaConsumer($conf);
$consumer->subscribe(['mi-primer-topic']);
echo "Consumer PHP esperando mensajes...\n";
while (true) {
$message = $consumer->consume(120 * 1000); // Timeout de 120 segundos
switch ($message->err) {
case RD_KAFKA_RESP_ERR_NO_ERROR:
$data = json_decode($message->payload, true);
echo "\n" . str_repeat("=", 60) . "\n";
echo "