首页 文章 精选 留言 我的

精选列表

搜索[时序预测],共10000篇文章
优秀的个人博客,低调大师

时序数据库Influx-IOx源码学习五(创建数据库)

欢迎关注公众号: 上篇介绍到:InfluxDB-IOx的Run命令启动过程,详情见:https://my.oschina.net/u/3374539/blog/5021654 这章记录一下Database create命令的执行过程。 在第三章命令行中介绍了,所有的子命令都有一个独立的参数或配置称为subcommand。 enum Command { Convert { // 省略 ...}, Meta {// 省略 ...}, Database(commands::database::Config), Run(Box<commands::run::Config>), Stats(commands::stats::Config), Server(commands::server::Config), Writer(commands::writer::Config), Operation(commands::operations::Config), } 这章我们打开看一眼commands::database下的config包含了什么。 pub struct Config { #[structopt(subcommand)] command: Command, } //见名知意,基本猜测一下就行了,慢慢使用到再回来看 enum Command { Create(Create), List(List), Get(Get), Write(Write), Query(Query), Chunk(chunk::Config), Partition(partition::Config), } 先来看一下create命令的执行。 Command::Create(command) => { //创建一个grpc的client let mut client = management::Client::new(connection); //设置基本的配置项 let rules = DatabaseRules { //数据库名字 name: command.name, //内存的各种配置,包含缓存大小,时间等等 lifecycle_rules: Some(LifecycleRules { //省略。。 }), //设置分区的策略 partition_template: Some(PartitionTemplate { //省略。。 }), //其它都填充default ..Default::default() }; //使用配置信息创建数据库,这里是生成了一个CreateDatabaseRequest去调用了远程服务器的方法 client.create_database(rules).await?; println!("Ok"); } 在上一章中提到了grpc的启动,这里就涉及到了之前提到的grpc的框架tonic,在tonic中使用#[tonic::async_trait]了标记一个服务器端的实现开始。我在ide中搜索,可以在src/influxdb_ioxd/rpc/management.rs:50行中找到ManagementService相关的实现。 有关tonic更多的资料请阅读:https://github.com/hyperium/tonic #[tonic::async_trait] impl<M> management_service_server::ManagementService for ManagementService<M> where M: ConnectionManager + Send + Sync + Debug + 'static, { //省略其它方法。。。 async fn create_database( &self, //这里就是接收CreateDatabaseRequest的请求 request: Request<CreateDatabaseRequest>, ) -> Result<Response<CreateDatabaseResponse>, Status> { //对数据进行一下校验,然后获得在上面配置的rules规则 let rules: DatabaseRules = request .into_inner() .rules .ok_or_else(|| FieldViolation::required("")) .and_then(TryInto::try_into) .map_err(|e| e.scope("rules"))?; //这里就是在第三章中提到的server_id,如果没配置就会报错了 let server_id = match self.server.require_id().ok() { Some(id) => id, None => return Err(NotFound::default().into()), }; //这里就是真正的去创建,在下面继续跟踪 match self.server.create_database(rules, server_id).await { Ok(_) => Ok(Response::new(CreateDatabaseResponse {})), Err(Error::DatabaseAlreadyExists { db_name }) => { return Err(AlreadyExists { resource_type: "database".to_string(), resource_name: db_name, ..Default::default() } .into()) } Err(e) => Err(default_server_error_handler(e)), } } } 接下来要继续查看数据库真正的被创建出来,我读到这里存在一个问题,文件格式是什么样子的? pub async fn create_database(&self, rules: DatabaseRules, server_id: NonZeroU32) -> Result<()> { //检查server_id self.require_id()?; //把数据库名字存储到内存中,最终保存到一个btreemap中 let db_reservation = self.config.create_db(rules)?; //对数据进行持久化保存 self.persist_database_rules(db_reservation.rules().clone()) .await?; //启动数据库后台线程,在内存中写入数据库状态 db_reservation.commit(server_id, Arc::clone(&self.store), Arc::clone(&self.exec)); Ok(()) } 来解答上面的疑问,文件是怎样持久化、格式是什么样子的。 pub async fn persist_database_rules<'a>(&self, rules: DatabaseRules) -> Result<()> { //生成一个新的数据库路径 let location = object_store_path_for_database_config(&self.root_path()?, &rules.name); //序列化DatabaseRules这个pb到byte流 let mut data = BytesMut::new(); rules.encode(&mut data).context(ErrorSerializing)?; let len = data.len(); let stream_data = std::io::Result::Ok(data.freeze()); //将pb的内容进行存储 self.store .put( &location, futures::stream::once(async move { stream_data }), Some(len), ) .await .context(StoreError)?; Ok(()) } 这里调用了rules.encode()转换到pb的格式,这里是rust语言的一个方法,实现了From特性的,就得到了一个into的方法,如:impl From<DatabaseRules> for management::DatabaseRules. 到这里数据库的一个描述文件rules.pb就被写入到磁盘中了,路径是启动命令中指定的--data-dir参数路径 + --writer-id + 数据库名字。 例如,我的启动和创建命令为: ./influxdb_iox run --writer-id 1 --object-store file --data-dir ~/influxtest/ ./influxdb_iox database create test 那么得到的路径就为:~/influxtest/1/test/rules.pb. 之后可以运行一个pb的脚本来反查rules.pb中的数据内容,如下: $ ./scripts/prototxt decode influxdata.iox.management.v1.DatabaseRules \ < ~/influxtest/1/test/rules.pb influxdata/iox/management/v1/service.proto:6:1: warning: Import google/protobuf/field_mask.proto is unused. name: "test" partition_template { parts { time: "%Y-%m-%d %H:00:00" } } lifecycle_rules { mutable_linger_seconds: 300 mutable_size_threshold: 10485760 buffer_size_soft: 52428800 buffer_size_hard: 104857600 sort_order { order: ORDER_ASC created_at_time { } } } 看到这里已经知道整个生成过程及文件内容。 祝玩儿的开心。

优秀的个人博客,低调大师

时序数据库Influx-IOx源码学习四(Run命令的执行)

欢迎关注公众号: 上篇介绍到:InfluxDB-IOx的命令行及配置,详情见:https://my.oschina.net/u/3374539/blog/5017858 这章记录一下Run命令的执行过程。 //根据用户在命令行配置的num_threads参数 //来选择创建一个多线程的模型,还是current_thread的模型 //后面有时间深入研究tokio的时候再来分析有什么异同 let tokio_runtime = get_runtime(config.num_threads)?; //block_on会让线程一直等待方法里的future执行完成 //这是让闭包中的方法占有了io driver 和 timer context tokio_runtime.block_on(async move { let host = config.host; match config.command { // 省略其它command ... Command::Run(config) => { //具体去子类型里执行,然后await一个结果 if let Err(e) = commands::run::command(logging_level, *config).await { eprintln!("Server command failed: {}", e); std::process::exit(ReturnCode::Failure as _) } } } }); 在influxdb_ioxd::main方法中,忽略一些不太需要重点关注的,分别是初始化log的管理、PanicsTracing、CancellationToken等。 //初始化对象存储 let object_store = ObjectStore::try_from(&config)?; //可以看到,目前已经支持了 //1.内存(在container环境运行时候使用) //2.Google //3.S3 //4.Azure //5.File 本地文件,方便开发者调试运行在云上时候的文件变化 fn try_from(config: &Config) -> Result<Self, Self::Error> { match config.object_store { Some(ObjStoreOpt::Memory) | None => { //创建一个btreemap用来缓存或者搜索 Ok(Self::new_in_memory(object_store::memory::InMemory::new())) } Some(ObjStoreOpt::Google) => { // 省略 } Some(ObjStoreOpt::S3) => { // 省略 } Some(ObjStoreOpt::Azure) => { // 省略 } Some(ObjStoreOpt::File) => match config.database_directory.as_ref() { Some(db_dir) => { //去递归创建这个配置路径中的文件夹 //context也是使用的snafu来处理错误的 fs::create_dir_all(db_dir) .context(CreatingDatabaseDirectory { path: db_dir })?; //都创建完成,并且没出错误,把路径保存起来 Ok(Self::new_file(object_store::disk::File::new(&db_dir))) } // 如果database_directory这个参数没有配置的时候 //使用snafu这个crate来返回一个错误 None => MissingObjectStoreConfig { object_store: ObjStoreOpt::File, missing: "data-dir", } .fail(), }, } } 关于错误处理的代码: #[snafu(display("Unable to create database directory {:?}: {}", path, source))] CreatingDatabaseDirectory { path: PathBuf, source: std::io::Error, }, #[snafu(display( "Specified {} for the object store, required configuration missing for {}", object_store, missing ))] MissingObjectStoreConfig { object_store: ObjStoreOpt, missing: String, }, 我们来测试一下错误的场景,来看看是否符合代码的预期。 // 不传入路径 cargo run run --object-store file Finished dev [unoptimized + debuginfo] target(s) in 0.42s Running `./influxdb_iox run --object-store file` Apr 15 13:38:34.352 INFO influxdb_iox::influxdb_ioxd: Using File for object storage Server command failed: Run: Specified File for the object store, required configuration missing for data-dir //传入一个创建不了的路径 cargo run run --object-store file --data-dir /root/1/1 Finished dev [unoptimized + debuginfo] target(s) in 0.47s Running `./influxdb_iox run --object-store file --data-dir /root/1/1` Apr 15 13:45:26.664 INFO influxdb_iox::influxdb_ioxd: Using File for object storage Server command failed: Run: Unable to create database directory "/root/1/1": Read-only file system (os error 30) 可以看到是符合预期的,bingo //创建一个空的结构体 let connection_manager = ConnectionManager {}; //创建AppServer结构体用来保存基本的信息 //server_config里就是保存的对象存储的信息及线程配置 //如果num_worker_threads没有填写,默认就使用cpu数量 let app_server = Arc::new(AppServer::new(connection_manager, server_config)); //不设置这个writer_id能启动,但是不能做任何操作 if let Some(id) = config.writer_id { //compare and set 一个非0的数值,错误就打印一个指定的panic app_server.set_id(id).expect("writer id already set"); //校验所有的配置 if let Err(e) = app_server.load_database_configs().await { error!( "unable to load database configurations from object storage: {}", e ) } } else { warn!("server ID not set. ID must be set via the INFLUXDB_IOX_ID config or API before writing or querying data."); } 接下来进入load_database_configs方法看看, let list_result = self .store //把write_id和配置的文件路径组合一下,作为一个目录 //遍历文件夹中的所有东西,用一个BTreeSet存所有子文件夹 //用Vec存下所有的文件信息,包括路径、修改时间、大小等 .list_with_delimiter(&self.root_path()?) .await .context(StoreError)?; //拿到配置的server的write_id let server_id = self.require_id()?; let handles: Vec<_> = list_result //配置的文件夹下的所有文件夹 .common_prefixes .into_iter() //全部进行map转换 .map(|mut path| { let store = Arc::clone(&self.store); let config = Arc::clone(&self.config); let exec = Arc::clone(&self.exec); //先找database的相关信息文件,名字叫rules.pb path.set_file_name(DB_RULES_FILE_NAME); //感觉是需要io来读取文件内容,所以开一个异步 tokio::task::spawn(async move { let mut res = get_store_bytes(&path, &store).await; //省略错误处理。。 let res = res.unwrap().freeze(); //解析文件内容,根据文件名可以看出是个pb文件。 match DatabaseRules::decode(res) { Err(e) => { //省略错误。。 } //根据解析出来的文件内容,在内存中恢复回来db的相关信息 Ok(rules) => match config.create_db(rules) { Err(e) => error!("error adding database to config: {}", e), //提交一个后台任务,用来不断的检测chunks的状态 //比如达到了某个大小,然后写入到存储等 Ok(handle) => handle.commit(server_id, store, exec), }, } }) }) .collect(); //等待所有任务完成 futures::future::join_all(handles).await; 这里就启动完成了一个基本的服务,创建了存储路径、初始化数据库的基本配置、启动了一个用来刷盘、整理chunk的后台任务。 接下来就是启动连接相关的了。 //从启动命令行中读取grpc的地址 let grpc_bind_addr = config.grpc_bind_address; //绑定这个地址 let socket = tokio::net::TcpListener::bind(grpc_bind_addr) .await .context(StartListeningGrpc { grpc_bind_addr })?; //真正的协议启动 let grpc_server = rpc::serve(socket, Arc::clone(&app_server), frontend_shutdown.clone()).fuse(); //同样的启动http相关的服务,使用的hyper库 let bind_addr = config.http_bind_address; let addr = AddrIncoming::bind(&bind_addr).context(StartListeningHttp { bind_addr })?; let http_server = http::serve(addr, Arc::clone(&app_server), frontend_shutdown.clone()).fuse(); //省略后面的停止流程。。。 然后看grpc的启动的服务 //启动起来健康检查的服务 let stream = TcpListenerStream::new(socket); let (mut health_reporter, health_service) = tonic_health::server::health_reporter(); //标识相对应的服务已经是可以提供服务的状态了 let services = [ generated_types::STORAGE_SERVICE, generated_types::IOX_TESTING_SERVICE, generated_types::ARROW_SERVICE, ]; for service in &services { health_reporter .set_service_status(service, tonic_health::ServingStatus::Serving) .await; } //增加一堆使用grpc的服务,并启动起来 tonic::transport::Server::builder() .add_service(health_service) .add_service(testing::make_server()) .add_service(storage::make_server(Arc::clone(&server))) .add_service(flight::make_server(Arc::clone(&server))) .add_service(write::make_server(Arc::clone(&server))) .add_service(management::make_server(Arc::clone(&server))) .add_service(operations::make_server(server)) .serve_with_incoming_shutdown(stream, shutdown.cancelled()) .await 然后是http相关的启动 pub async fn serve<M>( addr: AddrIncoming, server: Arc<AppServer<M>>, shutdown: CancellationToken, ) -> Result<(), hyper::Error> where M: ConnectionManager + Send + Sync + Debug + 'static, { //初始化路由相关的信息 let router = router(server); let service = RouterService::new(router).unwrap(); //启动服务 hyper::Server::builder(addr) .serve(service) .with_graceful_shutdown(shutdown.cancelled()) .await } 顺便看一下都提供了哪些地址可以被访问的: Router::builder() .data(server) //写了一个拦截,打印请求参数和返回结果 .middleware(Middleware::pre(|req| async move { debug!(request = ?req, "Processing request"); Ok(req) })) .middleware(Middleware::post(|res| async move { debug!(response = ?res, "Successfully processed request"); Ok(res) })) // this endpoint is for API backward compatibility with InfluxDB 2.x .post("/api/v2/write", write::<M>) .get("/health", health) .get("/metrics", handle_metrics) .get("/iox/api/v1/databases/:name/query", query::<M>) .get("/iox/api/v1/databases/:name/wal/meta", get_wal_meta::<M>) .get("/api/v1/partitions", list_partitions::<M>) .post("/api/v1/snapshot", snapshot_partition::<M>) //错误的时候调用的处理拦截 .err_handler_with_info(error_handler) .build() .unwrap() 做一个/health的测试: curl localhost:8080/health OK% 可以看到成功返回了值。 到这里基本启动就完成了,后面再用到的时候会继续对启动里的细节做研究,比如Panics,Log等等吧,欢迎持续关注。 祝玩儿的开心

资源下载

更多资源
腾讯云软件源

腾讯云软件源

为解决软件依赖安装时官方源访问速度慢的问题,腾讯云为一些软件搭建了缓存服务。您可以通过使用腾讯云软件源站来提升依赖包的安装速度。为了方便用户自由搭建服务架构,目前腾讯云软件源站支持公网访问和内网访问。

Nacos

Nacos

Nacos /nɑ:kəʊs/ 是 Dynamic Naming and Configuration Service 的首字母简称,一个易于构建 AI Agent 应用的动态服务发现、配置管理和AI智能体管理平台。Nacos 致力于帮助您发现、配置和管理微服务及AI智能体应用。Nacos 提供了一组简单易用的特性集,帮助您快速实现动态服务发现、服务配置、服务元数据、流量管理。Nacos 帮助您更敏捷和容易地构建、交付和管理微服务平台。

Rocky Linux

Rocky Linux

Rocky Linux(中文名:洛基)是由Gregory Kurtzer于2020年12月发起的企业级Linux发行版,作为CentOS稳定版停止维护后与RHEL(Red Hat Enterprise Linux)完全兼容的开源替代方案,由社区拥有并管理,支持x86_64、aarch64等架构。其通过重新编译RHEL源代码提供长期稳定性,采用模块化包装和SELinux安全架构,默认包含GNOME桌面环境及XFS文件系统,支持十年生命周期更新。

WebStorm

WebStorm

WebStorm 是jetbrains公司旗下一款JavaScript 开发工具。目前已经被广大中国JS开发者誉为“Web前端开发神器”、“最强大的HTML5编辑器”、“最智能的JavaScript IDE”等。与IntelliJ IDEA同源,继承了IntelliJ IDEA强大的JS部分的功能。

用户登录
用户注册