首页 文章 精选 留言 我的

精选列表

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

腾讯 APIJSON 插件 apijson-influxdb 开源,支持物联网时序数据库

腾讯 APIJSON 是一种专为 API 而生的 JSON 网络传输协议 以及 基于这套协议实现的 ORM 库。为各种增删改查提供了完全自动化的万能 API,零代码实时满足千变万化的各种新增和变更需求。 能大幅降低开发和沟通成本,简化开发流程,缩短开发周期。适合中小型前后端分离的项目。 自 2016 年开源 7 年来发展迅速,目前 16K+ Star 位居 400W Java 开源项目前 100。 国内 腾讯、华为、阿里巴巴、字节跳动、美团、拼多多、百度、京东、网易、快手、圆通 等 和 国外 Google, Apple, Microsoft, Amazon, Tesla, Meta(FB), Paypal 等数百个知名大厂员工 Star, 也有 腾讯、华为、字节跳动、Microsoft、Zoom、知乎 等 工程师/专家/架构师 提了 PR/Issue, 还被 腾讯、华为、百度、SHEIN、快手、中兴、传音、圆通、美图 等各大知名厂商用于各类项目。 apijson-influxdb 腾讯APIJSON6.1.0+ 的 InfluxDB 数据库插件,可通过 Maven, Gradle 等远程依赖。 An InfluxDB plugin for TencentAPIJSON6.1.0+ 添加依赖 Add Dependency Maven 1. 在 pom.xml 中添加 JitPack 仓库 <repositories> <repository> <id>jitpack.io</id> <url>https://jitpack.io</url> </repository> </repositories> 2. 在 pom.xml 中添加 apijson-influxdb 依赖 <dependency> <groupId>com.github.APIJSON</groupId> <artifactId>apijson-influxdb</artifactId> <version>LATEST</version> </dependency> 使用 Usage 在你项目继承 AbstractSQLExecutor 的子类重写方法 execute Override execute in your SQLExecutor extends AbstractSQLExecutor @Override public JSONObject execute(@NotNull SQLConfig<Long> config, boolean unknownType) throws Exception { if (config.isInfluxDB()) { return InfluxdbUtil.execute(config, null, unknownType); } return super.execute(config, unknownType); } 在你项目继承 AbstractSQLConfig 的子类重写方法 execute Override execute in your SQLConfig extends AbstractSQLConfig @Override public String getSchema() { return InfluxDBUtil.getSchema(super.getSchema(), DEFAULT_SCHEMA, isInfluxDB()); } @Override public String getSQLSchema() { return InfluxDBUtil.getSQLSchema(super.getSQLSchema(), isInfluxDB()); } 有问题可以去 Tencent/APIJSON 提 issue 点右上角 ⭐Star 支持一下,谢谢 ^_^ https://github.com/APIJSON/apijson-influxdb

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

时序数据库Influx-IOx源码学习六-1(数据写入之分区)

欢迎关注公众号: 上一章说到如何创建一个数据库,并且数据库的描述信息是如何保存的。详情见:https://my.oschina.net/u/3374539/blog/5025128 这一章记录一下,数据是如何写入并保存的,具体会分为两篇来写: 一篇介绍分区是如何完成的 一篇介绍具体的写入 说到数据写入,必然是需要能够连接到服务器。IOx项目为提供了多种方式可以于服务器进行交互,分别是Grpc和Http基于这两种通信方式,又扩展支持了influxdb2_client以及influxdb_iox_client。 基于influxdb_iox_client我写了一个数据写入及查询的示例来观测接口是如何组织的,代码如下: #[tokio::main] async fn main() { { let connection = Builder::default() .build("http://127.0.0.1:8081") .await .unwrap(); write::Client::new(connection) .write("a", r#"myMeasurement,tag1=value1,tag2=value2 fieldKey="123" 1556813561098000000"#) .await .expect("failed to write data"); } let connection = Builder::default() .build("http://127.0.0.1:8081") .await .unwrap(); let mut query = flight::Client::new(connection) .perform_query("a", "select * from myMeasurement") .await .expect("query request should work"); let mut batches = vec![]; while let Some(data) = query.next().await.expect("valid batches") { batches.push(data); } let format1 = format::QueryOutputFormat::Pretty; println!("{}", format1.format(&batches).unwrap()); } +------------+--------+--------+-------------------------+ | fieldKey | tag1 | tag2 | time | +------------+--------+--------+-------------------------+ | 123 | value1 | value2 | 2019-05-02 16:12:41.098 | | 123 | value1 | value2 | 2019-05-02 16:12:41.098 | | fieldValue | value1 | value2 | 2019-05-02 16:12:41.098 | | fieldValue | value1 | value2 | 2019-05-02 16:12:41.098 | | 123 | value1 | value2 | 2019-05-02 16:12:41.098 | +------------+--------+--------+-------------------------+ 因为我多运行了几次,所以能看到数据被重复插入了。 这里还需要说一下的是写入的语句格式可以参见: [LineProtocol] https://docs.influxdata.com/influxdb/v2.0/reference/syntax/line-protocol/#data-types-and-format write::Client中的write方法生成了一个WriteRequest结构,并使用RPC调用远程的write方法。打开src/influxdb_ioxd/rpc/write.rs : 22行可以看到方法的具体实现。 async fn write( &self, request: tonic::Request<WriteRequest>, ) -> Result<tonic::Response<WriteResponse>, tonic::Status> { let request = request.into_inner(); //得到上面在客户端中写入的数据库名字,在上面的例子中传入的"a" let db_name = request.db_name; //这里得到了写入的LineProtocol let lp_data = request.lp_data; let lp_chars = lp_data.len(); //解析LineProtocol的内容 //示例中的lp会被解析为: //measurement: "myMeasurement" //tag_set: [("tag1", "value1"), ("tag2", "value2")] //field_set: [("fieldKey", "123")] //timestamp: 1556813561098000000 let lines = parse_lines(&lp_data) .collect::<Result<Vec<_>, influxdb_line_protocol::Error>>() .map_err(|e| FieldViolation { field: "lp_data".into(), description: format!("Invalid Line Protocol: {}", e), })?; let lp_line_count = lines.len(); debug!(%db_name, %lp_chars, lp_line_count, "Writing lines into database"); //对数据进行保存 self.server .write_lines(&db_name, &lines) .await .map_err(default_server_error_handler)?; //返回成功 let lines_written = lp_line_count as u64; Ok(Response::new(WriteResponse { lines_written })) } 继续看self.server.write_lines的执行: pub async fn write_lines(&self, db_name: &str, lines: &[ParsedLine<'_>]) -> Result<()> { self.require_id()?; //验证一下名字,然后拿到之前创建数据库时候在内存中存储的相关信息 let db_name = DatabaseName::new(db_name).context(InvalidDatabaseName)?; let db = self .config .db(&db_name) .context(DatabaseNotFound { db_name: &*db_name })?; //这里就开始执行分片相关的策略 let (sharded_entries, shards) = { //读取创建数据库时候配置的分片策略 let rules = db.rules.read(); let shard_config = &rules.shard_config; //根据数据和shard策略,把逐个数据对应的分区找到 //写入到一个List<分区标识,List<数据>>这样的结构中 //具体的结构信息后面看 let sharded_entries = lines_to_sharded_entries(lines, shard_config.as_ref(), &*rules) .context(LineConversion)?; //再把所有分区的配置返回给调用者 let shards = shard_config .as_ref() .map(|cfg| Arc::clone(&cfg.shards)) .unwrap_or_default(); (sharded_entries, shards) }; //根据上面返回的集合进行map方法遍历,写到每个分区中 futures_util::future::try_join_all( sharded_entries .into_iter() .map(|e| self.write_sharded_entry(&db_name, &db, Arc::clone(&shards), e)), ) .await?; Ok(()) } 这里描述了写入一条数据的主逻辑:数据写入的时候,先把数据划分到具体的分区里(使用List结构存储下所有的分区对应的数据),然后并行的进行数据写入 接下来看,数据是如何进行分区的: pub fn lines_to_sharded_entries( lines: &[ParsedLine<'_>], sharder: Option<&impl Sharder>, partitioner: &impl Partitioner, ) -> Result<Vec<ShardedEntry>> { let default_time = Utc::now(); let mut sharded_lines = BTreeMap::new(); //对所有要插入的数据进行遍历 for line in lines { //先找到符合哪个shard let shard_id = match &sharder { Some(s) => Some(s.shard(line).context(GeneratingShardId)?), None => None, }; //再判断属于哪个分区 let partition_key = partitioner .partition_key(line, &default_time) .context(GeneratingPartitionKey)?; let table = line.series.measurement.as_str(); //最后存储到一个map中 //shard-> partition -> table -> List<data> 的映射关系 sharded_lines .entry(shard_id) .or_insert_with(BTreeMap::new) .entry(partition_key) .or_insert_with(BTreeMap::new) .entry(table) .or_insert_with(Vec::new) .push(line); } let default_time = Utc::now(); //最后遍历这个map 转换到之前提到的List结构中 let sharded_entries = sharded_lines .into_iter() .map(|(shard_id, partitions)| build_sharded_entry(shard_id, partitions, &default_time)) .collect::<Result<Vec<_>>>()?; Ok(sharded_entries) } 这里理解shard的概念就是一个或者一组机器,称为一个shard,他们负责真正的存储数据。 partition理解为一个个文件夹,在shard上具体的存储路径。 这里看一下是怎样完成shard的划分的: impl Sharder for ShardConfig { fn shard(&self, line: &ParsedLine<'_>) -> Result<ShardId, Error> { if let Some(specific_targets) = &self.specific_targets { //如果对数据进行匹配,如果符合规则就返回,可以采用当前的shard //官方的代码中只实现了根据表名进行shard的策略 //这个配置似乎只能通过grpc来进行设置,这样好处可能是将来有个什么管理界面能动态修改 if specific_targets.matcher.match_line(line) { return Ok(specific_targets.shard); } } //如果没有配置就使用hash的方式 //对整条数据进行hash,然后比较机器的hash,找到合适的节点 //如果没找到,就放在hashring的第一个节点 //hash算法见后面 if let Some(hash_ring) = &self.hash_ring { return hash_ring .shards .find(LineHasher { line, hash_ring }) .context(NoShardsDefined); } NoShardingRuleMatches { line: line.to_string(), } .fail() } } //具体的Hash算法,如果全配置的话分的就会特别散,几乎不同测点都放到了不同的地方 impl<'a, 'b, 'c> Hash for LineHasher<'a, 'b, 'c> { fn hash<H: Hasher>(&self, state: &mut H) { //如果配置了使用table名字就在hash中加入tablename if self.hash_ring.table_name { self.line.series.measurement.hash(state); } //然后按照配置的列的值进行hash for column in &self.hash_ring.columns { if let Some(tag_value) = self.line.tag_value(column) { tag_value.hash(state); } else if let Some(field_value) = self.line.field_value(column) { field_value.to_string().hash(state);t } state.write_u8(0); // column separator } } } 接下来看默认的partition分区方式: impl Partitioner for PartitionTemplate { fn partition_key(&self, line: &ParsedLine<'_>, default_time: &DateTime<Utc>) -> Result<String> { let parts: Vec<_> = self .parts .iter() //匹配分区策略,或者是单一的,或者是复合的 //目前支持基于表、值、时间 //其余还会支持正则表达式和strftime模式 .map(|p| match p { TemplatePart::Table => line.series.measurement.to_string(), TemplatePart::Column(column) => match line.tag_value(&column) { Some(v) => format!("{}_{}", column, v), None => match line.field_value(&column) { Some(v) => format!("{}_{}", column, v), None => "".to_string(), }, }, TemplatePart::TimeFormat(format) => match line.timestamp { Some(t) => Utc.timestamp_nanos(t).format(&format).to_string(), None => default_time.format(&format).to_string(), }, _ => unimplemented!(), }) .collect(); //最后返回一个组合文件名,或者是 a-b-c 或者是一个单一的值 Ok(parts.join("-")) } } 到这里分区的工作就完成了,下一篇继续分析是怎样写入的。 祝玩儿的开心

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

时序数据库Influx-IOx源码学习三(命令行及配置)

欢迎关注公众号: 上篇介绍到:InfluxDB-IOx的环境搭建,详情见:https://my.oschina.net/u/3374539/blog/5016798 本章开始,讲解启动的主流程! 打开src/main.rs文件可以找到下面的代码 fn main() -> Result<(), std::io::Error> { // load all environment variables from .env before doing anything load_dotenv(); let config = Config::from_args(); println!("{:?}", config); //省略 ..... Ok(()) } 在main方法中映入眼帘的第一行就是load_dotenv()方法,然后是Config::from_args()接下来就分别跟踪这两个方法,看明白是怎么工作的。 加载配置文件 在README文件中,我们可以看到这样一行: Should you desire specifying config via a file, you can do so using a .env formatted file in the working directory. You can use the provided example as a template if you want: cp docs/env.example .env 意思就是这个工程使用的配置文件,名字是.env。了解这个特殊的名字之后,我们看代码src/main.rs:276: fn load_dotenv() { //调用dotenv方法,并对其返回值进行判断 match dotenv() { //如果返回成功,程序什么都不做,继续执行。 Ok(_) => {} //返回的是错误,那么判断一下是否为'未找到'错误, //如果是未找到,那么就什么都不做(也就是有默认值填充) Err(dotenv::Error::Io(err)) if err.kind() == std::io::ErrorKind::NotFound => { } //这里就是真真正正必须要处理的错误了,直接退出程序 Err(e) => { eprintln!("FATAL Error loading config from: {}", e); eprintln!("Aborting"); std::process::exit(1); } }; } 然后跟踪dotenv()方法看看如何执行(这里就进入了dotenv这个crate了): 为了方便写,我就直接把所有调用,从上到下的顺序全都写出来了 //返回一个PathBuf的Result,之后再看这个Result pub fn dotenv() -> Result<PathBuf> { //new一个Finder结构并调用find方法 //?代表错误的时候直接抛出错误 let (path, iter) = Finder::new().find()?; //返回一个自定义的Iter结构,并调用load方法 iter.load()?; //成功返回 Ok(path) } //创建一个Finder结构体,filename使用`.env`填充 pub fn new() -> Self { Finder { filename: Path::new(".env"), } } //返回一个元组,多个返回值,(路径,文件读取相关记录) pub fn find(self) -> Result<(PathBuf, Iter<File>)> { //使用标准库中的current_dir()方法得到当前的路径 //出错就返回Error::Io错误,正常就调用find方法 let path = find(&env::current_dir().map_err(Error::Io)?, self.filename)?; //如果找到了.env文件就打开,打开错误就返回Error::Io错误 let file = File::open(&path).map_err(Error::Io)?; //使用打开的文件创建一个Iter的结构 let iter = Iter::new(file); //返回 Ok((path, iter)) } //递归查找.env文件 pub fn find(directory: &Path, filename: &Path) -> Result<PathBuf> { //拼装一个全路径 let candidate = directory.join(filename); //尝试打开这个文件 match fs::metadata(&candidate) { //成功打开了,说明找到了.env文件,就返回成功 //但我有个疑问文件内容为啥不校验一下呢? Ok(metadata) => if metadata.is_file() { return Ok(candidate); }, //除了没找到文件的错误之外,其它错误都直接返回异常 Err(error) => { if error.kind() != io::ErrorKind::NotFound { return Err(Error::Io(error)); } } } //没找到的时候,就返回到父级文件夹里,继续找,一直到根文件夹 if let Some(parent) = directory.parent() { find(parent, filename) } else { //一直到根文件夹,还没找到就返回一个NotFound的IO错误, //这个在上面的代码中提到,这个错误会被忽略 Err(Error::Io(io::Error::new(io::ErrorKind::NotFound, "path not found"))) } } //对应的iter.load()?;方法实现 pub fn load(self) -> Result<()> { //可以使用for是因为实现了Iterator 这个trait for item in self { //获取读取出来的一行一行的配置项 let (key, value) = item?; //验证key没有什么问题,就放到env中 if env::var(&key).is_err() { env::set_var(&key, value); } } Ok(()) } // 为了能够for循环,实现的Iterator impl<R: Read> Iterator for Iter<R> { type Item = Result<(String, String)>; fn next(&mut self) -> Option<Self::Item> { loop { //一行一行的读取文件内容 let line = match self.lines.next() { Some(Ok(line)) => line, Some(Err(err)) => return Some(Err(Error::Io(err))), None => return None, }; //解析配置项目,这里就不在深入跟了 match parse::parse_line(&line, &mut self.substitution_data) { Ok(Some(result)) => return Some(Ok(result)), Ok(None) => {} Err(err) => return Some(Err(err)), } } } } 研究这里的时候,我发现了一个比较好玩儿的东西就是返回值的Result<PathBuf>。标准库的定义中,Result是有两个值,分别是<T,E>。 自定义的类型,节省了Error这个模板代码 pub type Result<T> = std::result::Result<T, Error>; //Error也自己定义 pub enum Error { LineParse(String, usize), Io(io::Error), EnvVar(std::env::VarError), #[doc(hidden)] __Nonexhaustive } //实现一个not_found()的方法来判断是否为not_found的一个错误类型 impl Error { pub fn not_found(&self) -> bool { if let Error::Io(ref io_error) = *self { return io_error.kind() == io::ErrorKind::NotFound; } false } } //实现标准库中的error::Error这个trait impl error::Error for Error { //追踪错误的上一级,应该是打印堆栈这种功能 //如果内部有错误类型Err返回:Some(e),如果没有返回:None //关于'static这个生命周期的标注,我也不是很理解 //是指存储的错误生命周期足够长还是什么? fn source(&self) -> Option<&(dyn error::Error + 'static)> { match self { Error::Io(err) => Some(err), Error::EnvVar(err) => Some(err), _ => None, } } } //实现错误的打印 impl fmt::Display for Error { fn fmt(&self, fmt: &mut fmt::Formatter) -> fmt::Result { match self { Error::Io(err) => write!(fmt, "{}", err), Error::EnvVar(err) => write!(fmt, "{}", err), Error::LineParse(line, error_index) => write!(fmt, "Error parsing line: '{}', error at line index: {}", line, error_index), _ => unreachable!(), } } } 更详细的rust错误处理,可以参见:https://zhuanlan.zhihu.com/p/109242831 命令行参数 在main方法中我们可以看到第二行, let config = Config::from_args(); 这是influx使用了structopt这个crate,调用该方法后,程序会根据结构体上的#[structopt()]中的参数进行执行命令行解析。 #[derive(Debug, StructOpt)] #[structopt( //cargo的crate名字 name = "influxdb_iox", //打印出来介绍 about = "InfluxDB IOx server and command line tools", long_about = // 省略 ... )] struct Config { // from_occurrences代表出现了几次,就是-vvv的时候v出现的次数 #[structopt(short, long, parse(from_occurrences))] verbose: u64, #[structopt( short, long, global = true, env = "IOX_ADDR", default_value = "http://127.0.0.1:8082" )] host: String, #[structopt(long)] num_threads: Option<usize>, //subcommand代表是一个子类型的, //具体还有什么命令行要去子类型里继续解析, //这个字段不展示在命令行中 #[structopt(subcommand)] command: Command, } //在influx的命令行中提供了8个主要的命令, //在上一章中使用到的run参数就是属于Run(Box<commands::run::Config>)里的调用。 //这里都是subcommand,需要继续解析,这个在以后学习每个具体功能的时候再分析 #[derive(Debug, StructOpt)] 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), } 下面通过打印出来的例子来对应structopt中的内容。 $ ./influxdb_iox -vvvv run Config { verbose: 4, host: "http://127.0.0.1:8082", num_threads: None, command: Run(Config { rust_log: None, log_format: None, verbose_count: 0, writer_id: None, http_bind_address: 127.0.0.1:8080, grpc_bind_address: 127.0.0.1:8082, database_directory: None, object_store: None, bucket: None, aws_access_key_id: None, aws_secret_access_key: None, aws_default_region: "us-east-1", google_service_account: None, azure_storage_account: None, azure_storage_access_key: None, jaeger_host: None }) } 可以看到,我们执行了Run这个变体的Subcommand,并且指定了Config结构体中的verbose 4 次,IOx也成功的识别了。 后面继续学习程序的启动过程,祝玩儿的开心!

资源下载

更多资源
Mario

Mario

马里奥是站在游戏界顶峰的超人气多面角色。马里奥靠吃蘑菇成长,特征是大鼻子、头戴帽子、身穿背带裤,还留着胡子。与他的双胞胎兄弟路易基一起,长年担任任天堂的招牌角色。

Nacos

Nacos

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

Spring

Spring

Spring框架(Spring Framework)是由Rod Johnson于2002年提出的开源Java企业级应用框架,旨在通过使用JavaBean替代传统EJB实现方式降低企业级编程开发的复杂性。该框架基于简单性、可测试性和松耦合性设计理念,提供核心容器、应用上下文、数据访问集成等模块,支持整合Hibernate、Struts等第三方框架,其适用范围不仅限于服务器端开发,绝大多数Java应用均可从中受益。

WebStorm

WebStorm

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

用户登录
用户注册