基于脚手架微服务的视频点播系统-服务端开发部分[一]文件子服务实现
基于脚手架微服务的视频点播系统-服务端开发部分[一]文件子服务实现
项目源码 gitee地址
一.预备工作
上篇文章我们将服务端需要的数据库表和rpc通信以及缓存都给打了个样板。但实际上博主在实际开发时发现还需要进行部分修改,所以所有的数据库表和缓存及rpc通信设计都会进行些微修改,所以下面用到的时候我会把相关设计再给出一遍但凡是修改过的,就认为上篇文章是用来大致理解我们要实现哪些东西就行。
1.1双写策略实现Mysql数据库与Redis缓存之间的数据同步
因为我们上篇文章提到过,为什么可以参考上篇文章这里就不多赘述了,所以这里我们来实现下。
实现思路大致是这样的,因为我们需要对缓存与消息队列进行操作,所以需要有两个成员,一个缓存操作对象,一个消息队列客户端。再结合这个类的功能来看:外部传入一个key值,内部进行第一次删除之后将此key值的消息传入消息队列,等待一定时间后此类接收此消息再次对缓存进行此key值的清除。因此消息队列客户端需要被用来创造两个对象,一个消息发布者,一个消息订阅者,显然此时需要订阅一个延时队列。让消息发布者发布消息到普通队列,订阅者订阅死信队列,所以我们先来定义缓存及缓存删除相关的交换机,队列及绑定key(因为此服务用到缓存的只有会话表,所以这里的缓存设计是针对会话表而言的):
会话表缓存设计
过期时间:60~120min
同步策略:双写同步策略(通过key删除缓存,并发布延迟⼆次删除消息)
缓存类型:hash
- key:session_SESSION_ID
- 会话ID
- 用户ID:该字段是否存在,用于判断会话是否是临时会话
数据同步缓存删除类设计
| 交换机 | delete_cache_exchange | 延时消息类型(3s) |
|---|---|---|
| 队列 | delete_cache_queue | |
| BindingKey | delete_cache_queue | |
| 具体实现如下: |
//remove_cache.h
#pragma once
#include <lime_scaffold/limemq.h>
#include <lime_scaffold/limeredis.h>
#include <lime_scaffold/limelog.h>
#include "message.pb.h"
//实现双写缓冲区的删除操作
namespace limevp_data{
class RemoveCache{
public:
using Ptr = std::shared_ptr<RemoveCache>;
RemoveCache(const std::shared_ptr<sw::redis::Redis>& redis,
const limemq::MQClient::ptr& mqClient,
const limemq::declare_settings& mqsettings);
void sync(const std::string& key);//给入缓存的key,内部进行延时删除
private:
//处理删除缓存的消息回调函数
void callback(const char* body, size_t len);
private:
std::shared_ptr<sw::redis::Redis> _redis;//redis操作句柄
limemq::Publisher::ptr _publisher;//消息队列发布者
limemq::Subscriber::ptr _subscriber;//消息队列订阅者
};
}// namespace limevp_data
//remove_cache.cc
#include "remove_cache.h"
namespace limevp_data{
RemoveCache::RemoveCache(const std::shared_ptr<sw::redis::Redis>& redis,const limemq::MQClient::ptr& mqClient,const limemq::declare_settings& mqsettings)
: _redis(redis)
{
//初始化消息队列发布者和订阅者
_publisher = std::make_shared<limemq::Publisher>(mqClient, mqsettings);
_subscriber = std::make_shared<limemq::Subscriber>(mqClient, mqsettings);
//设置订阅者的回调函数
_subscriber->consume(std::bind(&RemoveCache::callback, this, std::placeholders::_1, std::placeholders::_2));
}
void RemoveCache::sync(const std::string& key)
{
//删除缓存中的key
auto rtx = _redis->transaction(false,false);
auto r = rtx.redis();
r.del(key);
//使用protobuf序列化消息体
CommunicationInterface::DeleteCacheMsg msg;
msg.add_key(key);
//发布消息
_publisher->publish(msg.SerializeAsString());
}
void RemoveCache::callback(const char* body, size_t len)
{
//解析消息体
CommunicationInterface::DeleteCacheMsg msg;
auto ret = msg.ParseFromArray(body, len);
if(!ret)
{
ERR("收到删除缓存延迟消息,但是消息反序列化失败!");
return;
}
//删除缓存中的key
auto rtx = _redis->transaction(false,false);
auto r = rtx.redis();
for(auto& key : msg.key())
{
r.del(key);
DBG("收到删除缓存延迟消息,删除缓存key: {}", key);
}
}
}//namespace limevp_data
二.数据库表定义及数据操作类实现
2.1文件元信息表
2.1.1数据库表设计与定义
首先文件元信息表的设计如下:
| 字段 | 类型 | 约束 | 空值 | 备注 |
|---|---|---|---|---|
| ID | BIGINT UNSIGNED | PK | NOT NULL | |
| 文件ID | VARCHAR(64) | UK | NOT NULL | |
| 上传用户ID | VARCHAR(64) | NK | NOT NULL | |
| 存储路径 | VARCHAR(128) | NOT NULL | FASTDFS存储ID | |
| 文件大小 | INT UNSIGNED | NOT NULL | ||
| 文件mime | VARCHAR(16) | NOT NULL | Content-Type |
接下来通过我们之前了解的odb操作来进行文件元信息表的定义,之后我们所有的数据库表定义均放到data目录下的data.h/data.cc下:
//data.h
#pragma once
#include <string>
#include <cstddef> // std::size_t
#include <boost/date_time/posix_time/posix_time.hpp>
#include <odb/nullable.hxx>
#include <odb/core.hxx>
#include <memory>
//定义所有的数据库表结构
namespace limevp_data{
/*文件元信息表*/
#pragma db object table("tbl_file_metadata")
class File{
public:
using Ptr = std::shared_ptr<File>;
File();
File(const std::string& file_id, const std::string& upload_user_id, const std::string& file_path, unsigned int file_size, const std::string& file_mime);
unsigned long get_id();
std::string get_file_id() const;
void set_file_id(const std::string& file_id);
std::string get_upload_user_id() const;
void set_upload_user_id(const std::string& upload_user_id);
std::string get_file_path() const;
void set_file_path(const std::string& file_path);
unsigned int get_file_size() const;
void set_file_size(unsigned int file_size);
std::string get_file_mime() const;
void set_file_mime(const std::string& file_mime);
private:
friend class odb::access;
#pragma db id auto column("id") type("BIGINT UNSIGNED")
unsigned long _id;
#pragma db column("file_id") type("VARCHAR(64)") not_null unique
std::string _file_id;
#pragma db column("upload_user_id") type("VARCHAR(64)") not_null index
std::string _upload_user_id;
#pragma db column("file_path") type("VARCHAR(128)") not_null
std::string _file_path;
#pragma db column("file_size") type("INT UNSIGNED") not_null
unsigned int _file_size;
#pragma db column("file_mime") type("VARCHAR(16)") not_null
std::string _file_mime;
};
}//namespace limevp_data
//data.cc
#include "data.h"
namespace limevp_data{
/*file metadata start*/
File::File() {}
File::File(const std::string& file_id, const std::string& upload_user_id, const std::string& file_path, unsigned int file_size, const std::string& file_mime)
: _file_id(file_id)
, _upload_user_id(upload_user_id)
, _file_path(file_path)
, _file_size(file_size)
, _file_mime(file_mime)
{}
unsigned long File::get_id() {return _id;}
std::string File::get_file_id() const {return _file_id;}
void File::set_file_id(const std::string& file_id) {_file_id = file_id;}
std::string File::get_upload_user_id() const {return _upload_user_id;}
void File::set_upload_user_id(const std::string& upload_user_id) {_upload_user_id = upload_user_id;}
std::string File::get_file_path() const {return _file_path;}
void File::set_file_path(const std::string& file_path) {_file_path = file_path;}
unsigned int File::get_file_size() const {return _file_size;}
void File::set_file_size(unsigned int file_size) {_file_size = file_size;}
std::string File::get_file_mime() const {return _file_mime;}
void File::set_file_mime(const std::string& file_mime) {_file_mime = file_mime;}
/*file metadata end*/
}//namespace limevp_data
2.1.2数据表操作类实现
因为文件元信息不需要添加到缓存中,所以我们这里实现此类时只需要给我这个类传入一个odb的mysql数据库操作句柄即可,需要注意是odb事务生成的数据库操作句柄,确保整个类在操作时处于同一个事务中同时外部使用时直接通过异常来判断并处理数据库表操作时产生的异常信息:
此实现位置为data/file.h-file.cc
//file.h
#pragma once
#include <lime_scaffold/limeodb.h>
#include "data.h"
#include "data-odb.hxx"
namespace limevp_data{
class FileData{
public:
using Ptr = std::shared_ptr<FileData>;
FileData(odb::database& db);
//新增文件信息
void addFile2Db(File& file);
//获取文件信息
File::Ptr getFileFromDb(const std::string& fileId);
//修改文件信息
void updateFile2Db(File& file);
//删除文件信息
void delFromDb(const std::string& fileId);
private:
//注意这里是一个引用对象
odb::database& _db;
};
}//namespace limevp_data
//file.cc
#include "file.h"
namespace limevp_data{
FileData::FileData(odb::database& db)
: _db(db)
{}
//新增文件信息
void FileData::addFile2Db(File& file)
{
_db.persist(file);
}
//获取文件信息
File::Ptr FileData::getFileFromDb(const std::string& fileId)
{
File::Ptr result(_db.query_one<File>(odb::query<File>::file_id == fileId));
return result;
}
//修改文件信息-必须是getFileFromDb获取到的对象
void FileData::updateFile2Db(File& file)
{
//先查询,再修改,如果没有此对象则返回
auto result = getFileFromDb(file.get_file_id());
if(!result)
{
return;
}
//更新对应数据
result->set_upload_user_id(file.get_upload_user_id());
result->set_file_path(file.get_file_path());
result->set_file_size(file.get_file_size());
result->set_file_mime(file.get_file_mime());
_db.update(result.get());
}
//删除文件信息
void FileData::delFromDb(const std::string& fileId)
{
_db.erase_query<File>(odb::query<File>::file_id == fileId);
}
}// namespace limevp_data
2.1.3文件表操作测试
#include "file.h"
#include <lime_scaffold/limelog.h>
const std::string USER = "root";
const std::string PASSWORD = "123456";
const std::string DATABASE = "vbtest";
const std::string HOST = "192.168.30.128";
const unsigned int PORT = 3306;
const std::string CHRSET = "utf8";
int main() {
limelog::limelog_init();
limeodb::mysql_settings settings{
.host = HOST,
.user = USER,
.password = PASSWORD,
.database = DATABASE,
.port = PORT,
.charset = CHRSET
};
std::shared_ptr<odb::mysql::database> handler = limeodb::DbFactory::create_mysqldb(settings);
try
{
// {
// odb::mysql::transaction t(handler->begin(),true);//后面的bool参数表示是否是在当前线程中创建事务对象
// //获取数据库操作句柄
// auto& db = t.database();
// limevp_data::FileData file_data(db);
// //新增文件信息
// limevp_data::File file("file123","user123","/path/to/file.mp4",1048576,"video/mp4");
// file_data.addFile2Db(file);
// DBG("新增文件信息成功,文件ID: {}", file.get_file_id());
// t.commit();
// }
// {
// odb::mysql::transaction t(handler->begin(),true);//后面的bool参数表示是否是在当前线程中创建事务对象
// //获取数据库操作句柄
// auto& db = t.database();
// limevp_data::FileData file_data(db);
// //修改文件信息
// limevp_data::File file("file123","user456","/path/to/file.mp4",1048576,"video/mp4");
// file_data.updateFile2Db(file);
// t.commit();
// }
// {
// odb::mysql::transaction t(handler->begin(),true);//后面的bool参数表示是否是在当前线程中创建事务对象
// //获取数据库操作句柄
// auto& db = t.database();
// limevp_data::FileData file_data(db);
// //获取文件信息
// auto file_ptr = file_data.getFileFromDb("file123");
// if(file_ptr)
// {
// DBG("获取文件信息成功,文件ID: {}, 上传用户ID: {}, 文件路径: {}, 文件大小: {}, 文件MIME: {}",
// file_ptr->get_file_id(),
// file_ptr->get_upload_user_id(),
// file_ptr->get_file_path(),
// file_ptr->get_file_size(),
// file_ptr->get_file_mime());
// }
// else
// {
// DBG("未找到对应的文件信息, 文件ID: file123");
// }
// t.commit();
// }
{
odb::mysql::transaction t(handler->begin(),true);//后面的bool参数表示是否是在当前线程中创建事务对象
//获取数据库操作句柄
auto& db = t.database();
limevp_data::FileData file_data(db);
//删除文件信息
file_data.delFromDb("file123");
DBG("删除文件信息成功, 文件ID: file123");
t.commit();
}
}catch(const std::exception& e)
{
DBG("文件元信息数据库操作时发生异常: {}", e.what());
}
return 0;
}
编译构建:
# 1. 声明cmake所需版本
cmake_minimum_required(VERSION 3.1.3)
# 2. 设置工程项目名称(内部会生成一系列的内置变量)
project(file_test VERSION 1.0)
set(CMAKE_CXX_FLAGS "${CMAKE_CXX_FLAGS} -g -std=c++17")
# 4. 执行外部命令,生成临时代码文件: cal.pb.h cal.pb.cc
set(data_dir ${CMAKE_CURRENT_SOURCE_DIR}/../../data)
set(proto_dir ${CMAKE_CURRENT_SOURCE_DIR}/../../proto)
set(odb_files ${data_dir}/data.h)
file(GLOB_RECURSE proto_files ${proto_dir}/*.proto)
execute_process(
COMMAND odb -d mysql --std c++11 --generate-query --generate-schema --profile boost/date-time ${odb_files}
COMMAND protoc --experimental_allow_proto3_optional -I ${proto_dir} --cpp_out=${CMAKE_CURRENT_BINARY_DIR} ${proto_files}
)
# 3. 添加生成目标(可执行程序)
file(GLOB_RECURSE data_files ${data_dir}/*.cc)
file(GLOB_RECURSE temp_files ${CMAKE_CURRENT_BINARY_DIR}/*.cc ${CMAKE_CURRENT_BINARY_DIR}/*.cxx)
set(main_file ${CMAKE_CURRENT_SOURCE_DIR}/file_test.cc)
add_executable(${PROJECT_NAME} ${temp_files} ${data_files} ${main_file})
# 5. 头文件路径设置
include_directories(${CMAKE_CURRENT_BINARY_DIR})
include_directories(${data_dir})
# 6. 依赖检查(依赖库的查找以及链接选项的添加)
find_package(lime_scaffold REQUIRED)
# 7. 设置链接库
target_link_libraries(${PROJECT_NAME} PRIVATE lime_scaffold::lime_scaffold)
2.2会话元信息表定义及数据操作类实现
2.2.1数据库表设计与定义
| 字段 | 类型 | 约束 | 空值 | 备注 |
|---|---|---|---|---|
| ID | BIGINT UNSIGNED | FK | NOT NULL | |
| 会话ID | VARCHAR(64) | UK | NOT NULL | |
| 用户ID | VARCHAR(64) | NK | NULL | 为空代表未登录 |
#pragma once
#include <string>
#include <cstddef> // std::size_t
#include <boost/date_time/posix_time/posix_time.hpp>
#include <odb/nullable.hxx>
#include <odb/core.hxx>
#include <memory>
//定义所有的数据库表结构
namespace limevp_data{
/*会话元信息表*/
#pragma db object table("tbl_session_metadata")
class Session{
public:
using Ptr = std::shared_ptr<Session>;
Session();
Session(const std::string& session_id);
Session(const std::string& session_id, const odb::nullable<std::string>& user_id);
unsigned long get_id();
std::string get_session_id() const;
void set_session_id(const std::string& session_id);
odb::nullable<std::string> get_user_id() const;
void set_user_id(const odb::nullable<std::string>& user_id);
private:
friend class odb::access;
#pragma db id auto column("id") type("BIGINT UNSIGNED")
unsigned long _id;
#pragma db column("session_id") type("VARCHAR(64)") not_null unique
std::string _session_id;
#pragma db column("user_id") type("VARCHAR(64)") null index
odb::nullable<std::string> _user_id;
};
#pragma db view object(Session) \
query((?))
struct SessionPtr {
using ptr = std::shared_ptr<SessionPtr>;
std::shared_ptr<Session> session;
};
}//namespace limevp_data
#include "data.h"
namespace limevp_data{
/*session metadata start*/
Session::Session() {}
Session::Session(const std::string& session_id)
: _session_id(session_id)
{}
Session::Session(const std::string& session_id, const odb::nullable<std::string>& user_id)
: _session_id(session_id)
, _user_id(user_id)
{}
unsigned long Session::get_id() {return _id;}
std::string Session::get_session_id() const {return _session_id;}
void Session::set_session_id(const std::string& session_id) {_session_id = session_id;}
odb::nullable<std::string> Session::get_user_id() const {return _user_id;}
void Session::set_user_id(const odb::nullable<std::string>& user_id) {_user_id = user_id;}
/*session metadata end*/
}//namespace limevp_data
2.2.2数据表操作类实现
这里因为会话元信息在获取之后是需要加载到缓存中的,但是我们先来思考一个问题,上面注意到我们在数据库表的定义里面加了一个视图查询操作,因为uid是NK的,当不同的设备都没有进行登录时第一次登录均为临时用户,如果他们使用同一个账号进行登录,便会将此临时用户的sessionId关联到uid上。不过此时便会出现多个sid对应一个uid的情况,所以我们在根据uid查找缓存时查到的数据不止有一个。因此需要视图操作批量查找。
再来说此类大致需要哪些操作吧:
- 向数据库中新增会话元信息,新增数据不需要加载到缓存中
- 依据sid从数据库中获取会话元信息(优先从缓存获取,缓存未命中则从数据库获取,并添加缓存)
- 更新数据库会话,并删除会话缓存,发布删除消息
- 通过会话ID删除数据库会话,并删除会话缓存,发布删除消息
- 通过用户ID删除数据库会话,并删除会话缓存,发布删除消息
因为缓存删除操作上面我们已经实现过了,所以这里需要的对象有mysql,redis数据库操作句柄及缓存删除操作句柄,具体实现如下:
namespace limevp_data{
class SessionData{
public:
using Ptr = std::shared_ptr<SessionData>;
SessionData(odb::database& db,
const std::shared_ptr<sw::redis::Redis>& redis,
const RemoveCache::Ptr& rmcahe);
//向数据库新增会话
void addSession2Db(Session& session);
//更新数据库会话,并删除会话缓存,发布删除消息-必须先调用getSession获取到对象
void updateSession2Db(Session& session);
//通过会话ID删除数据库会话,并删除会话缓存,发布删除消息
void delSessionFromDbBySessionId(const std::string& sessionId);
//通过用户ID删除数据库会话,并删除会话缓存,发布删除消息
void delSessionFromDbByUserId(const std::string& userId);
//获取会话信息(优先从缓存获取,缓存未命中则从数据库获取,并添加缓存)
Session::Ptr getSession(const std::string& sessionId);
private:
std::string getCacheKey(const std::string& sessionId);
//向数据库添加会话信息
void addSessionToDb(Session& session);
//从数据库获取会话信息-通过会话ID
Session::Ptr getSessionBySidFromDb(const std::string& sessionId);
//从数据库获取会话信息-通过用户ID
std::vector<Session::Ptr> getSessionsByUidFromDb(const std::string& userId);
//修改数据库会话信息
void updateSessionToDb(Session& session);
//通过会话ID删除数据库会话
void delSessionBySidFromDb(const std::string& sessionId);
//通过用户ID删除数据库会话
void delSessionByUidFromDb(const std::string& userId);
//向redis添加会话缓存
void addSessionToCache(const Session::Ptr& session);
//从redis获取会话缓存
Session::Ptr getSessionFromCache(const std::string& sessionId);
//通过会话ID删除缓存会话
void delSessionFromCache(const std::string& sessionId);
private:
static const std::string _cache_key_prefix;//缓存KEY前缀
static const int _cache_expire_seconds;//缓存过期时间
static const std::string _session_id;
static const std::string _user_id;
odb::database& _db;//mysql操作句柄
std::shared_ptr<sw::redis::Redis> _redis;//redis操作句柄
RemoveCache::Ptr _rmcahe;//缓存同步句柄
};
}//limevp_data
#include "session.h"
namespace limevp_data{
const std::string SessionData::_cache_key_prefix = "vp_session_";
const int SessionData::_cache_expire_seconds = 3600; //1小时
const std::string SessionData::_session_id = "sessionid";
const std::string SessionData::_user_id = "userid";//与设计的RESTful API一致
SessionData::SessionData(odb::database& db,const std::shared_ptr<sw::redis::Redis>& redis,const RemoveCache::Ptr& rmcahe)
: _db(db),_redis(redis),_rmcahe(rmcahe)
{}
//向数据库新增会话
void SessionData::addSession2Db(Session& session)
{
//此时因为是数据库新增,所以不需要添加到缓存
addSessionToDb(session);
}
//更新数据库会话,并删除会话缓存,发布删除消息
void SessionData::updateSession2Db(Session& session)
{
updateSessionToDb(session);
//删除缓存
_rmcahe->sync(getCacheKey(session.get_session_id()));
}
//通过会话ID删除数据库会话,并删除会话缓存,发布删除消息
void SessionData::delSessionFromDbBySessionId(const std::string& sessionId)
{
delSessionBySidFromDb(sessionId);
//删除缓存
_rmcahe->sync(getCacheKey(sessionId));
}
//通过用户ID删除数据库会话,并删除会话缓存,发布删除消息
void SessionData::delSessionFromDbByUserId(const std::string& userId)
{
//先通过用户ID获取所有会话ID
auto sessions = getSessionsByUidFromDb(userId);
if(sessions.empty())
{
return;
}
//删除数据库会话
delSessionByUidFromDb(userId);
//删除缓存
for(auto& session : sessions)
{
_rmcahe->sync(getCacheKey(session->get_session_id()));
DBG("通过用户ID删除会话,删除缓存key: {}", getCacheKey(session->get_session_id()));
}
}
//获取会话信息(优先从缓存获取,缓存未命中则从数据库获取,并添加缓存)
Session::Ptr SessionData::getSession(const std::string& sessionId)
{
auto session = getSessionFromCache(sessionId);
if(session)
{
return session;
}
session = getSessionBySidFromDb(sessionId);
if(session)
{
addSessionToCache(session);
return session;
}
return nullptr;
}
std::string SessionData::getCacheKey(const std::string& sessionId)
{
return _cache_key_prefix + sessionId;
}
//向数据库添加会话信息
void SessionData::addSessionToDb(Session& session)
{
_db.persist(session);
}
//从数据库获取会话信息-通过会话ID
Session::Ptr SessionData::getSessionBySidFromDb(const std::string& sessionId)
{
Session::Ptr result(_db.query_one<Session>(odb::query<Session>::session_id == sessionId));
return result;
}
//从数据库获取会话信息-通过用户ID
std::vector<Session::Ptr> SessionData::getSessionsByUidFromDb(const std::string& userId)
{
//因为一个用户可能对应多个会话,比如不同客户端登录同一个账号时,就会有多个sessionid对应一个userid
typedef odb::query<SessionPtr> Query;
typedef odb::result<SessionPtr> Result;
Result r = _db.query<SessionPtr>(Query::user_id == userId);
std::vector<Session::Ptr> result;
for (auto& item : r)
{
result.push_back(item.session);
}
return result;
}
//修改数据库会话信息-内部直接进行查找,方便外部调用
void SessionData::updateSessionToDb(Session& session)
{
auto old_session = getSessionBySidFromDb(session.get_session_id());
//如果为空则插入
if (!old_session)
{
addSessionToDb(session);
return;
}
//如果不为空则进行更新
old_session->set_user_id(session.get_user_id());
_db.update(old_session.get());
}
//通过会话ID删除数据库会话
void SessionData::delSessionBySidFromDb(const std::string& sessionId)
{
_db.erase_query<Session>(odb::query<Session>::session_id == sessionId);
}
//通过用户ID删除数据库会话
void SessionData::delSessionByUidFromDb(const std::string& userId)
{
_db.erase_query<Session>(odb::query<Session>::user_id == userId);
}
//向redis添加会话缓存
void SessionData::addSessionToCache(const Session::Ptr& session)
{
auto rtx = _redis->transaction(false,false);
auto r = rtx.redis();
std::string key = getCacheKey(session->get_session_id());
std::unordered_map<std::string, std::string> values;
values[_session_id] = session->get_session_id();
values[_user_id] = session->get_user_id().get();
r.hmset(key, values.begin(), values.end());
r.expire(key, std::chrono::seconds(_cache_expire_seconds));
}
//从redis获取会话缓存
Session::Ptr SessionData::getSessionFromCache(const std::string& sessionId)
{
auto rtx = _redis->transaction(false,false);
auto r = rtx.redis();
std::string key = getCacheKey(sessionId);
auto value = r.hget(key,sessionId);
if(!value)
{
return nullptr;
}
Session::Ptr session = std::make_shared<Session>();
session->set_session_id(sessionId);
session->set_user_id(*value);
return session;
}
//通过会话ID删除缓存会话
void SessionData::delSessionFromCache(const std::string& sessionId)
{
auto rtx = _redis->transaction(false,false);
auto r = rtx.redis();
std::string key = getCacheKey(sessionId);
r.del(key);
}
} // namespace limevp_data
2.2.3会话表操作测试
#include "session.h"
#include "remove_cache.h"
#include <lime_scaffold/limelog.h>
const std::string USER = "root";
const std::string PASSWORD = "123456";
const std::string DATABASE = "vbtest";
const std::string HOST = "192.168.30.128";
const unsigned int PORT = 3306;
const std::string CHRSET = "utf8";
int main()
{
limelog::limelog_init();
limeodb::mysql_settings msettings{
.host = HOST,
.user = USER,
.password = PASSWORD,
.database = DATABASE,
.port = PORT,
.charset = CHRSET
};
std::shared_ptr<odb::mysql::database> handler = limeodb::DbFactory::create_mysqldb(msettings);
limeredis::redis_settings rsettings{
.host = "192.168.30.128",
.password = "123456",
.connection_pool_size = 3
};
auto redis = limeredis::RedisFactory::create(rsettings);
//进行相关初始设定
std::string url = "amqp://admin:123456@192.168.30.128:5672/";
limemq::declare_settings mqsettings
{
.exchange = "delete_cache_exchange",
.exchange_type = "delayed",
.queue = "delete_cache_queue",
.binding_key = "delete_cache_queue",
.delayed_ttl = 3000
};
//创建客户端
limemq::MQClient::ptr client = std::make_shared<limemq::MQClient>(url);
limevp_data::RemoveCache::Ptr rmcache = std::make_shared<limevp_data::RemoveCache>(redis, client, mqsettings);
try
{
// {
// //向数据库新增数据
// odb::mysql::transaction mtx(handler->begin(),true);
// auto& db = mtx.database();
// limevp_data::SessionData sessionData(db, redis, rmcache);
// limevp_data::Session session("session_67890", std::string("user_12345"));
// sessionData.addSession2Db(session);
// mtx.commit();
// DBG("向数据库新增会话数据成功");
// }
// {
// //从数据库获取数据
// odb::mysql::transaction mtx(handler->begin(),true);
// auto& db = mtx.database();
// limevp_data::SessionData sessionData(db, redis, rmcache);
// auto session = sessionData.getSession("session_67890");
// DBG("从数据库/缓存获取会话数据成功, session_id: {}, user_id: {}", session->get_session_id(), session->get_user_id().get());
// mtx.commit();
// }
// {
// //更新数据库数据
// odb::mysql::transaction mtx(handler->begin(),true);
// auto& db = mtx.database();
// limevp_data::SessionData sessionData(db, redis, rmcache);
// limevp_data::Session session("session_12345",std::string("user_12345"));
// sessionData.updateSession2Db(session);
// mtx.commit();
// DBG("更新数据库会话数据成功");
// }
// {
// //通过用户ID删除数据库数据
// odb::mysql::transaction mtx(handler->begin(),true);
// auto& db = mtx.database();
// limevp_data::SessionData sessionData(db, redis, rmcache);
// sessionData.delSessionFromDbByUserId("user_12345");
// mtx.commit();
// DBG("删除数据库会话数据成功");
// }
{
//通过会话ID删除数据库数据
odb::mysql::transaction mtx(handler->begin(),true);
auto& db = mtx.database();
limevp_data::SessionData sessionData(db, redis, rmcache);
sessionData.delSessionFromDbBySessionId("session_67890");
mtx.commit();
DBG("删除数据库会话数据成功");
}
}catch(const std::exception& e)
{
DBG("会话元信息数据库操作时发生异常: {}", e.what());
}
DBG("会话元信息数据库操作成功,按回车健退出");
getchar();
return 0;
}
编译构建:
# 1. 声明cmake所需版本
cmake_minimum_required(VERSION 3.1.3)
# 2. 设置工程项目名称(内部会生成一系列的内置变量)
project(test VERSION 1.0)
set(CMAKE_CXX_FLAGS "${CMAKE_CXX_FLAGS} -g -std=c++17")
# 4. 执行外部命令,生成临时代码文件: cal.pb.h cal.pb.cc
set(data_dir ${CMAKE_CURRENT_SOURCE_DIR}/../../data)
set(proto_dir ${CMAKE_CURRENT_SOURCE_DIR}/../../proto)
set(odb_files ${data_dir}/data.h)
file(GLOB_RECURSE proto_files ${proto_dir}/*.proto)
execute_process(
COMMAND odb -d mysql --std c++11 --generate-query --generate-schema --profile boost/date-time ${odb_files}
COMMAND protoc --experimental_allow_proto3_optional -I ${proto_dir} --cpp_out=${CMAKE_CURRENT_BINARY_DIR} ${proto_files}
)
# 3. 添加生成目标(可执行程序)
file(GLOB_RECURSE data_files ${data_dir}/*.cc)
file(GLOB_RECURSE temp_files ${CMAKE_CURRENT_BINARY_DIR}/*.cc ${CMAKE_CURRENT_BINARY_DIR}/*.cxx)
set(main_file ${CMAKE_CURRENT_SOURCE_DIR}/test.cc)
add_executable(${PROJECT_NAME} ${temp_files} ${data_files} ${main_file})
# 5. 头文件路径设置
include_directories(${CMAKE_CURRENT_BINARY_DIR})
include_directories(${data_dir})
# 6. 依赖检查(依赖库的查找以及链接选项的添加)
find_package(lime_scaffold REQUIRED)
# 7. 设置链接库
target_link_libraries(${PROJECT_NAME} PRIVATE lime_scaffold::lime_scaffold)
三.业务逻辑实现
整个文件子服务的业务流程如下:

所以我们需要一个数据库操作类(svc_data),接收消息队列的文件删除消息并进行对应文件,数据库表操作的类(svc_mq),处理rpc请求的rpc服务器(svc_rpc),整合前面三者为一个子服务的类(svc_server),图示关系如下:

这里使用了一个建造者模式,接下来我们来进行具体实现:
3.0错误类封装
如果我们所有的错误都使用魔幻数字进行报错返回,显然是不合理的,所以我们需要针对文件子服务的错误情况封装一个error类:
#include <iostream>
namespace limevp_error{
enum ErrorCode{
SUCCESS = 0,
INVALID_SESSION,
UPLOAD_FILE_TO_FDFS_FAILED,
DATABASE_OP_FAILED,
FILE_NOT_EXIST,
DOWNLOAD_FILE_FROM_FDFS_FAILED
};
const char* get_error_message(ErrorCode error_code);
class ErrorException : public std::exception{
public:
ErrorException(ErrorCode code);
virtual const char* what() const noexcept override;
int getCode() const noexcept;
private:
ErrorCode _code;
};
}//namespace limevp_error
#include "error.h"
namespace limevp_error{
const char* get_error_message(ErrorCode error_code)
{
switch (error_code)
{
case SUCCESS: return "成功";
case INVALID_SESSION: return "无效会话";
case UPLOAD_FILE_TO_FDFS_FAILED: return "上传文件到FDFS失败";
case DATABASE_OP_FAILED: return "数据库操作失败";
case FILE_NOT_EXIST: return "文件不存在";
case DOWNLOAD_FILE_FROM_FDFS_FAILED: return "下载文件从FDFS失败";
}
return "未知错误";
}
ErrorException::ErrorException(ErrorCode code)
{
_code = code;
}
const char* ErrorException::what() const noexcept
{
return get_error_message(_code);
}
int ErrorException::getCode() const noexcept
{
return static_cast<int>(_code);
}
} // namespace limevp_error
3.1svc_data
#pragma once
#include "file.h"
#include "session.h"
namespace svc_file{
//对文件元信息数据库进行相关操作
class SvcData{
public:
using Ptr = std::shared_ptr<SvcData>;
SvcData(const std::shared_ptr<odb::database>& db,
const std::shared_ptr<sw::redis::Redis>& redis,
const limemq::MQClient::ptr& mqClient,
const limemq::declare_settings& mqsettings);
//依据sessionId获取对应上传文件用户ID
odb::nullable<std::string> getUplUidBySid(const std::string& session_id);
//依据文件ID获取文件元信息
limevp_data::File::Ptr getFileByFid(const std::string& file_id);
//新增文件元信息
bool addFileMeta(limevp_data::File& file);
//删除文件元信息
bool delFileMeta(const std::string& file_id);
private:
//mysql数据库对象,redis缓存对象,缓存同步对象
std::shared_ptr<odb::database> _db;
std::shared_ptr<sw::redis::Redis> _redis;
limevp_data::RemoveCache::Ptr _rmcache;
};
}
#include "svc_data.h"
#include <lime_scaffold/limelog.h>
namespace svc_file{
SvcData::SvcData(const std::shared_ptr<odb::database>& db,
const std::shared_ptr<sw::redis::Redis>& redis,
const limemq::MQClient::ptr& mqClient,
const limemq::declare_settings& mqsettings)
: _db(db), _redis(redis),
_rmcache(std::make_shared<limevp_data::RemoveCache>(redis, mqClient, mqsettings))
{}
//依据sessionId获取对应上传文件用户ID
odb::nullable<std::string> SvcData::getUplUidBySid(const std::string& session_id)
{
try{
odb::transaction mtx(_db->begin(),true);
auto& db = mtx.database();
limevp_data::SessionData sessionData(db, _redis, _rmcache);
auto session = sessionData.getSession(session_id);
mtx.commit();
if(session){
return session->get_user_id();
}
}catch(const std::exception& e){
ERR("依据sessionId获取对应上传文件用户ID失败:{}-{}", session_id, e.what());
}
return odb::nullable<std::string>();
}
//依据文件ID获取文件元信息
limevp_data::File::Ptr SvcData::getFileByFid(const std::string& file_id)
{
try{
odb::transaction mtx(_db->begin());
auto& db = mtx.database();
limevp_data::FileData fileData(db);
auto file = fileData.getFileFromDb(file_id);
mtx.commit();
if(file){
return file;
}
}catch(const std::exception& e){
ERR("依据文件ID获取文件元信息失败:{}-{}", file_id, e.what());
}
return nullptr;
}
//新增文件元信息
bool SvcData::addFileMeta(limevp_data::File& file)
{
try{
odb::transaction mtx(_db->begin());
auto& db = mtx.database();
limevp_data::FileData fileData(db);
fileData.addFile2Db(file);
mtx.commit();
}catch(const std::exception& e){
ERR("新增文件元信息失败:{}-{}", file.get_file_id(), e.what());
return false;
}
return true;
}
//删除文件元信息
bool SvcData::delFileMeta(const std::string& file_id)
{
try{
odb::transaction mtx(_db->begin());
auto& db = mtx.database();
limevp_data::FileData fileData(db);
fileData.delFromDb(file_id);
mtx.commit();
}catch(const std::exception& e){
ERR("删除文件元信息失败:{}-{}", file_id, e.what());
return false;
}
return true;
}
}//namespace svc_file
3.2svc_mq
#pragma once
#include "svc_data.h"
#include <lime_scaffold/limemq.h>
namespace svc_file{
//订阅消息队列,进行文件删除操作
class SvcDelFileMQ{
public:
using Ptr = std::shared_ptr<SvcDelFileMQ>;
SvcDelFileMQ(const SvcData::Ptr& svc_data,
const limemq::MQClient::ptr& mqClient,
const limemq::declare_settings& mqsettings);
private:
void callback(const char *msg, size_t len);
private:
limemq::Subscriber::ptr _sub;
SvcData::Ptr _svc_data;
};
}//namespace svc_file
#include "svc_mq.h"
#include <lime_scaffold/limelog.h>
#include <lime_scaffold/limefds.h>
namespace svc_file{
SvcDelFileMQ::SvcDelFileMQ(const SvcData::Ptr& svc_data,
const limemq::MQClient::ptr& mqClient,
const limemq::declare_settings& mqsettings)
:_sub(std::make_shared<limemq::Subscriber>(mqClient, mqsettings)),
_svc_data(svc_data)
{
_sub->consume(std::bind(&SvcDelFileMQ::callback, this, std::placeholders::_1, std::placeholders::_2));
}
void SvcDelFileMQ::callback(const char *msg, size_t len)
{
CommunicationInterface::DeleteFileMsg delete_msg;
auto ret = delete_msg.ParseFromArray(msg, len);
if(!ret)
{
ERR("解析文件删除消息失败");
return;
}
//对文件进行删除操作
int sz = delete_msg.file_id_size();
for(int i = 0; i < sz; i++)
{
auto file_id = delete_msg.file_id(i);
auto file = _svc_data->getFileByFid(file_id);
if(file)
{
_svc_data->delFileMeta(file_id);
//通过fds删除文件
limefds::FdfsClient::deleteFile(file->get_file_path());
DBG("删除文件成功:{}-{}", file_id, file->get_file_path());
}
else
{
ERR("文件不存在:{}", file_id);
}
}
}
}//namespace svc_file
3.3svc_rpc
#pragma once
#include "base.pb.h"
#include "file.pb.h"
#include "svc_data.h"
namespace svc_file {
class SvcFileRpc : public CommunicationInterface::FileService{
public:
SvcFileRpc(const SvcData::Ptr& svc_data);
virtual void uploadPhoto(google::protobuf::RpcController* controller,
const ::CommunicationInterface::UploadPhotoReq* request,
::CommunicationInterface::UploadPhotoResp* response,
::google::protobuf::Closure* done);
virtual void downloadPhoto(google::protobuf::RpcController* controller,
const ::CommunicationInterface::DownloadPhotoReq* request,
::CommunicationInterface::DownloadPhotoResp* response,
::google::protobuf::Closure* done);
virtual void uploadVideo(google::protobuf::RpcController* controller,
const ::CommunicationInterface::UploadVideoReq* request,
::CommunicationInterface::UploadVideoResp* response,
::google::protobuf::Closure* done);
virtual void downloadVideo(google::protobuf::RpcController* controller,
const ::CommunicationInterface::DownloadVideoReq* request,
::CommunicationInterface::DownloadVideoResp* response,
::google::protobuf::Closure* done);
private:
SvcData::Ptr _svc_data;
static const size_t uuid_length;//生成文件ID的长度
};
} // namespace svc_file
#include <lime_scaffold/limeutil.h>
#include <lime_scaffold/limerpc.h>
#include "svc_rpc.h"
#include "../../common/error.h"
#include <lime_scaffold/limefds.h>
namespace svc_file {
const size_t SvcFileRpc::uuid_length = 16;
SvcFileRpc::SvcFileRpc(const SvcData::Ptr& svc_data)
: _svc_data(svc_data)
{}
void SvcFileRpc::uploadPhoto(google::protobuf::RpcController* controller,
const ::CommunicationInterface::UploadPhotoReq* request,
::CommunicationInterface::UploadPhotoResp* response,
::google::protobuf::Closure* done)
{
//只有done->Run()才会真正的结束rpc调用
brpc::ClosureGuard done_guard(done);
// 1. 获取请求元素: 请求ID, 会话ID(⽤于获取用户ID), ⽂件数据, ⽂件mime
std::string request_id = request->requestid();
std::string session_id = request->sessionid();
auto file_info = request->fileinfo();
std::string file_data = file_info.filedata();
std::string file_mime = file_info.filemime();
try{
// 2. 从缓存通过会话ID查询会话信息, 获得上传者用户ID(缓存未命中则从数据库获取)
auto uploader_uid = _svc_data->getUplUidBySid(session_id);
if(!uploader_uid)
{
throw limevp_error::ErrorException(limevp_error::INVALID_SESSION);
}
// 3. 将文件数据上传到FDFS进行存储,并获取FDFS文件ID--作为存储路径PATH
auto file_path = limefds::FdfsClient::uploadFromBuffer(file_data);
if(!file_path)
{
throw limevp_error::ErrorException(limevp_error::UPLOAD_FILE_TO_FDFS_FAILED);
}
// 4. ⽣成文件ID, 向数据库新增文件信息(文件ID,文件大小,文件mime,上传用户ID,存储路径)
limevp_data::File file_info;
file_info.set_file_id(limeutil::LimeRandom::code(uuid_length));
file_info.set_upload_user_id(*uploader_uid);
file_info.set_file_path(*file_path);
file_info.set_file_size(file_data.size());
file_info.set_file_mime(file_mime);
bool ret = _svc_data->addFileMeta(file_info);
if(!ret)
{
throw limevp_error::ErrorException(limevp_error::DATABASE_OP_FAILED);
}
// 5. 返回响应信息: 文件ID
response->set_requestid(request_id);
response->set_errorcode(limevp_error::SUCCESS);
response->set_errormsg("");
auto result = response->mutable_result();
result->set_fileid(file_info.get_file_id());
}catch(const limevp_error::ErrorException& e){
response->set_requestid(request_id);
response->set_errorcode(e.getCode());
response->set_errormsg(e.what());
}
}
void SvcFileRpc::downloadPhoto(google::protobuf::RpcController* controller,
const ::CommunicationInterface::DownloadPhotoReq* request,
::CommunicationInterface::DownloadPhotoResp* response,
::google::protobuf::Closure* done)
{
brpc::ClosureGuard done_guard(done);
// 1. 获取请求元素: 文件ID,请求ID
std::string file_id = request->fileid();
std::string request_id = request->requestid();
try{
// 2. 从数据库通过文件ID查询文件信息(文件ID,文件大小,文件mime,上传用户ID,存储路径)
auto file_meta = _svc_data->getFileByFid(file_id);
if(!file_meta)
{
throw limevp_error::ErrorException(limevp_error::FILE_NOT_EXIST);
}
size_t file_size = file_meta->get_file_size();
std::string file_mime = file_meta->get_file_mime();
std::string file_upl_uid = file_meta->get_upload_user_id();
std::string file_path = file_meta->get_file_path();
// 3. 通过文件信息中的路径,从FDFS进行文件数据下载到内存中
std::string file_data;
bool ret = limefds::FdfsClient::downloadToBuffer(file_path, file_data);
if(!ret)
{
throw limevp_error::ErrorException(limevp_error::DOWNLOAD_FILE_FROM_FDFS_FAILED);
}
// 4. 返回响应信息: 文件ID, 上传者ID, mime, 文件大小, 文件内容
response->set_requestid(request_id);
response->set_errorcode(limevp_error::SUCCESS);
response->set_errormsg("");
auto result = response->mutable_result();
result->set_fileid(file_id);
result->set_uploader(file_upl_uid);
result->set_filemime(file_mime);
result->set_filesize(file_size);
result->set_filedata(file_data);
}catch(const limevp_error::ErrorException& e){
response->set_requestid(request_id);
response->set_errorcode(e.getCode());
response->set_errormsg(e.what());
}
}
void SvcFileRpc::uploadVideo(google::protobuf::RpcController* controller,
const ::CommunicationInterface::UploadVideoReq* request,
::CommunicationInterface::UploadVideoResp* response,
::google::protobuf::Closure* done)
{
brpc::ClosureGuard done_guard(done);
// 1. 获取请求元素: 请求ID, 会话ID(⽤于获取用户ID), ⽂件数据, ⽂件mime
std::string request_id = request->requestid();
std::string session_id = request->sessionid();
auto file_info = request->fileinfo();
std::string file_data = file_info.filedata();
std::string file_mime = file_info.filemime();
try{
// 2. 从缓存通过会话ID查询会话信息, 获得上传者用户ID(缓存未命中则从数据库获取)
auto uploader_uid = _svc_data->getUplUidBySid(session_id);
if(!uploader_uid)
{
throw limevp_error::ErrorException(limevp_error::INVALID_SESSION);
}
// 3. 将文件数据上传到FDFS进行存储,并获取FDFS文件ID--作为存储路径PATH
auto file_path = limefds::FdfsClient::uploadFromBuffer(file_data);
if(!file_path)
{
throw limevp_error::ErrorException(limevp_error::UPLOAD_FILE_TO_FDFS_FAILED);
}
// 4. ⽣成文件ID, 向数据库新增文件信息(文件ID,文件大小,文件mime,上传用户ID,存储路径)
limevp_data::File file_info;
file_info.set_file_id(limeutil::LimeRandom::code(uuid_length));
file_info.set_upload_user_id(*uploader_uid);
file_info.set_file_path(*file_path);
file_info.set_file_size(file_data.size());
file_info.set_file_mime(file_mime);
bool ret = _svc_data->addFileMeta(file_info);
if(!ret)
{
throw limevp_error::ErrorException(limevp_error::DATABASE_OP_FAILED);
}
// 5. 返回响应信息: 文件ID
response->set_requestid(request_id);
response->set_errorcode(limevp_error::SUCCESS);
response->set_errormsg("");
auto result = response->mutable_result();
result->set_fileid(file_info.get_file_id());
}catch(const limevp_error::ErrorException& e){
response->set_requestid(request_id);
response->set_errorcode(e.getCode());
response->set_errormsg(e.what());
}
}
void SvcFileRpc::downloadVideo(google::protobuf::RpcController* controller,
const ::CommunicationInterface::DownloadVideoReq* request,
::CommunicationInterface::DownloadVideoResp* response,
::google::protobuf::Closure* done)
{
brpc::ClosureGuard done_guard(done);
// 1. 获取请求元素: 文件ID,请求ID
std::string file_id = request->fileid();
std::string request_id = request->requestid();
try{
// 2. 从数据库通过文件ID查询文件信息(文件ID,文件大小,文件mime,上传用户ID,存储路径)
auto file_meta = _svc_data->getFileByFid(file_id);
if(!file_meta)
{
throw limevp_error::ErrorException(limevp_error::FILE_NOT_EXIST);
}
size_t file_size = file_meta->get_file_size();
std::string file_mime = file_meta->get_file_mime();
std::string file_upl_uid = file_meta->get_upload_user_id();
std::string file_path = file_meta->get_file_path();
// 3. 通过文件信息中的路径,从FDFS进行文件数据下载到内存中
std::string file_data;
bool ret = limefds::FdfsClient::downloadToBuffer(file_path, file_data);
if(!ret)
{
throw limevp_error::ErrorException(limevp_error::DOWNLOAD_FILE_FROM_FDFS_FAILED);
}
// 4. 返回响应信息: 文件ID, 上传者ID, mime, 文件大小, 文件内容
response->set_requestid(request_id);
response->set_errorcode(limevp_error::SUCCESS);
response->set_errormsg("");
auto result = response->mutable_result();
result->set_fileid(file_id);
result->set_uploader(file_upl_uid);
result->set_filemime(file_mime);
result->set_filesize(file_size);
result->set_filedata(file_data);
}catch(const limevp_error::ErrorException& e){
response->set_requestid(request_id);
response->set_errorcode(e.getCode());
response->set_errormsg(e.what());
}
}
} // namespace svc_file
3.4svc_server
#pragma once
#include <lime_scaffold/limerpc.h>
#include "svc_mq.h"
#include "svc_rpc.h"
#include <lime_scaffold/limeetcd.h>
#include <lime_scaffold/limefds.h>
namespace svc_file {
class FileServer{
public:
using Ptr = std::shared_ptr<FileServer>;
FileServer(const std::shared_ptr<brpc::Server>& rpc_server,
const SvcDelFileMQ::Ptr& mq_del_server,
const limeetcd::SvcProvider::ptr& etcd_provider);
void start();
private:
//rpc服务器,用于接收外部的rpc请求,需要一直保留
std::shared_ptr<brpc::Server> _rpc_server;
//mq服务器,用于接收外部传入的删除文件请求,需要一直保留
SvcDelFileMQ::Ptr _mq_del_server;
//ectd注册中心对象,用于向注册中心注册服务,并进行服务保活,需要一直保留
limeetcd::SvcProvider::ptr _etcd_provider;
};
struct registry_settings {
std::string registry_center_addr; //注册中心地址
std::string service_name; //服务名称
std::string service_addr; //服务地址
};
//建造者模式构建FileServer对象并调用
class FileServerBuilder{
public:
FileServerBuilder() = default;
~FileServerBuilder() = default;
FileServerBuilder& set_listen_port(unsigned int port);
FileServerBuilder& set_mq_url(const std::string& url);
FileServerBuilder& set_registry_settings(const registry_settings& settings);
FileServerBuilder& set_mysql_settings(const limeodb::mysql_settings& settings);
FileServerBuilder& set_redis_settings(const limeredis::redis_settings& settings);
FileServerBuilder& set_mq_del_settings(const std::string& settings);
FileServerBuilder& set_mq_cache_settings(const std::string& settings);
FileServerBuilder& set_fdfs_settings(const std::vector<std::string>& tracker_servers);
FileServer::Ptr build();
private:
unsigned int _listen_port;//监听端口
std::string _mq_url;
registry_settings _reg_settings;//注册中心设置
limeodb::mysql_settings _mysql_settings;//mysql设置
limeredis::redis_settings _redis_settings;//redis设置
limemq::declare_settings _mq_del_settings;//订阅删除文件消息的消息队列设置
limemq::declare_settings _mq_cache_settings;//双写缓存消息队列设置
limefds::fdfs_config _fdfs_settings;//fast-fds全局设置
};
}//namespace svc_file
#include <lime_scaffold/limeutil.h>
#include "svc_server.h"
namespace svc_file{
FileServer::FileServer(const std::shared_ptr<brpc::Server>& rpc_server,
const SvcDelFileMQ::Ptr& mq_del_server,
const limeetcd::SvcProvider::ptr& etcd_provider)
: _rpc_server(rpc_server),
_mq_del_server(mq_del_server),
_etcd_provider(etcd_provider)
{}
void FileServer::start()
{
//启动FileServer
_etcd_provider->registry();
_rpc_server->RunUntilAskedToQuit();
}
FileServerBuilder& FileServerBuilder::set_listen_port(unsigned int port)
{
_listen_port = port;
return *this;
}
FileServerBuilder& FileServerBuilder::set_mq_url(const std::string& url)
{
_mq_url = url;
return *this;
}
FileServerBuilder& FileServerBuilder::set_registry_settings(const registry_settings& settings)
{
_reg_settings = settings;
return *this;
}
FileServerBuilder& FileServerBuilder::set_mysql_settings(const limeodb::mysql_settings& settings)
{
_mysql_settings = settings;
return *this;
}
FileServerBuilder& FileServerBuilder::set_redis_settings(const limeredis::redis_settings& settings)
{
_redis_settings = settings;
return *this;
}
FileServerBuilder& FileServerBuilder::set_mq_del_settings(const std::string& settings)
{
auto mq_del_settings = limeutil::LimeJson::deserialize(settings);
if (!mq_del_settings)
{
ERR("mq设置解析失败!");
abort();//
}
const Json::Value& mq_settings = *mq_del_settings;
_mq_del_settings.exchange = mq_settings["exchange"].asString();
_mq_del_settings.exchange_type = mq_settings["exchange_type"].asString();
_mq_del_settings.queue = mq_settings["queue"].asString();
_mq_del_settings.binding_key = mq_settings["binding_key"].asString();
_mq_del_settings.delayed_ttl = mq_settings["delayed_ttl"].asUInt64();
return *this;
}
FileServerBuilder& FileServerBuilder::set_mq_cache_settings(const std::string& settings)
{
auto mq_cache_settings = limeutil::LimeJson::deserialize(settings);
if (!mq_cache_settings)
{
ERR("mq设置解析失败!");
abort();//
}
const Json::Value& mq_settings = *mq_cache_settings;
_mq_cache_settings.exchange = mq_settings["exchange"].asString();
_mq_cache_settings.exchange_type = mq_settings["exchange_type"].asString();
_mq_cache_settings.queue = mq_settings["queue"].asString();
_mq_cache_settings.binding_key = mq_settings["binding_key"].asString();
_mq_cache_settings.delayed_ttl = mq_settings["delayed_ttl"].asUInt64();
return *this;
}
FileServerBuilder& FileServerBuilder::set_fdfs_settings(const std::vector<std::string>& tracker_servers)
{
_fdfs_settings.tracker_servers = tracker_servers;
return *this;
}
FileServer::Ptr FileServerBuilder::build()
{
//初始化mqclient
auto mq_client = std::make_shared<limemq::MQClient>(_mq_url);
//初始化etcdclient
auto provider = std::make_shared<limeetcd::SvcProvider>(_reg_settings.registry_center_addr, _reg_settings.service_name, _reg_settings.service_addr);
//初始化mysql客户端
auto mysql_client = limeodb::DbFactory::create_mysqldb(_mysql_settings);
//初始化redis客户端
auto redis_client = limeredis::RedisFactory::create(_redis_settings);
//全局初始化fds
limefds::FdfsClient::init(_fdfs_settings);
//初始化数据库操作句柄svc_data
auto svc_data = std::make_shared<SvcData>(mysql_client, redis_client,mq_client,_mq_cache_settings);
//初始化删除文件mq
auto svcdel_filemq = std::make_shared<SvcDelFileMQ>(svc_data,mq_client, _mq_del_settings);
//初始化rpc服务
SvcFileRpc* svcfilerpc = new SvcFileRpc(svc_data);
auto rpc_server = limerpc::RpcServer::create(_listen_port,svcfilerpc);
//构建FileServer对象
auto file_server = std::make_shared<FileServer>(rpc_server,svcdel_filemq,provider);
return file_server;
}
}
3.5主控流程
#include "svc_server.h"
#include <gflags/gflags.h>
#include <lime_scaffold/limelog.h>
//监听端口
DEFINE_int32(listen_port, 9001, "监听端口");
//日志相关设置
DEFINE_bool(log_async, false, "是否开启异步日志");
DEFINE_int32(log_level, 1, "日志输出等级: 1-debug;2-info;3-warn;4-error; 6-off");
DEFINE_string(log_format, "[%H:%M:%S][%-7l]%v", "日志输出格式");
DEFINE_string(log_target, "stdout", "日志输出目标");
//mysql数据库连接参数
DEFINE_string(mysql_host, "192.168.30.128", "mysql数据库地址");
DEFINE_string(mysql_user, "root", "mysql数据库用户名");
DEFINE_string(mysql_password, "123456", "mysql数据库密码");
DEFINE_string(mysql_database, "vbtest", "mysql数据库名称");
DEFINE_int32(mysql_port, 3306, "mysql数据库端口");
DEFINE_string(mysql_charset, "utf8", "mysql数据库字符集");
DEFINE_int32(mysql_connect_pool_size, 5, "mysql数据库连接池大小");
//redis数据库连接参数
DEFINE_string(redis_host, "192.168.30.128", "redis数据库地址");
DEFINE_string(redis_password,"123456", "redis数据库密码");
DEFINE_int32(redis_connection_pool_size, 5, "redis数据库连接池大小");
//注册中心配置
DEFINE_string(reg_center_addr, "http://192.168.30.128:2379", "注册中心地址");
DEFINE_string(service_name, "file_server", "服务名称");
DEFINE_string(service_addr, "192.168.30.128:9001", "服务地址");
//mq相关配置
DEFINE_string(mq_host, "amqp://admin:123456@192.168.30.128:5672/", "mq服务器地址");
//mq订阅删除文件相关配置-json格式
DEFINE_string(mq_del_settings,R"({"exchange":"delete_file_exchange","exchange_type":"direct","queue":"delete_file_queue","binding_key":"delete_file_queue"})", "mq订阅删除文件相关配置-json格式");
//mq订阅缓存删除消息的相关配置-json格式
DEFINE_string(mq_cache_del_settings,R"({"exchange":"delete_cache_exchange","exchange_type":"delayed","queue":"delete_cache_queue","binding_key":"delete_cache_queue","delayed_ttl":3000})", "mq订阅缓存删除消息的相关配置-json格式");
//fds相关配置-数组
DEFINE_string(fds_tracker_server,"192.168.30.128:22122", "fds_tracker_server服务器地址");
//文件子服务的主控流程
int main(int argc, char* argv[])
{
//进⾏参数解析
google::ParseCommandLineFlags(&argc, &argv, true);
//初始化日志配置
limelog::log_settings settings{
.async = FLAGS_log_async,
.level = FLAGS_log_level,
.format = FLAGS_log_format,
.path = FLAGS_log_target
};
limelog::limelog_init(settings);
//组织FileServer的配置参数
//1.初始化mysql数据库连接参数
limeodb::mysql_settings mysql_settings{
.host = FLAGS_mysql_host,
.user = FLAGS_mysql_user,
.password = FLAGS_mysql_password,
.database = FLAGS_mysql_database,
.port = FLAGS_mysql_port,
.charset = FLAGS_mysql_charset,
.connection_pool_size = FLAGS_mysql_connect_pool_size
};
//2.初始化redis数据库连接参数
limeredis::redis_settings redis_settings{
.host = FLAGS_redis_host,
.password = FLAGS_redis_password,
.connection_pool_size = FLAGS_redis_connection_pool_size
};
//3.初始化fds相关配置
std::vector<std::string> fds_tracker_servers;
fds_tracker_servers.push_back(FLAGS_fds_tracker_server);
//4.初始化注册中心配置
svc_file::registry_settings reg_settings{
.registry_center_addr = FLAGS_reg_center_addr,
.service_name = FLAGS_service_name,
.service_addr = FLAGS_service_addr
};
//使用初始化数据构造FileServerBuilder
svc_file::FileServerBuilder builder;
auto file_server = builder.set_listen_port(FLAGS_listen_port)\
.set_mq_url(FLAGS_mq_host)\
.set_registry_settings(reg_settings)\
.set_mysql_settings(mysql_settings)\
.set_redis_settings(redis_settings)\
.set_mq_del_settings(FLAGS_mq_del_settings)\
.set_mq_cache_settings(FLAGS_mq_cache_del_settings)\
.set_fdfs_settings(fds_tracker_servers)\
.build();
//启动FileServer
file_server->start();
return 0;
}
3.6测试客户端样例编写
#include <iostream>
#include <lime_scaffold/limelog.h>
#include <lime_scaffold/limeetcd.h>
#include <lime_scaffold/limemq.h>
#include <lime_scaffold/limerpc.h>
#include <lime_scaffold/limeutil.h>
#include <gflags/gflags.h>
#include <data-odb.hxx>
#include <message.pb.h>
#include <base.pb.h>
#include <file.pb.h>
#include <message.pb.h>
//定义mq相关配置
DEFINE_string(mq_host, "amqp://admin:123456@192.168.30.128:5672/", "mq服务器地址");
DEFINE_string(mq_del_settings,R"({"exchange":"delete_file_exchange","exchange_type":"direct","queue":"delete_file_queue","binding_key":"delete_file_queue"})", "mq订阅删除文件相关配置-json格式");
//定义etcd相关配置
DEFINE_string(reg_center_addr, "http://192.168.30.128:2379", "注册中心地址");
DEFINE_string(svc_name, "file_server", "目标服务名称");
void UploadPhotoTest(CommunicationInterface::FileService_Stub& stub)
{
//上传图片测试
CommunicationInterface::UploadPhotoReq* req = new CommunicationInterface::UploadPhotoReq();
req->set_requestid("上传图片测试");
req->set_sessionid("114514");
auto file_info = req->mutable_fileinfo();
file_info->set_filemime("image/jpeg");
std::string file_data;
if(!limeutil::LimeFile::read("/home/dev/workspace/cpp-videoplayer/server/svc_file/client/111.jpg",file_data))
{
ERR("文件读取失败,请检查文件路径是否正确!");
abort();
}
file_info->set_filedata(file_data);
//发起rpc调用-异步调用
CommunicationInterface::UploadPhotoResp* resp = new CommunicationInterface::UploadPhotoResp();
brpc::Controller* controller = new brpc::Controller();
//补充:设置Controller的timeout时间,默认是3秒
controller->set_timeout_ms(4000);
//设置回调函数
auto done = limerpc::ClosureFactory::create([controller,req,resp](){
std::unique_ptr<brpc::Controller> cntl_guard(controller);
std::unique_ptr<CommunicationInterface::UploadPhotoReq> req_guard(req);
std::unique_ptr<CommunicationInterface::UploadPhotoResp> res_guard(resp);
if (cntl_guard->Failed()) {
std::cerr << "rpc远程调用失败: " << cntl_guard->ErrorText() << std::endl;
return;
}
auto result = resp->mutable_result();
DBG("上传图片测试成功,返回fileId为:{} ",result->fileid());
});
stub.uploadPhoto(controller, req, resp, done);
}
void DownloadPhotoTest(CommunicationInterface::FileService_Stub& stub,const std::string file_id)
{
//下载图片测试
CommunicationInterface::DownloadPhotoReq* req = new CommunicationInterface::DownloadPhotoReq();
req->set_requestid("下载图片测试");
req->set_sessionid("114514");
req->set_fileid(file_id);
//发起rpc调用-异步调用
CommunicationInterface::DownloadPhotoResp* resp = new CommunicationInterface::DownloadPhotoResp();
brpc::Controller* controller = new brpc::Controller();
//补充:设置Controller的timeout时间,默认是3秒
controller->set_timeout_ms(4000);
//设置回调函数
auto done = limerpc::ClosureFactory::create([controller,req,resp](){
std::unique_ptr<brpc::Controller> cntl_guard(controller);
std::unique_ptr<CommunicationInterface::DownloadPhotoReq> req_guard(req);
std::unique_ptr<CommunicationInterface::DownloadPhotoResp> res_guard(resp);
if (cntl_guard->Failed()) {
std::cerr << "rpc远程调用失败: " << cntl_guard->ErrorText() << std::endl;
return;
}
auto result = resp->mutable_result();
DBG("下载图片测试成功,返回fileId为:{},保存文件到当前目录",result->fileid());
std::string file_data = result->filedata();
if(!limeutil::LimeFile::write("./111bak.jpg",file_data))
{
ERR("文件保存失败,请检查文件路径是否正确!");
}
});
stub.downloadPhoto(controller, req, resp, done);
}
void UploadVideoTest(CommunicationInterface::FileService_Stub& stub)
{
//上传视频测试
CommunicationInterface::UploadVideoReq* req = new CommunicationInterface::UploadVideoReq();
req->set_requestid("上传视频测试");
req->set_sessionid("114514");
auto file_info = req->mutable_fileinfo();
file_info->set_filemime("video/mp4");
std::string file_data;
if(!limeutil::LimeFile::read("/home/dev/workspace/cpp-videoplayer/server/svc_file/client/test.mp4",file_data))
{
ERR("文件读取失败,请检查文件路径是否正确!");
abort();
}
file_info->set_filedata(file_data);
//发起rpc调用-异步调用
CommunicationInterface::UploadVideoResp* resp = new CommunicationInterface::UploadVideoResp();
brpc::Controller* controller = new brpc::Controller();
//补充:设置Controller的timeout时间,默认是3秒
controller->set_timeout_ms(4000);
//设置回调函数
auto done = limerpc::ClosureFactory::create([controller,req,resp](){
std::unique_ptr<brpc::Controller> cntl_guard(controller);
std::unique_ptr<CommunicationInterface::UploadVideoReq> req_guard(req);
std::unique_ptr<CommunicationInterface::UploadVideoResp> res_guard(resp);
if (cntl_guard->Failed()) {
std::cerr << "rpc远程调用失败: " << cntl_guard->ErrorText() << std::endl;
return;
}
auto result = resp->mutable_result();
DBG("上传视频测试成功,返回fileId为:{} ",result->fileid());
});
stub.uploadVideo(controller, req, resp, done);
}
void DownloadVideoTest(CommunicationInterface::FileService_Stub& stub,const std::string file_id)
{
//下载视频测试
CommunicationInterface::DownloadVideoReq* req = new CommunicationInterface::DownloadVideoReq();
req->set_requestid("下载视频测试");
req->set_sessionid("114514");
req->set_fileid(file_id);
//发起rpc调用-异步调用
CommunicationInterface::DownloadVideoResp* resp = new CommunicationInterface::DownloadVideoResp();
brpc::Controller* controller = new brpc::Controller();
//补充:设置Controller的timeout时间,默认是3秒
controller->set_timeout_ms(4000);
//设置回调函数
auto done = limerpc::ClosureFactory::create([controller,req,resp](){
std::unique_ptr<brpc::Controller> cntl_guard(controller);
std::unique_ptr<CommunicationInterface::DownloadVideoReq> req_guard(req);
std::unique_ptr<CommunicationInterface::DownloadVideoResp> res_guard(resp);
if (cntl_guard->Failed()) {
std::cerr << "rpc远程调用失败: " << cntl_guard->ErrorText() << std::endl;
return;
}
auto result = resp->mutable_result();
DBG("下载图片测试成功,返回fileId为:{},保存文件到当前目录",result->fileid());
std::string file_data = result->filedata();
if(!limeutil::LimeFile::write("./testbak.mp4",file_data))
{
ERR("文件保存失败,请检查文件路径是否正确!");
}
});
stub.downloadVideo(controller, req, resp, done);
}
void mqDeleteFildTest(const std::string& file_id)
{
//通过消息队列向目标服务发送删除文件请求
//1.创建消息队列客户端
auto mq_client = std::make_shared<limemq::MQClient>(FLAGS_mq_host);
//2.解析json格式的mq配置
auto del_setting_json = limeutil::LimeJson::deserialize(FLAGS_mq_del_settings);
if(!del_setting_json)
{
ERR("解析mq配置失败");
return;
}
limemq::declare_settings del_settings{
.exchange = (*del_setting_json)["exchange"].asString(),
.exchange_type = (*del_setting_json)["exchange_type"].asString(),
.queue = (*del_setting_json)["queue"].asString(),
.binding_key = (*del_setting_json)["binding_key"].asString(),
};
//3.创建消息队列发布者
auto publisher = std::make_shared<limemq::Publisher>(mq_client, del_settings);
//4.发布消息-到文件删除队列
CommunicationInterface::DeleteFileMsg delete_file_msg;
delete_file_msg.add_file_id(file_id);
publisher->publish(delete_file_msg.SerializeAsString());
DBG("文件删除请求已发送,等待服务端响应...");
}
int main(int argc, char* argv[])
{
limelog::limelog_init();
//三步走:1.在服务注册中心找到目标服务;
//创建服务管理类
limerpc::SvcRpcChannels svc_rpc_channels;
std::string svc_name = FLAGS_svc_name;
//添加服务关心
svc_rpc_channels.set_match(svc_name);
//通过etcd模拟服务发现-自动添加服务关心的结点
std::string reg_center_addr = FLAGS_reg_center_addr;
auto online_callback = std::bind(&limerpc::SvcRpcChannels::add_node, &svc_rpc_channels, std::placeholders::_1, std::placeholders::_2);
auto offline_callback = std::bind(&limerpc::SvcRpcChannels::remove_node, &svc_rpc_channels, std::placeholders::_1, std::placeholders::_2);
limeetcd::SvcWatcher svc_watcher(reg_center_addr, online_callback, offline_callback);
svc_watcher.watch();
//获取服务信道
limerpc::ChannelPtr channel;
while(!channel){
WRN("无合适的服务信道,等待服务上线...");
std::this_thread::sleep_for(std::chrono::seconds(1));
channel = svc_rpc_channels.get_channel(svc_name);
}
DBG("已获取到服务信道");
//创建stub对象-用于发起rpc调用
CommunicationInterface::FileService_Stub stub(channel.get());
//3.进行如下测试
if(argc > 1 && std::string(argv[1]) == "upload_photo")
UploadPhotoTest(stub);
else if(argc > 2 && std::string(argv[1]) == "download_photo")
DownloadPhotoTest(stub,argv[2]);
else if(argc > 1 && std::string(argv[1]) == "upload_video")
UploadVideoTest(stub);
else if(argc > 2 && std::string(argv[1]) == "download_video")
DownloadVideoTest(stub,argv[2]);
//4.通过消息队列向目标服务发送删除文件请求
else if(argc > 2 && std::string(argv[1]) == "mq_delete_file")
mqDeleteFildTest(argv[2]);
DBG("测试消息已发出,等待服务端响应,按回车键退出...");
getchar();
return 0;
}
编译构建
# 1. 声明cmake所需版本
cmake_minimum_required(VERSION 3.1.3)
# 2. 设置工程项目名称(内部会生成一系列的内置变量)
project(main VERSION 1.0)
set(CMAKE_CXX_FLAGS "${CMAKE_CXX_FLAGS} -g -std=c++17")
# 4. 执行外部命令,生成临时代码文件: cal.pb.h cal.pb.cc
set(data_dir ${CMAKE_CURRENT_SOURCE_DIR}/../data)
set(proto_dir ${CMAKE_CURRENT_SOURCE_DIR}/../proto)
set(common_dir ${CMAKE_CURRENT_SOURCE_DIR}/../common)
set(source_dir ${CMAKE_CURRENT_SOURCE_DIR}/source)
set(client_dir ${CMAKE_CURRENT_SOURCE_DIR}/client)
set(odb_files ${data_dir}/data.h)
file(GLOB_RECURSE proto_files ${proto_dir}/base.proto ${proto_dir}/file.proto ${proto_dir}/message.proto)
execute_process(
COMMAND odb -d mysql --std c++11 --generate-query --generate-schema --profile boost/date-time ${odb_files}
COMMAND protoc --experimental_allow_proto3_optional -I ${proto_dir} --cpp_out=${CMAKE_CURRENT_BINARY_DIR} ${proto_files}
)
# 3. 添加生成目标(可执行程序)
file(GLOB_RECURSE common_files ${common_dir}/*.cc)
file(GLOB_RECURSE data_files ${data_dir}/*.cc)
file(GLOB_RECURSE temp_files ${CMAKE_CURRENT_BINARY_DIR}/*.cc ${CMAKE_CURRENT_BINARY_DIR}/*.cxx)
file(GLOB_RECURSE source_files ${source_dir}/*.cc)
set(main_file ${CMAKE_CURRENT_SOURCE_DIR}/source/main.cc)
add_executable(${PROJECT_NAME} ${temp_files} ${data_files} ${main_file} ${common_files} ${source_files})
# 4. 添加生成目标(测试客户端程序)
file(GLOB_RECURSE client_files ${client_dir}/*.cc)
add_executable(file_client ${temp_files} ${data_files} ${client_files})
target_link_libraries(file_client PRIVATE lime_scaffold::lime_scaffold)
# 5. 头文件路径设置
include_directories(${CMAKE_CURRENT_BINARY_DIR})
include_directories(${data_dir})
include_directories(${common_dir})
include_directories(${source_dir})
# 6. 依赖检查(依赖库的查找以及链接选项的添加)
find_package(lime_scaffold REQUIRED)
# 7. 设置链接库
target_link_libraries(${PROJECT_NAME} PRIVATE lime_scaffold::lime_scaffold)
火山引擎视频云技术社区,是面向 AI 音视频开发者的技术交流平台。这里汇聚源自抖音、豆包等亿级 DAU 产品的 RTC、直播、点播、AI 媒体处理、音视频互动技术,提供接入指南、最佳实践、性能调优、场景案例、Demo 代码、开源项目、白皮书和 API 文档。社区汇聚官方工程师与一线开发者,为 AI 视频通话、数字人、AI 视频处理等应用的开发与落地提供技术支持。
更多推荐
所有评论(0)