即时通讯项目---消息转发子服务
前言
本篇是我介绍Chat-Im(仿微信的即时通讯项目)中的一篇,大家可以有兴趣的可以参照我的gitee代码和博客自己动手试一试。
gitee:https://gitee.com/qi-haozhe/chat-im
点我直接转到gitee

本篇主要介绍消息转发子服务的实现,具体来说就是client在一个聊天会话中发送了一条消息,server收到之后就要把该消息转发给聊天会话中的人。该服务要做的就是找到聊天会话中有些谁,然后交给网关服务去逐个进行消息转发。
一、功能设计
转发子服务,主要用于针对一条消息内容,组织消息的 ID 以及各项所需要素,然后告诉网关服务器一条消息应该发给谁。
通常消息都是以聊天会话为基础进行发送的,根据会话找到它的所有成员,这就是转发的目标。聊天会话又分为单人会话和群聊,单人会话只需要把消息转发给一个人即可,群聊会话的话,需要把该群聊中的所有人全部获取。
除此之外,转发子服务将收到的消息,放入消息队列中,由消息存储管理子服务进行消费存储
- 获取消息转发目标:针对消息内容,组织消息,并告知网关转发目标。
二、模块划分
- 参数/配置文件解析模块:基于 gflags 框架直接使用进行参数/配置文件解析。
- 日志模块:基于 spdlog 框架封装的模块直接使用进行日志输出。
- 服务注册模块:基于 etcd 框架封装的注册模块直接使用进行消息转发服务的服务注册。
- 数据库数据操作模块:基于 odb-mysql 数据管理封装的模块,从数据库获取会话成员。
- 服务发现与调用模块:基于 etcd 框架与 brpc 框架封装的服务发现与调用模块,从用户子服务获取消息发送者的用户信息。
- rpc 服务模块:基于 brpc 框架搭建 rpc 服务器。
- MQ 发布模块:基于 rabbitmq-client 封装的模块将消息发布到消息队列,让消息存储子服务进行消费,对消息进行存储。

三、前置操作
1. odb文件
对应gitee文件:https://gitee.com/qi-haozhe/chat-im/blob/master/server/odb/chat_session_member.hxx
点我直接跳转
该文件定义了一个名为 ChatSessionMember 的 C++ 类,它是用于表示聊天会话成员的数据模型,并通过 ODB(ORM 工具)映射到数据库表。
数据模型定义:
- 表示聊天会话中的成员关系,记录哪个用户(user_id)属于哪个会话(session_id)。
- 对应数据库表 chat_session_member。
ODB ORM 映射:
- 通过 #pragma db object 指令将类映射到数据库表。
- 成员变量通过 ODB 指令(如 #pragma db id)定义数据库字段属性。
//chat_session_member.hxx
// 聊天会话成员表映射对象
#pragma once
#include <string>
#include <cstddef>
#include <odb/core.hxx>
namespace im
{
#pragma db object table("chat_session_member")
class ChatSessionMember
{
public:
ChatSessionMember() {}
ChatSessionMember(const std::string &ssid, const std::string &uid) : _session_id(ssid), _user_id(uid) {}
~ChatSessionMember() {}
std::string session_id() const { return _session_id; }
void session_id(std::string &ssid) { _session_id = ssid; }
std::string user_id() const { return _user_id; }
void user_id(std::string &uid) { _user_id = uid; }
private:
friend class odb::access;
#pragma db id auto
unsigned long _id;
#pragma db type("varchar(64)") index
std::string _session_id;
#pragma db type("varchar(64)")
std::string _user_id;
};
}
2. sql文件
这是执行获得的sql文件
//odb -d mysql --std c++11 --generate-query --generate-schema --profile boost/date-time chat_session_member.hxx
/* This file was generated by ODB, object-relational mapping (ORM)
* compiler for C++.
*/
DROP TABLE IF EXISTS `chat_session_member`;
CREATE TABLE `chat_session_member` (
`id` BIGINT UNSIGNED NOT NULL PRIMARY KEY AUTO_INCREMENT,
`session_id` varchar(64) NOT NULL,
`user_id` varchar(64) NOT NULL)
ENGINE=InnoDB;
CREATE INDEX `session_id_i`
ON `chat_session_member` (`session_id`);
建立的表如图所示:

3. MySQL表增删改查文件
对应gitee代码:https://gitee.com/qi-haozhe/chat-im/blob/master/server/common/mysql_chat_session_member.hpp
点我直接跳转
这个文件定义了一个 ChatSessionMemeberTable 类,用于封装对聊天会话成员表(chat_session_member)的数据库操作。
核心功能
-
数据库操作封装:
- 提供对
chat_session_member表的 CRUD 操作(增删查)。 - 基于 ODB ORM 框架,通过
odb::database执行持久化操作。
- 提供对
-
主要接口:
- 新增成员:
append(ChatSessionMember&):插入单个成员。append(vector<ChatSessionMember>&):批量插入成员。
- 删除成员:
remove(ChatSessionMember&):删除指定会话中的指定成员。remove(const string& ssid):删除会话的所有成员。
- 查询成员:
members(const string& ssid):获取会话的所有成员用户ID列表。
- 新增成员:
代码设计分析
- 事务管理
- 每个操作通过
odb::transaction包裹,确保原子性:odb::transaction trans(_db->begin()); // 数据库操作... trans.commit(); - 异常时自动回滚(通过 RAII 机制)。
- 错误处理
- 捕获
std::exception并记录错误日志(通过LOG_ERROR):catch (std::exception &e) { LOG_ERROR("操作失败: {}!", e.what()); return false; } - 查询失败时返回空结果(如
members()返回空vector)。
- 查询优化
- 使用 ODB 的查询接口(
odb::query)按条件筛选:_db->query<ChatSessionMember>(query::session_id == ssid); - 对
session_id建立索引(见前文模型定义),提升查询性能。
/**
* @file mysql_chat_session_member.hpp
* @brief 封装对于会话中成员的管理操作
* @author qhz (2695432062@qq.com)
*/
#pragma once
#include "mysql.hpp"
#include "chat_session_member.hxx"
#include "chat_session_member-odb.hxx"
#include "logger.hpp"
namespace im
{
class ChatSessionMemeberTable
{
public:
using ptr = std::shared_ptr<ChatSessionMemeberTable>;
ChatSessionMemeberTable(const std::shared_ptr<odb::core::database> &db) : _db(db) {}
// 单个会话成员的新增 --- ssid & uid
bool append(ChatSessionMember &csm)
{
try
{
odb::transaction trans(_db->begin());
_db->persist(csm);
trans.commit();
}
catch (std::exception &e)
{
LOG_ERROR("新增单会话成员失败 {}-{}:{}!",
csm.session_id(), csm.user_id(), e.what());
return false;
}
return true;
}
bool append(std::vector<ChatSessionMember> &csm_lists)
{
try
{
odb::transaction trans(_db->begin());
for (auto &csm : csm_lists)
{
_db->persist(csm);
}
trans.commit();
}
catch (std::exception &e)
{
LOG_ERROR("新增多会话成员失败 {}-{}:{}!",
csm_lists[0].session_id(), csm_lists.size(), e.what());
return false;
}
return true;
}
// 删除指定会话中的指定成员 -- ssid & uid
bool remove(ChatSessionMember &csm)
{
try
{
odb::transaction trans(_db->begin());
typedef odb::query<ChatSessionMember> query;
typedef odb::result<ChatSessionMember> result;
_db->erase_query<ChatSessionMember>(query::session_id == csm.session_id() &&
query::user_id == csm.user_id());
trans.commit();
}
catch (std::exception &e)
{
LOG_ERROR("删除单会话成员失败 {}-{}:{}!",
csm.session_id(), csm.user_id(), e.what());
return false;
}
return true;
}
// 删除会话的所有成员信息
bool remove(const std::string &ssid)
{
try
{
odb::transaction trans(_db->begin());
typedef odb::query<ChatSessionMember> query;
typedef odb::result<ChatSessionMember> result;
_db->erase_query<ChatSessionMember>(query::session_id == ssid);
trans.commit();
}
catch (std::exception &e)
{
LOG_ERROR("删除会话所有成员失败 {}:{}!", ssid, e.what());
return false;
}
return true;
}
std::vector<std::string> members(const std::string &ssid)
{
std::vector<std::string> res;
try
{
odb::transaction trans(_db->begin());
typedef odb::query<ChatSessionMember> query;
typedef odb::result<ChatSessionMember> result;
result r(_db->query<ChatSessionMember>(query::session_id == ssid));
for (result::iterator i(r.begin()); i != r.end(); ++i)
{
res.push_back(i->user_id());
}
trans.commit();
}
catch (std::exception &e)
{
LOG_ERROR("获取会话成员失败:{}-{}!", ssid, e.what());
}
return res;
}
private:
std::shared_ptr<odb::core::database> _db;
};
}
四、接口实现流程
对应gitee文件:https://gitee.com/qi-haozhe/chat-im/blob/master/server/transmit/source/transmite_server.hpp
点我直接跳转
1. 实现流程
获取消息转发目标与消息处理
- 从请求中取出消息内容,会话 ID, 用户 ID
- 根据用户 ID 从用户子服务获取当前发送者用户信息
- 根据消息内容构造完成的消息结构(分配消息 ID,填充发送者信息,填充消息产生时间)
- 将消息序列化后发布到 MQ 消息队列中,让消息存储子服务对消息进行持久化存储
- 从数据库获取目标会话所有成员 ID
- 组织响应(完整消息+目标用户 ID),发送给网关,告知网关该将消息发送给谁。
只需要重写MsgTransmitService服务即可,在该服务内部,我们需要对发来的消息进行进一步的加工,比如填充该条消息发送者的信息,发送时间,该条消息属于哪个聊天会话等等,所以需要以下字段。
而GetTransmitTargetRsp则需要把封装好的消息进行返回并把要转发的成员列表发送回去,所以需要message字段和target_id_list字段。
//transmite.proto
//这个用于和网关进行通信
message NewMessageReq {
string request_id = 1; //请求ID -- 全链路唯一标识
optional string user_id = 2;
optional string session_id = 3;//客户端身份识别信息 -- 这就是消息发送者
string chat_session_id = 4; //聊天会话ID -- 标识了当前消息属于哪个会话,应该转发给谁
MessageContent message = 5; // 消息内容--消息类型+内容
}
message NewMessageRsp {
string request_id = 1;
bool success = 2;
string errmsg = 3;
}
//这个用于内部的通信,生成完整的消息信息,并获取消息的转发人员列表
message GetTransmitTargetRsp {
string request_id = 1;
bool success = 2;
string errmsg = 3;
MessageInfo message = 4; // 组织好的消息结构 --
repeated string target_id_list = 5; //消息的转发目标列表
}
service MsgTransmitService {
rpc GetTransmitTarget(NewMessageReq) returns (GetTransmitTargetRsp);
}
流程主要就是先把uid和消息内容从req中拿出来,通过uid去请求user子服务获取个人信息,然后对消息内容做一个封装,添加发送者、发送时间、属于哪个聊天会话等等。然后从mysql中去查找chat_session_member表,该表中有个session_id是会话id,属于同一个会话中的id,session_id是一样的,比如a和b是好友,他俩聊天了有一个session_id是c,那么数据库里就有两条记录如下:
id session_id user_id
1 c a
2 c b
标识a、b属于同一个会话。
查完该表之后,把session_id相同的uid全拿出来,然后存到rsp中,把消息添加到消息队列中为了让消息持久化服务从消息队列中获取消息进行持久化,然后返回rsp即可,返回之后网关会根据转发消息列表进行逐个转发消息的。

void GetTransmitTarget(google::protobuf::RpcController* controller,
const ::ymm_im::NewMessageReq* request,
::ymm_im::GetTransmitTargetRsp* response,
::google::protobuf::Closure* done) override {
brpc::ClosureGuard rpc_guard(done);
auto err_response = [this, response](const std::string &rid,
const std::string &errmsg) -> void {
response->set_request_id(rid);
response->set_success(false);
response->set_errmsg(errmsg);
return;
};
//从请求中获取关键信息:用户ID,所属会话ID,消息内容
std::string rid = request->request_id();
std::string uid = request->user_id();
std::string chat_ssid = request->chat_session_id();
const ymm_im::MessageContent &content = request->message();
// 进行消息组织:发送者-用户子服务获取信息,所属会话,消息内容,产生时间,消息ID
auto channel = _mm_channels->choose(_user_service_name);
if (!channel) {
LOG_ERROR("{}-{} 没有可供访问的用户子服务节点!", rid, _user_service_name);
return err_response(rid, "没有可供访问的用户子服务节点!");
}
ymm_im::UserService_Stub stub(channel.get());
ymm_im::GetUserInfoReq req;
ymm_im::GetUserInfoRsp rsp;
req.set_request_id(rid);
req.set_user_id(uid);
brpc::Controller cntl;
stub.GetUserInfo(&cntl, &req, &rsp, nullptr);
if (cntl.Failed() == true || rsp.success() == false) {
LOG_ERROR("{} - 用户子服务调用失败:{}!", request->request_id(), cntl.ErrorText());
return err_response(request->request_id(), "用户子服务调用失败!");
}
LOG_DEBUG("{} - 用户子服务调用成功,获取用户信息:{}", request->request_id(), rsp.user_info().ShortDebugString());
ymm_im::MessageInfo message;
message.set_message_id(uuid());
message.set_chat_session_id(chat_ssid);
message.set_timestamp(time(nullptr));
message.mutable_sender()->CopyFrom(rsp.user_info());
message.mutable_message()->CopyFrom(content);
// 获取消息转发客户端用户列表
auto target_list = _mysql_session_member_table->members(chat_ssid);
// 将封装完毕的消息,发布到消息队列,待消息存储子服务进行消息持久化
bool ret = _mq_client->publish(_exchange_name, message.SerializeAsString(), _routing_key);
if (ret == false) {
LOG_ERROR("{} - 持久化消息发布失败:{}!", request->request_id(), cntl.ErrorText());
return err_response(request->request_id(), "持久化消息发布失败:!");
}
LOG_DEBUG("{} - 持久化消息发布成功,消息内容:{}", request->request_id(), message.ShortDebugString());
//组织响应
response->set_request_id(rid);
response->set_success(true);
response->mutable_message()->CopyFrom(message);
for (const auto &id : target_list) {
response->add_target_id_list(id);
}
}
更多推荐



所有评论(0)