uniqid(), 'event' => 'test_event', 'data' => [ 'user_id' => 12345, 'username' => 'test_user', 'action' => 'login', 'timestamp' => date('Y-m-d H:i:s') ] ]; // 发送消息 $kafka->publish('new_sdk_topic', $message); return $this->success('消息发送成功', [ 'message' => $message ]); } catch (\Exception $e) { return $this->error('消息发送失败:' . $e->getMessage()); } } /** * 测试自定义消息发送 */ public function sendCustomMessage() { $data = $this->request->post(); try { $kafka = Kafka::getInstance(); // 构造消息 $message = [ 'id' => uniqid(), 'event' => $data['event'] ?? 'custom_event', 'data' => $data['data'] ?? [], 'timestamp' => date('Y-m-d H:i:s') ]; // 发送消息 $kafka->publish($data['topic'] ?? 'new_sdk_topic', $message); return $this->success('自定义消息发送成功', [ 'message' => $message ]); } catch (\Exception $e) { return $this->error('消息发送失败:' . $e->getMessage()); } } }