[微服务即时通讯系统]2.服务端-环境搭建
本专栏内容为:项目专栏
💓博主csdn个人主页:小小unicorn
⏩专栏分类:微服务即时通讯系统
🚚代码仓库:小小unicorn的代码仓库🚚
🌹🌹🌹关注我带你学习编程知识
环境安装
brpc
介绍
brpc是用 c++语言编写的工业级 RPC 框架,常用于搜索、存储、机器学习、广告、推荐等高性能系统。
你可以使用它:
• 搭建能在一个端口支持多协议的服务, 或访问各种服务restful http/https, h2/gRPC。使用 brpc 的 http 实现比 libcurl 方便多了。从其他语言通过 HTTP/h2+json 访问基于 protobuf 的协议.
redis和memcached,线程安全,比官方 client更方便。rtmp/flv/hls,可用于搭建流媒体服务.- 支持
thrift, 线程安全,比官方client更方便 - 各种百度内使用的协议:
baidu_std, streaming_rpc, hulu_pbrpc, sofa_pbrpc, nova_pbrpc, public_pbrpc, ubrpc和使用nshead的各种协议. - 基于工业级的
RAFT算法实现搭建高可用分布式系统,已在braft开源。
•Server能同步或异步处理请求。
•Client支持同步、 异步、 半同步,或使用组合channels简化复杂的分库或并发访
问。
• 通过http界面调试服务, 使用cpu, heap, contention profilers.
• 获得更好的延时和吞吐.
• 把你组织中使用的协议快速地加入brpc,或定制各类组件, 包括命名服务 (dns, zk, etcd), 负载均衡 (rr, random, consistent hashing)
安装
先安装依赖:
sudo apt-get install -y git g++ make libssl-dev libprotobuf-dev libprotoc-dev protobuf-compiler libleveldb-dev
安装完成后:

接下来安装brpc:
git clone https://github.com/apache/brpc.git
安装完成后会自动生成一个brpc目录:

在这个目录下:
mkdir build && cd build
cmake -DCMAKE_INSTALL_PREFIX=/usr .. && cmake --build . -j6
make && sudo make install
运行结果如下:

使用
我们可以先进行一个Protobuf的测试:
syntax="proto3";
package example;
option cc_generic_services = true;
message EchoRequest {
string message = 1;
}
message EchoResponse {
string message = 1;
}
service EchoService {
rpc Echo(EchoRequest) returns (EchoResponse);
}
运行一下代码:
protoc --cpp_out=./ main.proto
运行结果:

会在当前目录下生成两个文件。
使用样例:
测试用例:
server.cc
#include <brpc/server.h>
#include <butil/logging.h>
#include "main.pb.h"
//1. 继承于EchoService创建一个子类,并实现rpc调用的业务功能
class EchoServiceImpl : public example::EchoService {
public:
EchoServiceImpl(){}
~EchoServiceImpl(){}
void Echo(google::protobuf::RpcController* controller,
const ::example::EchoRequest* request,
::example::EchoResponse* response,
::google::protobuf::Closure* done) {//业务逻辑
brpc::ClosureGuard rpc_guard(done);
std::cout << "收到消息:" << request->message() << std::endl;
//构造响应
std::string str = request->message() + "--这是响应!!";
response->set_message(str);
//done->Run();
}
};
int main(int argc, char *argv[])
{
//关闭brpc的默认日志输出
logging::LoggingSettings settings;
settings.logging_dest = logging::LoggingDestination::LOG_TO_NONE;
logging::InitLogging(settings);
//2. 构造服务器对象
brpc::Server server;
//3. 向服务器对象中,新增EchoService服务
EchoServiceImpl echo_service;
int ret = server.AddService(&echo_service, brpc::ServiceOwnership::SERVER_DOESNT_OWN_SERVICE);
if (ret == -1) {
std::cout << "添加Rpc服务失败!\n";
return -1;
}
//4. 启动服务器
brpc::ServerOptions options;
options.idle_timeout_sec = -1; //连接空闲超时时间-超时后连接被关闭
options.num_threads = 1; // io线程数量
ret = server.Start(8080, &options);
if (ret == -1) {
std::cout << "启动服务器失败!\n";
return -1;
}
server.RunUntilAskedToQuit();//修改等待运行结束
return 0;
}
client.cc
#include <brpc/channel.h>
#include <thread>
#include "main.pb.h"
void callback(brpc::Controller* cntl, ::example::EchoResponse* response) {
std::unique_ptr<brpc::Controller> cntl_guard(cntl);
std::unique_ptr<example::EchoResponse> resp_guard(response);
if (cntl->Failed() == true) {
std::cout << "Rpc调用失败:" << cntl->ErrorText() << std::endl;
return;
}
std::cout << "收到响应: " << response->message() << std::endl;
}
int main(int argc, char *argv[])
{
//1. 构造Channel信道,连接服务器
brpc::ChannelOptions options;
options.connect_timeout_ms = -1;// 连接等待超时时间,-1表示一直等待
options.timeout_ms = -1; //rpc请求等待超时时间,-1表示一直等待
options.max_retry = 3;//请求重试次数
options.protocol = "baidu_std"; //序列化协议,默认使用baidu_std(百度内定协议)
brpc::Channel channel;
int ret = channel.Init("127.0.0.1:8080", &options);
if (ret == -1) {
std::cout << "初始化信道失败!\n";
return -1;
}
//2. 构造EchoService_Stub对象,用于进行rpc调用
example::EchoService_Stub stub(&channel);
//3. 进行Rpc调用/
example::EchoRequest req;
req.set_message("你好~比特~!");
brpc::Controller *cntl = new brpc::Controller();
example::EchoResponse *rsp = new example::EchoResponse();
// stub.Echo(cntl, &req, rsp, nullptr);
// if (cntl->Failed() == true) {
// std::cout << "Rpc调用失败:" << cntl->ErrorText() << std::endl;
// return -1;
// }
// std::cout << "收到响应: " << rsp->message() << std::endl;
// delete cntl;
// delete rsp;
auto clusure = google::protobuf::NewCallback(callback, cntl, rsp);
stub.Echo(cntl, &req, rsp, clusure); //异步调用
std::cout << "异步调用结束!\n";
std::this_thread::sleep_for(std::chrono::seconds(3));
return 0;
}
makefile
all : server client
server : server.cc main.pb.cc
g++ -std=c++17 $^ -o $@ -lbrpc -lgflags -lssl -lcrypto -lprotobuf -lleveldb
client : client.cc main.pb.cc
g++ -std=c++17 $^ -o $@ -lbrpc -lgflags -lssl -lcrypto -lprotobuf -lleveldb
编译完成后,启动两个终端,分别运行服务端和客户端:


这里可以看到异步通信是完全Ok的!
brpc信道管理封装
channel.hpp
#pragma once
#include <brpc/channel.h>
#include <string>
#include <vector>
#include <unordered_map>
#include <mutex>
#include "logger.hpp"
namespace bite_im
{
// 1. 封装单个服务的信道管理类:
class ServiceChannel
{
public:
using ptr = std::shared_ptr<ServiceChannel>;
using ChannelPtr = std::shared_ptr<brpc::Channel>;
ServiceChannel(const std::string &name) : _service_name(name), _index(0) {}
// 服务上线了一个节点,则调用append新增信道
void append(const std::string &host)
{
auto channel = std::make_shared<brpc::Channel>();
brpc::ChannelOptions options;
options.connect_timeout_ms = -1; // 连接等待超时时间,-1表示一直等待
options.timeout_ms = -1; // rpc请求等待超时时间,-1表示一直等待
options.max_retry = 3; // 请求重试次数
options.protocol = "baidu_std"; // 序列化协议,默认使用baidu_std
int ret = channel->Init(host.c_str(), &options);
if (ret == -1)
{
LOG_ERROR("初始化{}-{}信道失败!", _service_name, host);
return;
}
std::unique_lock<std::mutex> lock(_mutex);
_hosts.insert(std::make_pair(host, channel));
_channels.push_back(channel);
}
// 服务下线了一个节点,则调用remove释放信道
void remove(const std::string &host)
{
std::unique_lock<std::mutex> lock(_mutex); // 先加锁
auto it = _hosts.find(host);
if (it == _hosts.end())
{
LOG_WARN("{}-{}节点删除信道时,没有找到信道信息!", _service_name, host); // 警告
return;
}
for (auto vit = _channels.begin(); vit != _channels.end(); ++vit)
{
if (*vit == it->second)
{
_channels.erase(vit);
break;
}
}
_hosts.erase(it);
}
// 通过RR轮转策略,获取一个Channel用于发起对应服务的Rpc调用
ChannelPtr choose()
{
std::unique_lock<std::mutex> lock(_mutex);
if (_channels.size() == 0)
{
LOG_ERROR("当前没有能够提供 {} 服务的节点!", _service_name);
return ChannelPtr();
}
int32_t idx = _index++ % _channels.size();
return _channels[idx]; // 注意不能越界
}
private:
std::mutex _mutex; // 互斥锁
int32_t _index; // 当前轮转下标计数器
std::string _service_name; // 服务名称
std::vector<ChannelPtr> _channels; // 当前服务对应的信道集合
std::unordered_map<std::string, ChannelPtr> _hosts; // 主机地址与信道映射关系
};
// 总体的服务信道管理类
class ServiceManager
{
public:
using ptr = std::shared_ptr<ServiceManager>; // 只能指针
ServiceManager() {}
// 获取指定服务的节点信道
ServiceChannel::ChannelPtr choose(const std::string &service_name)
{
std::unique_lock<std::mutex> lock(_mutex); // 先加锁
auto sit = _services.find(service_name);
if (sit == _services.end())
{
LOG_ERROR("当前没有能够提供 {} 服务的节点!", service_name);
return ServiceChannel::ChannelPtr();
}
return sit->second->choose();
}
// 先声明,我关注哪些服务的上下线,不关心的就不需要管理了
void declared(const std::string &service_name)
{
std::unique_lock<std::mutex> lock(_mutex);
_follow_services.insert(service_name);
}
// 服务上线时调用的回调接口,将服务节点管理起来
void onServiceOnline(const std::string &service_instance, const std::string &host)
{
std::string service_name = getServiceName(service_instance);
ServiceChannel::ptr service;
{
std::unique_lock<std::mutex> lock(_mutex);
auto fit = _follow_services.find(service_name);
if (fit == _follow_services.end())
{
LOG_DEBUG("{}-{} 服务上线了,但是当前并不关心!", service_name, host);
return;
}
// 先获取管理对象,没有则创建,有则添加节点
auto sit = _services.find(service_name);
if (sit == _services.end())
{
service = std::make_shared<ServiceChannel>(service_name);
_services.insert(std::make_pair(service_name, service));
}
else
{
service = sit->second;
}
}
if (!service)
{
LOG_ERROR("新增 {} 服务管理节点失败!", service_name);
return;
}
service->append(host);
LOG_DEBUG("{}-{} 服务上线新节点,进行添加管理!", service_name, host);
}
// 服务下线时调用的回调接口,从服务信道管理中,删除指定节点信道
void onServiceOffline(const std::string &service_instance, const std::string &host)
{
std::string service_name = getServiceName(service_instance);
ServiceChannel::ptr service;
{
std::unique_lock<std::mutex> lock(_mutex);
auto fit = _follow_services.find(service_name);
if (fit == _follow_services.end())
{
LOG_DEBUG("{}-{} 服务下线了,但是当前并不关心!", service_name, host);
return;
}
// 先获取管理对象,没有则创建,有则添加节点
auto sit = _services.find(service_name);
if (sit == _services.end())
{
LOG_WARN("删除{}服务节点时,没有找到管理对象", service_name);
return;
}
service = sit->second;
}
service->remove(host);
LOG_DEBUG("{}-{} 服务下线节点,进行删除管理!", service_name, host);
}
private:
std::string getServiceName(const std::string &service_instance)
{
auto pos = service_instance.find_last_of('/');
if (pos == std::string::npos)
return service_instance;
return service_instance.substr(0, pos);
}
private:
std::mutex _mutex;
std::unordered_set<std::string> _follow_services;
std::unordered_map<std::string, ServiceChannel::ptr> _services;
};
}
brpc&&etcd联调测试
discover.cc
#include "../common/etcd.hpp"
#include "../common/channel.hpp"
#include <gflags/gflags.h>
#include <thread>
#include "main.pb.h"
DEFINE_bool(run_mode, false, "程序的运行模式,false-调试; true-发布;");
DEFINE_string(log_file, "", "发布模式下,用于指定日志的输出文件");
DEFINE_int32(log_level, 0, "发布模式下,用于指定日志输出等级");
DEFINE_string(etcd_host, "http://127.0.0.1:2379", "服务注册中心地址");
DEFINE_string(base_service, "/service", "服务监控根目录");
DEFINE_string(call_service, "/service/echo", "服务监控根目录");
int main(int argc, char *argv[])
{
google::ParseCommandLineFlags(&argc, &argv, true);
bite_im::init_logger(FLAGS_run_mode, FLAGS_log_file, FLAGS_log_level);
// 1. 先构造Rpc信道管理对象
auto sm = std::make_shared<ServiceManager>();
sm->declared(FLAGS_call_service);
auto put_cb = std::bind(&ServiceManager::onServiceOnline, sm.get(), std::placeholders::_1, std::placeholders::_2);
auto del_cb = std::bind(&ServiceManager::onServiceOffline, sm.get(), std::placeholders::_1, std::placeholders::_2);
// 2. 构造服务发现对象
Discovery::ptr dclient = std::make_shared<Discovery>(FLAGS_etcd_host, FLAGS_base_service, put_cb, del_cb);
while (1)
{
// 3. 通过Rpc信道管理对象,获取提供Echo服务的信道
auto channel = sm->choose(FLAGS_call_service);
if (!channel)
{
std::this_thread::sleep_for(std::chrono::seconds(1));
return -1;
}
// 4. 发起EchoRpc调用
example::EchoService_Stub stub(channel.get());
example::EchoRequest req;
req.set_message("你好~比特~!");
brpc::Controller *cntl = new brpc::Controller();
example::EchoResponse *rsp = new example::EchoResponse();
stub.Echo(cntl, &req, rsp, nullptr);
if (cntl->Failed() == true)
{
std::cout << "Rpc调用失败:" << cntl->ErrorText() << std::endl;
delete cntl;
delete rsp;
std::this_thread::sleep_for(std::chrono::seconds(1));
continue;
}
std::cout << "收到响应: " << rsp->message() << std::endl;
std::this_thread::sleep_for(std::chrono::seconds(1));
}
return 0;
}
registry.cc
#include "../common/etcd.hpp"
#include <gflags/gflags.h>
#include <thread>
#include <brpc/server.h>
#include <butil/logging.h>
#include "main.pb.h"
DEFINE_bool(run_mode, false, "程序的运行模式,false-调试; true-发布;");
DEFINE_string(log_file, "", "发布模式下,用于指定日志的输出文件");
DEFINE_int32(log_level, 0, "发布模式下,用于指定日志输出等级");
DEFINE_string(etcd_host, "http://127.0.0.1:2379", "服务注册中心地址");
DEFINE_string(base_service, "/service", "服务监控根目录");
DEFINE_string(instance_name, "/echo/instance", "当前实例名称");
DEFINE_string(access_host, "127.0.0.1:7070", "当前实例的外部访问地址");
DEFINE_int32(listen_port, 7070, "Rpc服务器监听端口");
class EchoServiceImpl : public example::EchoService
{
public:
EchoServiceImpl() {}
~EchoServiceImpl() {}
void Echo(google::protobuf::RpcController *controller,
const ::example::EchoRequest *request,
::example::EchoResponse *response,
::google::protobuf::Closure *done)
{
brpc::ClosureGuard rpc_guard(done);
std::cout << "收到消息:" << request->message() << std::endl;
std::string str = request->message() + "--这是响应!!";
response->set_message(str);
// done->Run();
}
};
int main(int argc, char *argv[])
{
google::ParseCommandLineFlags(&argc, &argv, true);
bite_im::init_logger(FLAGS_run_mode, FLAGS_log_file, FLAGS_log_level);
// 关闭brpc的默认日志输出
logging::LoggingSettings settings;
settings.logging_dest = logging::LoggingDestination::LOG_TO_NONE;
logging::InitLogging(settings);
// 2. 构造服务器对象
brpc::Server server;
// 3. 向服务器对象中,新增EchoService服务
EchoServiceImpl echo_service;
int ret = server.AddService(&echo_service, brpc::ServiceOwnership::SERVER_DOESNT_OWN_SERVICE);
if (ret == -1)
{
std::cout << "添加Rpc服务失败!\n";
return -1;
}
// 4. 启动服务器
brpc::ServerOptions options;
options.idle_timeout_sec = -1; // 连接空闲超时时间-超时后连接被关闭
options.num_threads = 1; // io线程数量
ret = server.Start(FLAGS_listen_port, &options);
if (ret == -1)
{
std::cout << "启动服务器失败!\n";
return -1;
}
// 4. 注册服务;
Registry::ptr rclient = std::make_shared<Registry>(FLAGS_etcd_host);
// 注册的时 /service/echo/instance ;; 发现的时候 /service/echo
rclient->registry(FLAGS_base_service + FLAGS_instance_name, FLAGS_access_host);
server.RunUntilAskedToQuit(); // 休眠等待运行结束
return 0;
}
makfile
all:discovery registry
discovery : discovery.cc main.pb.cc
g++ -g -std=c++17 $^ -o $@ -lspdlog -lfmt -lgflags -letcd-cpp-api -lcpprest -lbrpc -lssl -lcrypto -lprotobuf -lleveldb
registry : registry.cc main.pb.cc
g++ -g -std=c++17 $^ -o $@ -lspdlog -lfmt -lgflags -letcd-cpp-api -lcpprest -lbrpc -lssl -lcrypto -lprotobuf -lleveldb
编译运行:


这里的联调测试是没有问题的!!
es
介绍
Elasticsearch, 简称 ES,它是个开源分布式搜索引擎,它的特点有:分布式,零配
置,自动发现,索引自动分片,索引副本机制,restful 风格接口,多数据源,自动搜
索负载等。它可以近乎实时的存储、检索数据;本身扩展性很好,可以扩展到上百台
服务器,处理 PB 级别的数据。 es 也使用Java开发并使用Lucene作为其核心来实现
所有索引和搜索的功能,但是它的目的是通过简单的RESTful API来隐藏 Lucene 的
复杂性,从而让全文搜索变得简单。
Elasticsearch 是面向文档(document oriented)的,这意味着它可以存储整个对象或文
档(document)。然而它不仅仅是存储,还会索引(index)每个文档的内容使之可以被搜
索。在 Elasticsearch 中,你可以对文档(而非成行成列的数据)进行索引、搜索、排
序、过滤
安装es:
添加仓库密钥
wget -qO - https://artifacts.elastic.co/GPG-KEY-elasticsearch | sudo apt-key add -

里面有一个警告不用管
添加镜像源仓库
echo "deb https://artifacts.elastic.co/packages/7.x/apt stable main" | sudo tee /etc/apt/sources.list.d/elasticsearch.list
再更新软件包列表:
sudo apt update

更新完成后:

安装es
sudo apt-get install elasticsearch=7.17.21
启动es并安装中文分词
sudo systemctl start elasticsearch
sudo /usr/share/elasticsearch/bin/elasticsearch-plugin install https://get.infini.cloud/elasticsearch/analysis-ik/7.17.21
启动es:

安装分词插件结果:

这里可以看到我们的es访问成功的
验证es是否安装成功
sudo netstat -anptu |grep 9200
curl -X GET "http://localhost:9200/"

设置外网访问:如果新配置完成的话,默认只能在本机进行访问
vim /etc/elasticsearch/elasticsearch.yml
# 新增配置
network.host: 0.0.0.0
http.port: 9200
cluster.initial_master_nodes: ["node-1"]
安装kibana
安装 Kibana:
使用 apt 命令安装 Kibana
sudo apt install kibana

配置 Kibana(可选):
根据需要配置 Kibana。
配置文件通常位于/etc/kibana/kibana.yml。可能需要设置如服务器地址、端口、 Elasticsearch URL 等。
sudo vim /etc/kibana/kibana.yml

启动 Kibana 服务:
安装完成后,启动 Kibana 服务。
sudo systemctl start kibana
设置开机自启(可选):
如果你希望 Kibana 在系统启动时自动启动,可以使用以下命令来启用自启动
sudo systemctl enable kibana
验证安装:
使用以下命令检查 Kibana 服务的状态。
sudo systemctl status kibana

一般情况下,监听的是5601端口,我们可以查看一下:
sudo netstat -anptu |grep 5601

访问 Kibana:
在浏览器中访问 Kibana,通常是 http://yourip:5601
es客户端的安装
ES C++的客户端选择并不多, 我们这里使用 elasticlient 库, 下面进行安装
# 克隆代码
git clone https://github.com/seznam/elasticlient

# 切换目录
cd elasticlient
# 更新子模块
git submodule update --init --recursive
# 编译代码
mkdir build
cd build
#安装 MicroHTTPD 库
sudo apt-get install libmicrohttpd-dev
#在cmake之前先手动安装子模块
cd ../external/googletest/
mkdir cmake && cd cmake/
cmake -DCMAKE_INSTALL_PREFIX=/usr ..
#然后在返回到build目录下进行cmake
cmake -DCMAKE_INSTALL_PREFIX=/usr ..

#cmake结束后,在进行Make和sudo make install
make && sudo make install

到这里,我们的环境就彻底安装结束了。。。
Kibana 访问 es 进行测试
通过网页访问 kibana:
创建索引库,
POST /user/_doc
{
"settings" : {
"analysis" : {
"analyzer" : {
"ik" : {
"tokenizer" : "ik_max_word"
}
}
}
},
"mappings" : {
"dynamic" : true,
"properties" : {
"nickname" : {
"type" : "text",
"analyzer" : "ik_max_word"
},
"user_id" : {
"type" : "keyword",
"analyzer" : "standard"
},
"phone" : {
"type" : "keyword",
"analyzer" : "standard"
},
"description" : {
"type" : "text",
"enabled" : false
},
"avatar_id" : {
"type" : "keyword"
"enabled" : false
}
}
}
}
各个字段功能:

运行如下:

这里我们的索引就已经创建好了。
新增数据:
POST /user/_doc/_bulk
{"index":{"_id":"1"}}
{"user_id" : "USER4b862aaa-2df8654a-7eb4bb65-e3507f66","nickname" : "昵称 1","phone" : "手机号 1","description" :"签名 1","avatar_id" : "头像 1"}
{"index":{"_id":"2"}}
{"user_id" : "USER14eeeaa5-442771b9-0262e455-e4663d1d","nickname" : "昵称 2","phone" : "手机号 2","description" :"签名 2","avatar_id" : "头像 2"}
{"index":{"_id":"3"}}
{"user_id" : "USER484a6734-03a124f0-996c169dd05c1869","nickname" : "昵称 3","phone" : "手机号 3","description" :"签名 3","avatar_id" : "头像 3"}
{"index":{"_id":"4"}}
{"user_id" : "USER186ade83-4460d4a6-8c08068f-83127b5d","nickname" : "昵称 4","phone" : "手机号 4","description" :"签名 4","avatar_id" : "头像 4"}
{"index":{"_id":"5"}}
{"user_id" : "USER6f19d074-c33891cf-23bf5a83-57189a19","nickname" : "昵称 5","phone" : "手机号 5","description" :"签名 5","avatar_id" : "头像 5"}
{"index":{"_id":"6"}}
{"user_id" : "USER97605c64-9833ebb7-d0455353-35a59195","nickname" : "昵称 6","phone" : "手机号 6","description" :"签名 6","avatar_id" : "头像 6"}
运行结果:

我们查询一下所有数据:
POST /user/_doc/_search
{
"query":{
"match_all":{}
}
}

这里我们有多少条数据,就会查出来多少数据。
我们还可以进行过滤搜索:
GET /user/_doc/_search?pretty
{
"query" : {
"bool" : {
"must_not" : [
{
"terms" : {
"user_id.keyword" : [
"USER4b862aaa-2df8654a-7eb4bb65-e3507f66",
"USER14eeeaa5-442771b9-0262e455-e4663d1d",
"USER484a6734-03a124f0-996c169dd05c1869"
]
}
}
],
"should" : [
{
"match" : {
"user_id" : "昵称"
}
},
{
"match" : {
"nickname" : "昵称"
}
},
{
"match" : {
"phone" : "昵称"
}
}
]
}
}
}

删除索引:
DELETE /user

jsoncpp的使用
main.cc
#include <json/json.h>
#include <iostream>
#include <sstream>
#include <memory>
bool Serialize(const Json::Value &val, std::string &dst)
{
// 先定义Json::StreamWriter 工厂类 Json::StreamWriterBuilder
Json::StreamWriterBuilder swb;
swb.settings_["emitUTF8"] = true;
std::unique_ptr<Json::StreamWriter> sw(swb.newStreamWriter());
// 通过Json::StreamWriter中的write接口进行序列化
std::stringstream ss;
int ret = sw->write(val, &ss);
if (ret != 0)
{
std::cout << "Json反序列化失败!\n";
return false;
}
dst = ss.str();
return true;
}
bool UnSerialize(const std::string &src, Json::Value &val)
{
Json::CharReaderBuilder crb;
std::unique_ptr<Json::CharReader> cr(crb.newCharReader());
std::string err;
bool ret = cr->parse(src.c_str(), src.c_str() + src.size(), &val, &err);
if (ret == false)
{
std::cout << "json反序列化失败: " << err << std::endl;
return false;
}
return true;
}
int main()
{
char name[] = "张三";
int age = 18;
float score[3] = {88, 89.5, 99};
Json::Value stu;
stu["姓名"] = name;
stu["年龄"] = age;
stu["成绩"].append(score[0]);
stu["成绩"].append(score[1]);
stu["成绩"].append(score[2]);
std::string stu_str;
bool ret = Serialize(stu, stu_str);
if (ret == false)
return -1;
std::cout << stu_str << std::endl;
Json::Value val;
ret = UnSerialize(stu_str, val);
if (ret == false)
return -1;
std::cout << val["姓名"].asString() << std::endl;
std::cout << val["年龄"].asInt() << std::endl;
int sz = val["成绩"].size();
for (int i = 0; i < sz; i++)
{
std::cout << val["成绩"][i].asFloat() << std::endl;
}
return 0;
}
makefile
main : main.cc
g++ -std=c++17 $^ -o $@ /usr/lib/x86_64-linux-gnu/libjsoncpp.so.19

es客户端二次封装
icsearch.hpp
#pragma once
#include <elasticlient/client.h>
#include <cpr/cpr.h>
#include <json/json.h>
#include <iostream>
#include <memory>
#include "logger.hpp"
bool Serialize(const Json::Value &val, std::string &dst)
{
// 先定义Json::StreamWriter 工厂类 Json::StreamWriterBuilder
Json::StreamWriterBuilder swb;
swb.settings_["emitUTF8"] = true;
std::unique_ptr<Json::StreamWriter> sw(swb.newStreamWriter());
// 通过Json::StreamWriter中的write接口进行序列化
std::stringstream ss;
int ret = sw->write(val, &ss);
if (ret != 0)
{
std::cout << "Json反序列化失败!\n";
return false;
}
dst = ss.str();
return true;
}
bool UnSerialize(const std::string &src, Json::Value &val)
{
Json::CharReaderBuilder crb;
std::unique_ptr<Json::CharReader> cr(crb.newCharReader());
std::string err;
bool ret = cr->parse(src.c_str(), src.c_str() + src.size(), &val, &err);
if (ret == false)
{
std::cout << "json反序列化失败: " << err << std::endl;
return false;
}
return true;
}
class ESIndex
{
public:
ESIndex(std::shared_ptr<elasticlient::Client> &client,
const std::string &name,
const std::string &type = "_doc") : _name(name), _type(type), _client(client)
{
Json::Value analysis;
Json::Value analyzer;
Json::Value ik;
Json::Value tokenizer;
tokenizer["tokenizer"] = "ik_max_word";
ik["ik"] = tokenizer;
analyzer["analyzer"] = ik;
analysis["analysis"] = analyzer;
_index["settings"] = analysis;
}
ESIndex &append(const std::string &key,
const std::string &type = "text",
const std::string &analyzer = "ik_max_word",
bool enabled = true)
{
Json::Value fields;
fields["type"] = type;
fields["analyzer"] = analyzer;
if (enabled == false)
fields["enabled"] = enabled;
_properties[key] = fields;
return *this;
}
bool create(const std::string &index_id = "default_index_id")
{
Json::Value mappings;
mappings["dynamic"] = true;
mappings["properties"] = _properties;
_index["mappings"] = mappings;
std::string body;
bool ret = Serialize(_index, body);
if (ret == false)
{
LOG_ERROR("索引序列化失败!");
return false;
}
LOG_DEBUG("{}", body);
// 2. 发起搜索请求
try
{
auto rsp = _client->index(_name, _type, index_id, body);
if (rsp.status_code < 200 || rsp.status_code >= 300)
{
LOG_ERROR("创建ES索引 {} 失败,响应状态码异常: {}", _name, rsp.status_code);
return false;
}
}
catch (std::exception &e)
{
LOG_ERROR("创建ES索引 {} 失败: {}", _name, e.what());
return false;
}
return true;
}
private:
std::string _name;
std::string _type;
Json::Value _properties;
Json::Value _index;
std::shared_ptr<elasticlient::Client> _client;
};
class ESInsert
{
public:
ESInsert(std::shared_ptr<elasticlient::Client> &client,
const std::string &name,
const std::string &type = "_doc") : _name(name), _type(type), _client(client) {}
template <typename T>
ESInsert &append(const std::string &key, const T &val)
{
_item[key] = val;
return *this;
}
bool insert(const std::string id = "")
{
std::string body;
bool ret = Serialize(_item, body);
if (ret == false)
{
LOG_ERROR("索引序列化失败!");
return false;
}
LOG_DEBUG("{}", body);
// 2. 发起搜索请求
try
{
auto rsp = _client->index(_name, _type, id, body);
if (rsp.status_code < 200 || rsp.status_code >= 300)
{
LOG_ERROR("新增数据 {} 失败,响应状态码异常: {}", body, rsp.status_code);
return false;
}
}
catch (std::exception &e)
{
LOG_ERROR("新增数据 {} 失败: {}", body, e.what());
return false;
}
return true;
}
private:
std::string _name;
std::string _type;
Json::Value _item;
std::shared_ptr<elasticlient::Client> _client;
};
class ESRemove
{
public:
ESRemove(std::shared_ptr<elasticlient::Client> &client,
const std::string &name,
const std::string &type = "_doc") : _name(name), _type(type), _client(client) {}
bool remove(const std::string &id)
{
try
{
auto rsp = _client->remove(_name, _type, id);
if (rsp.status_code < 200 || rsp.status_code >= 300)
{
LOG_ERROR("删除数据 {} 失败,响应状态码异常: {}", id, rsp.status_code);
return false;
}
}
catch (std::exception &e)
{
LOG_ERROR("删除数据 {} 失败: {}", id, e.what());
return false;
}
return true;
}
private:
std::string _name;
std::string _type;
std::shared_ptr<elasticlient::Client> _client;
};
class ESSearch
{
public:
ESSearch(std::shared_ptr<elasticlient::Client> &client,
const std::string &name,
const std::string &type = "_doc") : _name(name), _type(type), _client(client) {}
ESSearch &append_must_not_terms(const std::string &key, const std::vector<std::string> &vals)
{
Json::Value fields;
for (const auto &val : vals)
{
fields[key].append(val);
}
Json::Value terms;
terms["terms"] = fields;
_must_not.append(terms);
return *this;
}
ESSearch &append_should_match(const std::string &key, const std::string &val)
{
Json::Value field;
field[key] = val;
Json::Value match;
match["match"] = field;
_should.append(match);
return *this;
}
ESSearch &append_must_term(const std::string &key, const std::string &val)
{
Json::Value field;
field[key] = val;
Json::Value term;
term["term"] = field;
_must.append(term);
return *this;
}
ESSearch &append_must_match(const std::string &key, const std::string &val)
{
Json::Value field;
field[key] = val;
Json::Value match;
match["match"] = field;
_must.append(match);
return *this;
}
Json::Value search()
{
Json::Value cond;
if (_must_not.empty() == false)
cond["must_not"] = _must_not;
if (_should.empty() == false)
cond["should"] = _should;
if (_must.empty() == false)
cond["must"] = _must;
Json::Value query;
query["bool"] = cond;
Json::Value root;
root["query"] = query;
std::string body;
bool ret = Serialize(root, body);
if (ret == false)
{
LOG_ERROR("索引序列化失败!");
return Json::Value();
}
LOG_DEBUG("{}", body);
// 2. 发起搜索请求
cpr::Response rsp;
try
{
rsp = _client->search(_name, _type, body);
if (rsp.status_code < 200 || rsp.status_code >= 300)
{
LOG_ERROR("检索数据 {} 失败,响应状态码异常: {}", body, rsp.status_code);
return Json::Value();
}
}
catch (std::exception &e)
{
LOG_ERROR("检索数据 {} 失败: {}", body, e.what());
return Json::Value();
}
// 3. 需要对响应正文进行反序列化
LOG_DEBUG("检索响应正文: [{}]", rsp.text);
Json::Value json_res;
ret = UnSerialize(rsp.text, json_res);
if (ret == false)
{
LOG_ERROR("检索数据 {} 结果反序列化失败", rsp.text);
return Json::Value();
}
return json_res["hits"]["hits"];
}
private:
std::string _name;
std::string _type;
Json::Value _must_not;
Json::Value _should;
Json::Value _must;
std::shared_ptr<elasticlient::Client> _client;
};
main.cc
#include "../common/icsearch.hpp"
#include <gflags/gflags.h>
DEFINE_bool(run_mode, false, "程序的运行模式,false-调试; true-发布;");
DEFINE_string(log_file, "", "发布模式下,用于指定日志的输出文件");
DEFINE_int32(log_level, 0, "发布模式下,用于指定日志输出等级");
int main(int argc, char *argv[])
{
google::ParseCommandLineFlags(&argc, &argv, true);
bite_im::init_logger(FLAGS_run_mode, FLAGS_log_file, FLAGS_log_level);
std::vector<std::string> host_list = {"http://127.0.0.1:9200/"};
auto client = std::make_shared<elasticlient::Client>(host_list);
std::cout << (uint64_t)client.get() << std::endl;
bool ret = ESIndex(client, "test_user").append("nickname").append("phone", "keyword", "standard", true).create();
if (ret == false)
{
LOG_INFO("索引创建失败!");
return -1;
}
else
{
LOG_INFO("索引创建成功!");
}
// 数据的新增
ret = ESInsert(client, "test_user")
.append("nickname", "张三")
.append("phone", "15566667777")
.insert("00001");
if (ret == false)
{
LOG_ERROR("数据插入失败!");
return -1;
}
else
{
LOG_INFO("数据新增成功!");
}
// 数据的修改
ret = ESInsert(client, "test_user")
.append("nickname", "张三")
.append("phone", "13344445555")
.insert("00001");
if (ret == false)
{
LOG_ERROR("数据更新失败!");
return -1;
}
else
{
LOG_INFO("数据更新成功!");
}
Json::Value user = ESSearch(client, "test_user")
.append_should_match("phone.keyword", "13344445555")
//.append_must_not_terms("nickname.keyword", {"张三"})
.search();
if (user.empty() || user.isArray() == false)
{
LOG_ERROR("结果为空,或者结果不是数组类型");
return -1;
}
else
{
LOG_INFO("数据检索成功!");
}
int sz = user.size();
LOG_DEBUG("检索结果条目数量:{}", sz);
for (int i = 0; i < sz; i++)
{
LOG_INFO("nickname: {}", user[i]["_source"]["nickname"].asString());
}
ret = ESRemove(client, "test_user").remove("00001");
if (ret == false)
{
LOG_ERROR("删除数据失败");
return -1;
}
else
{
LOG_INFO("数据删除成功!");
}
return 0;
}
makefile
main : main.cc
g++ -std=c++17 $^ -o $@ -lcpr -lelasticlient -lspdlog -lfmt -lgflags -ljsoncpp
运行结果:

cpp-httplib
安装
git clone https://github.com/yhirose/cpp-httplib.git

我们将里面的头文件拷贝在我们的项目中:
mv ./cpp-httplib/httplib.h ../common/
测试使用
main.cc
#include "../common/httplib.h"
int main()
{
// 1. 实例化服务器对象
httplib::Server server;
// 2. 注册回调函数 void(const httplib::Request &, httplib::Response &)
server.Get("/hi", [](const httplib::Request &req, httplib::Response &rsp)
{
std::cout << req.method << std::endl;
std::cout << req.path << std::endl;
for (auto it : req.headers) {
std::cout << it.first << ": " << it.second << std::endl;
}
std::string body = "<html><body><h1>Hello Bite</h1></body></html>";
rsp.set_content(body, "text/html");
rsp.status = 200; });
// server.Post()
// 3. 启动服务器
server.listen("0.0.0.0", 9090);
return 0;
}
makefile
main : main.cc
g++ -std=c++17 $^ -o $@ -lpthread
运行结果:


Websocketpp
安装
sudo apt-get install libboost-dev libboost-system-dev libwebsocketpp-de

安装完毕后,若在 /usr/include 下有了 websocketpp 目录就表示安装成功了

测试使用
main.cc
#include <websocketpp/config/asio_no_tls.hpp>
#include <websocketpp/server.hpp>
// 0. 定义server类型
typedef websocketpp::server<websocketpp::config::asio> server_t;
void onOpen(websocketpp::connection_hdl hdl)
{
std::cout << "websocket长连接建立成功!\n";
}
void onClose(websocketpp::connection_hdl hdl)
{
std::cout << "websocket长连接断开!\n";
}
void onMessage(server_t *server, websocketpp::connection_hdl hdl, server_t::message_ptr msg)
{
// 1. 获取有效消息载荷数据,进行业务处理
std::string body = msg->get_payload();
std::cout << "收到消息:" << body << std::endl;
// 2. 对客户端进行响应
// 获取通信连接
auto conn = server->get_con_from_hdl(hdl);
// 发送数据
conn->send(body + "-Hello!", websocketpp::frame::opcode::value::text);
}
int main()
{
// 1. 实例化服务器对象
server_t server;
// 2. 初始化日志输出 --- 关闭日志输出
server.set_access_channels(websocketpp::log::alevel::none);
// 3. 初始化asio框架
server.init_asio();
// 4. 设置消息处理/连接握手成功/连接关闭回调函数
server.set_open_handler(onOpen);
server.set_close_handler(onClose);
auto msg_hadler = std::bind(onMessage, &server, std::placeholders::_1, std::placeholders::_2);
server.set_message_handler(msg_hadler);
// 5. 启用地址重用
server.set_reuse_addr(true);
// 5. 设置监听端口
server.listen(9090);
// 6. 开始监听
server.start_accept();
// 7. 启动服务器
server.run();
return 0;
}
为了方便测试,我们需要搭建一个客户端,我们直接给出一个页面:
<!DOCTYPE html><html lang="en">
<head>
<meta charset="UTF-8">
<meta http-equiv="X-UA-Compatible" content="IE=edge">
<meta name="viewport" content="width=device-width,
initial-scale=1.0">
<title>Test Websocket</title>
</head>
<body>
<input type="text" id="message">
<button id="submit">提交</button>
<script>
// 创建 websocket 实例
// ws://192.168.51.100:8888
// 类比 http
// ws 表示 websocket 协议
// 192.168.51.100 表示服务器地址
// 8888 表示服务器绑定的端口
let websocket = new WebSocket("ws://43.142.188.181:9090");
// 处理连接打开的回调函数
websocket.onopen = function() {
console.log("连接建立");
}
// 处理收到消息的回调函数
// 控制台打印消息
websocket.onmessage = function(e) {
console.log("收到消息: " + e.data);
}
// 处理连接异常的回调函数
websocket.onerror = function() {
console.log("连接异常");
}
// 处理连接关闭的回调函数
websocket.onclose = function() {
console.log("连接关闭");
}
// 实现点击按钮后, 通过 websocket 实例 向服务器发送请求
let input = document.querySelector('#message');
let button = document.querySelector('#submit');
button.onclick = function() {console.log("发送消息: " + input.value);
websocket.send(input.value);
}
</script>
</body>
</html>
启动我们的服务之后,打开这个页面,按住F12,打开控制台:

对应的服务这边:

redis
介绍
Redis(Remote Dictionary Server)是一个开源的高性能键值对(key-value)数据库。
它通常用作数据结构服务器,因为除了基本的键值存储功能外, Redis 还支持多种类型的数据结构,如字符串(strings)、哈希(hashes)、列表(lists)、集合(sets)、有序集合(sorted sets)以及范围查询、位图、超日志和地理空间索引等。
以下是Redis的一些主要特性:
- 内存中数据库:
Redis将所有数据存储在内存中,这使得读写速度非常快。 - 持久化:尽管
Redis是内存数据库,但它提供了持久化选项,可以将内存中的数
据保存到磁盘上,以防系统故障导致数据丢失。 - 支持多种数据结构:
Redis不仅支持基本的键值对,还支持列表、集合、有序集合
等复杂的数据结构。 - 原子操作:
Redis支持原子操作,这意味着多个操作可以作为一个单独的原子步骤
执行,这对于并发控制非常重要。 - 发布/订阅功能:
Redis支持发布订阅模式,允许多个客户端订阅消息,当消息发
布时,所有订阅者都会收到消息。 - 高可用性:通过
Redis哨兵(Sentinel)和Redis集群,Redis可以提供高可用性
和自动故障转移。 - 复制:
Redis支持主从复制,可以提高数据的可用性和读写性能。 - 事务:
Redis提供了事务功能,可以保证一系列操作的原子性执行。 - Lua 脚本:
Redis支持使用Lua脚本进行复杂的数据处理,可以在服务器端执行
复杂的逻辑。 - 客户端库: Redis 拥有丰富的客户端库,支持多种编程语言,如
Python、Ruby、
Java、C#等。 - 性能监控:
Redis提供了多种监控工具和命令,可以帮助开发者监控和优化性能。 - 易于使用:
Redis有一个简单的配置文件和命令行界面,使得设置和使用变得容易。
Redis广泛用于缓存、会话存储、消息队列、排行榜、实时分析等领域。由于其高性
能和灵活性,Redis成为了现代应用程序中非常流行的数据存储解决方案之一。
安装
sudo apt install -y redis

支持远程连接
修改 /etc/redis/redis.conf
• 修改 bind 127.0.0.1 为 bind 0.0.0.0
• 修改 protected-mode yes 为 protected-mode no


设置开机重启
sudo systemctl enable redis-server

安装 redis-plus-plus
C++ 操作 redis 的库有很多. 咱们此处使用 redis-plus-plus.
这个库的功能强大, 使用简单.
Github 地址: https://github.com/sewenew/redis-plus-plus
安装 hiredis
redis-plus-plus 是基于 hiredis 实现的.
hiredis 是一个 C 语言实现的 redis 客户端.
因此需要先安装 hiredis. 直接使用包管理器安装即可
sudo apt install libhiredis-dev

下载 redis-plus-plus 源码
git clone https://github.com/sewenew/redis-plus-plus.git

cd redis-plus-plus
mkdir build
cd build
cmake -DCMAKE_INSTALL_PREFIX=/usr ..
make
make install # 这一步操作需要管理员权限. 如果是非 root 用户, 使用
sudo make install 执行



redis基础操作实践
main.cc
#include <sw/redis++/redis.h>
#include <gflags/gflags.h>
#include <iostream>
#include <thread>
DEFINE_string(ip, "127.0.0.1", "这是服务器的IP地址,格式:127.0.0.1");
DEFINE_int32(port, 6379, "这是服务器的端口, 格式: 8080");
DEFINE_int32(db, 0, "库的编号:默认0号");
DEFINE_bool(keep_alive, true, "是否进行长连接保活");
void print(sw::redis::Redis &client)
{
auto user1 = client.get("会话ID1");
if (user1)
std::cout << *user1 << std::endl;
auto user2 = client.get("会话ID2");
if (user2)
std::cout << *user2 << std::endl;
auto user3 = client.get("会话ID3");
if (user3)
std::cout << *user3 << std::endl;
auto user4 = client.get("会话ID4");
if (user4)
std::cout << *user4 << std::endl;
auto user5 = client.get("会话ID5");
if (user5)
std::cout << *user5 << std::endl;
}
void add_string(sw::redis::Redis &client)
{
client.set("会话ID1", "用户ID1");
client.set("会话ID2", "用户ID2");
client.set("会话ID3", "用户ID3");
client.set("会话ID4", "用户ID4");
client.set("会话ID5", "用户ID5");
client.del("会话ID3");
client.set("会话ID5", "用户ID555"); // 数据已存在则进行修改,不存在则新增
print(client);
}
void expired_test(sw::redis::Redis &client)
{
// 这次的新增,数据其实已经有了,因此本次是修改
// 不仅仅修改了val,而且还给键值对新增了过期时间
client.set("会话ID1", "用户ID1111", std::chrono::milliseconds(1000));
print(client);
std::cout << "------------休眠2s-----------\n";
std::this_thread::sleep_for(std::chrono::seconds(2));
print(client);
}
void list_test(sw::redis::Redis &client)
{
client.rpush("群聊1", "成员1");
client.rpush("群聊1", "成员2");
client.rpush("群聊1", "成员3");
client.rpush("群聊1", "成员4");
client.rpush("群聊1", "成员5");
std::vector<std::string> users;
client.lrange("群聊1", 0, -1, std::back_inserter(users));
for (auto user : users)
{
std::cout << user << std::endl;
}
}
int main(int argc, char *argv[])
{
google::ParseCommandLineFlags(&argc, &argv, true);
// 功能接口演示中:
// 1. 构造连接选项,实例化Redis对象,连接服务器
sw::redis::ConnectionOptions opts;
opts.host = FLAGS_ip;
opts.port = FLAGS_port;
opts.db = FLAGS_db;
opts.keep_alive = FLAGS_keep_alive;
sw::redis::Redis client(opts);
// 2. 添加字符串键值对,删除字符串键值对,获取字符串键值对
add_string(client);
// 3. 实践控制数据有效时间的操作
expired_test(client);
// 4. 列表的操作,主要实现数据的右插,左获取
std::cout << "--------------------------\n";
list_test(client);
return 0;
}
makefile
main : main.cc
g++ -std=c++17 $^ -o $@ -lhiredis -lredis++ -lgflags
运行结果:

更多推荐
所有评论(0)