| 123456789101112131415161718192021222324252627282930313233343536373839404142 |
- <?php
- namespace app\common\example;
- use app\common\library\Kafka;
- class KafkaExample
- {
- /**
- * 生产者示例
- */
- public function producer()
- {
- $kafka = Kafka::getInstance();
-
- $message = [
- 'id' => 1,
- 'event' => 'user_login',
- 'data' => [
- 'user_id' => 123,
- 'timestamp' => time()
- ]
- ];
- $kafka->publish('new_sdk_topic', $message);
- }
- /**
- * 消费者示例
- */
- public function consumer()
- {
- $kafka = Kafka::getInstance();
-
- $topics = ['new_sdk_topic'];
-
- $kafka->consume($topics, function($message) {
- $data = json_decode($message->payload, true);
- // 处理消息
- echo "Received message: " . $message->payload . "\n";
- });
- }
- }
|