KafkaTest.php 1.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869
  1. <?php
  2. namespace app\api\controller\v1;
  3. use app\api\controller\Api;
  4. use app\common\library\Kafka;
  5. use think\Request;
  6. class KafkaTest extends Api
  7. {
  8. /**
  9. * 测试发送消息
  10. */
  11. public function testProducer()
  12. {
  13. try {
  14. $kafka = Kafka::getInstance();
  15. // 构造测试消息
  16. $message = [
  17. 'id' => uniqid(),
  18. 'event' => 'test_event',
  19. 'data' => [
  20. 'user_id' => 12345,
  21. 'username' => 'test_user',
  22. 'action' => 'login',
  23. 'timestamp' => date('Y-m-d H:i:s')
  24. ]
  25. ];
  26. // 发送消息
  27. $kafka->publish('new_sdk_topic', $message);
  28. return $this->success('消息发送成功', [
  29. 'message' => $message
  30. ]);
  31. } catch (\Exception $e) {
  32. return $this->error('消息发送失败:' . $e->getMessage());
  33. }
  34. }
  35. /**
  36. * 测试自定义消息发送
  37. */
  38. public function sendCustomMessage()
  39. {
  40. $data = $this->request->post();
  41. try {
  42. $kafka = Kafka::getInstance();
  43. // 构造消息
  44. $message = [
  45. 'id' => uniqid(),
  46. 'event' => $data['event'] ?? 'custom_event',
  47. 'data' => $data['data'] ?? [],
  48. 'timestamp' => date('Y-m-d H:i:s')
  49. ];
  50. // 发送消息
  51. $kafka->publish($data['topic'] ?? 'new_sdk_topic', $message);
  52. return $this->success('自定义消息发送成功', [
  53. 'message' => $message
  54. ]);
  55. } catch (\Exception $e) {
  56. return $this->error('消息发送失败:' . $e->getMessage());
  57. }
  58. }
  59. }