首页 文章 精选 留言 我的

精选列表

搜索[es],共3901篇文章
优秀的个人博客,低调大师

Elasticsearch地理坐标类型(Geo-point)在Spring Data ES中的常见使用问题整理解答

下文整理的几个问答,本人在实际应用中亲身经历或解决过的,主要涉及Elasticsearch地理坐标类型(Geo-point)在Java应用中的一些特殊使用场景,核心依赖如下: <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-elasticsearch</artifactId> </dependency> A1. elasticsearch的geo_point类型对应java中的哪种数据类型? Q1. spring data elasticsearch中定义了GeoPoint这个类来实现两者之间的类型映射,此外还需要为当前字段添加@GeoPointField注解进行标志,注意GeoPoint应该使用org.springframework.data.elasticsearch.core.geo包下的。 /** * 坐标位置 */ @GeoPointField private GeoPoint location; A2. spring data elasticsearch中,如何以某坐标点为中心搜索指定范围的其它点? Q2. 建议尽可能通过继承ElasticsearchRepository<T, ID extends Serializable>来简化完成相关查询; ElasticsearchRepository 实现以某点为中心并搜索指定范围,首先定义如下: public interface TestRepository extends ElasticsearchRepository<Test, String> { } 其次可通过QueryBuilder接口来实现上述功能,参考如下: @Service public class TestService { @Resource private TestRepository testRepository; public Page<Test> findPage(double latitude, double longitude, String distance, Pageable pageable) { // 间接实现了QueryBuilder接口 BoolQueryBuilder boolQueryBuilder = new BoolQueryBuilder(); // 以某点为中心,搜索指定范围 GeoDistanceQueryBuilder distanceQueryBuilder = new GeoDistanceQueryBuilder("location"); distanceQueryBuilder.point(latitude, longitude); // 定义查询单位:公里 distanceQueryBuilder.distance(distance, DistanceUnit.KILOMETERS); boolQueryBuilder.filter(distanceQueryBuilder); return testRepository.search(boolQueryBuilder, pageable); } } A3. spring data elasticsearch中,如何计算两个给定坐标点之间的距离? Q3. 在GeoDistance类中定义了相关的计算方法,参考如下: GeoDistance // 计算两点距离 double distance = GeoDistance.ARC.calculate(srcLat, srcLon, dstLat, dstLon, DistanceUnit.KILOMETERS); 关于GeoDistance.ARC和GeoDistance.PLANE,前者比后者计算起来要慢,但精确度要比后者高,具体区别可以看这里。 A4. spring data elasticsearch应用中,如何以某个坐标点为中心,按距离近远排序搜索指定范围? Q4. 通过SearchQuery来实现,参考下面这段代码中GeoDistanceSortBuilder的使用: @Service public class TestService { @Resource private TestRepository testRepository; public Page<Test> findPage(double latitude, double longitude, String distance, Pageable pageable) { // 实现了SearchQuery接口,用于组装QueryBuilder和SortBuilder以及Pageable等 NativeSearchQueryBuilder nativeSearchQueryBuilder = new NativeSearchQueryBuilder(); nativeSearchQueryBuilder.withPageable(pageable) // 间接实现了QueryBuilder接口 BoolQueryBuilder boolQueryBuilder = new BoolQueryBuilder(); // 以某点为中心,搜索指定范围 GeoDistanceQueryBuilder distanceQueryBuilder = new GeoDistanceQueryBuilder("location"); distanceQueryBuilder.point(latitude, longitude); // 定义查询单位:公里 distanceQueryBuilder.distance(distance, DistanceUnit.KILOMETERS); boolQueryBuilder.filter(distanceQueryBuilder); nativeSearchQueryBuilder.withQuery(boolQueryBuilder); // 按距离升序 GeoDistanceSortBuilder distanceSortBuilder = new GeoDistanceSortBuilder("location", latitude, longitude); distanceSortBuilder.unit(DistanceUnit.KILOMETERS); distanceSortBuilder.order(SortOrder.ASC); nativeSearchQueryBuilder.withSort(distanceSortBuilder); return testRepository.search(nativeSearchQueryBuilder.build()); } } 如果这对您有帮助,欢迎分享和点赞,转载请注明出处!

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

网易云信 x Doris:降本70%、提速11倍, 统一 ES/InfluxDB/Hive 多技术栈的落地实践

导读:网易云信引入 Apache Doris 统一了原有 Elasticsearch、InfluxDB 和 Hive 多技术栈系统。凭借其高性能和易扩展的特点,提供一站式的数据存储和分析服务。实现机器成本降低 70%、实时场景查询提速 11 倍、离线任务耗时缩短 80% 的显著收益。 网易云信是网易旗下 ToB 的通信与视频云服务品牌,依托网易 20 多年的技术沉淀,为企业和开发者提供稳定、安全、高效的通信与视频云服务,包含 IM 即时通讯、视频云、短信、轻舟微服务、中间件 PaaS 等。截止当前,网易云信产品已覆盖用户 10 亿+,覆盖 196 个国家,覆盖地区 567 个。 云信数据平台承载了多条业务线的数据,数据存储于 Elasticsearch、InfluxDB 和 Hive 这三类主要数据库中,最终为各类业务场景提供数据支撑与服务。但随着业务发展,数据量急剧增加,这种多系统并存的方式带来了高昂的成本和复杂的维护难题。 为解决这些问题,网易云信引入 Apache Doris 统一了原有多技术栈系统。凭借其高性能和易扩展的特点,提供一站式的数据存储和分析服务。实现机器成本降低 70%、实时场景查询提速 11 倍、离线任务耗时缩短 80% 的显著收益。 早期架构和挑战 如上图云信数据平台早期架构,IM、RTC、短信、直播、点播等业务数据通过主动上报、日志采集,数据同步等方式上报至数据平台中。数据平台通过流式数据清洗或批处理对数据进行结构化或半结构化计算,并将结果存储到对应的存储计算引擎中,最终服务于多种使用场景。 离线场景下,主要采用 Hadoop 生态的 HDFS、Hive 以及 Impala 进行数据存储与分析;实时日志检索场景中,主要使用 Elasticsearch;而监控类数据则主要存储于 InfluxDB。该架构存在三大痛点: 存储冗余: 原有架构为了满足不同场景的数据查询分析需求,数据存储在 InfluxDB、Elasticsearch 和 Hive 多个组件中,同一数据冗余存储,造成存储资源严重浪费。 查询效率低: Elasticsearch 的非标准查询语法学习门槛高,且在处理大规模数据时响应慢。InfluxDB 在高基数数据或大量数据序列处理时消耗资源,查询体验不佳。 资源抢占严重: Hive 的资源队列紧张,运行效率低,常出现排队现象。InfluxDB 在高并发查询时,写入和查询资源无法隔离,导致写入延迟甚至错误。 Apache Doris 选型思考 基于上述痛点问题,网易云信开始寻找新的解决方案,并期望使用单一数据库服务多种数据服务,满足高效易用、低成本的升级及使用要求。 01 Doris 核心优势 经过探索,Apache Doris 符合我们的选型要求。其具备以下优势: 统一存储:Doris 具备替代 InfluxDB、Elasticsearch、Hive 等多种数据库的能力,仅需存储一份数据,实现数仓查询出口的统一,可实现数倍存储成本的降低。 实时 OLAP 能力:Doris 是一款基于 MPP 架构的高性能、实时的分析型数据库,强调即时数据分析,具备优秀的并发查询能力,能支持高并发点查询和高吞吐复杂分析场景。 冷热分层优势:借助 Doris 冷热分层功能,既能满足近期热点数据的查询性能,又能将历史数据冷备在便宜的存储介质中,实现查询性能和存储成本的平衡。 02 Doris vs ClickHouse 调研 Doris 的同时,我们也关注到了 ClickHouse,两个系统的关键点对比如下。 Doris 具备运维简单 、丰富的监控与诊断工具,以及支持配置热加载等优势,支持倒排索引和全文检索,适用于高并发实时分析、监控数据处理和日志分析等场景。它支持动态分区,简化数据接入流程。采用标准 SQL 语法并支持 MySQL 协议,降低了开发人员的学习成本。 统一高效的新架构 综合以上优势,网易云信将 Apache Doris 作为全新架构核心引擎,对原架构的存储计算引擎进行替换和统一,包括 InfluxDB、Elasticsearch、Hive 等多种存储系统。升级后,所有处理后的统一存储在 Doris 中,由其为各使用场景提供服务,降低了数据存储和系统运维的成本。 同时,借助 Doris 冷热分层能力,进一步优化了存储成本。之前使用 Hive 处理半结构化数据时,由于采用 HDD 存储且缺乏索引,查询效率较低。现在切换到 Doris 之后,前一天的数据存储在 SSD 上并构建了索引,离线分析效率显著高于原先 HDFS + Hive 组合。 此外,通过合并离线和在线集群和分时段资源调度机制,我们实现了计算资源的高效利用。白天(开发人员活跃期)集中承载数据查询分析需求,晚高峰时段(20:00-00:00)则全力支持云信业务峰值写入。离线计算任务通常安排在凌晨(00:00 之后)进行,用于处理前一天的数据(即 T+1 模式)。这种机制确保了从日间到夜间的计算能力全覆盖,使计算资源全天候高效运转,大幅提升了资源利用率。 落地实践 在业务迁移的过程中,网易云信也遇到一些问题及挑战。借此机会,将这些宝贵的优化经验整理并分享,希望对大家的使用有所指引及帮助。 01 写入优化 云信 - 数据平台在业务高峰期时,面临着百万级别以上的写入 TPS 以及高吞吐写入流量,这对系统性能提出了极高的要求。 经过多年的发展,云信已积累了大量数据源,不同的数据源分别写入 Doris 中对应一张表,提升 Doris 写入吞吐的关键在于提高攒批,而 Stream Load 仅支持单表批量数据写入,大量的数据源将造成写入服务为每个数据源分别攒批,导致内存消耗大。 为解决这一问题,我们尝试使用 HashMap 作为底层集合进行数据聚合,但在多线程环境下会引发哈希冲突。于是改用数组作为底层集合进行批量聚合,成功避免了哈希冲突,并支持进程内横向扩展线程数,可有效提升数据处理的能力。但若 Doris 端出现波动或业务量激增,写入端仍可能面临压力。 虽然通过扩容可以提升整体吞吐量,但可能导致数据分散,降低批量聚合效果。例如,原本可以聚合 1000 条数据一次性发送,现在可能只能聚合 100 条,增加了 Doris 数据处理的压力。 为此,我们进一步优化,将同一数据源的数据集中在一个写入集群内处理,非本集群的数据则投递到其他 Topic,由对应集群进行消费。这种方式提升了批量聚合能力,缓解了写入过程中小文件过多的问题,从而降低了 Doris Compaction 的压力。在数据格式上,将原本的 JSON 格式改为 CSV,并结合 GZIP 压缩,显著降低了带宽和内存的消耗。 最终,经过这些优化,线上写入峰值已达 600 万 QPS,集群总写入带宽最大达到 11GB/s,稳定控制在实际可承受范围内,整体运行非常平稳。 02 稳定性保障 为了确保上层业务的顺利迁移和数据库的稳定运行,我们采取了多重保障措施,包括资源隔离、集群隔离和 UDF 开发。这些措施旨在优化资源分配,减少不同业务之间的干扰,提升系统的整体性能和安全性。 资源隔离: 资源隔离主要应用于在线读、离线读、风险读这几类场景中,根据在线场景、离线场景以及风险账号的不同类型,设置多个资源组,确保不同资源组之间的用户互不影响,提升系统的稳定性和安全性。在线读的频次较高,但大多为简单查询;离线读的频次较低,主要为复杂查询,因此分配的内存和 CPU 资源相对较多;风险读是容易导致 Doris 抖动的高风险查询,实际应用中分配的资源较少。详细使用可参考:https://doris.apache.org/zh-CN/docs/admin-manual/workload-management/workload-group 数据隔离: 当前线上环境包含两类业务,数据源和业务类型各不相同。目前线上环境采用的是同一集群内的 BE Group 进行隔离,实现两个业务的数据隔离,且两个业务的写入服务和 Doris BE 节点是互相隔离的。这样既保证了资源的高效利用,也实现了多业务的并行承载。可参考:https://doris.apache.org/zh-CN/docs/admin-manual/workload-management/resource-group UDF 开发和应用: 用户从文件中读取数据并维护一个静态 HashMap 字典后,可以通过传入 Key 作为参数来返回 Value 作为查询结果。由于 Doris 中 的 JAR 包是按需加载的,只有在真正执行时才会加载到 JVM 中,执行完毕后则会卸载。当字典数据庞大且并发较高时,每个实例独立加载一份会显著增加内存占用。为此,我们选择提前将读取文件维护字典进行拆分,并打包字典生成 JAR,将其放入 custom_lib 中,使其能够在 BE 启动时由 JVM 加载。这样,可以在全局只维护一个实例,避免频繁的加载和释放,有效节约内存资源。 03 自动化运维体系 自动化运维是我们管理服务的核心能力之一,特别体现在动态分区和自动降冷功能上。借助配置中心和定时任务,能够动态调整 Doris 表的配置,例如分区分桶数的设置、历史数据迁移至冷备集群的时间以及未来分区大小的设置等。这些自动化运维能力极大地简化了 Doris 数据存储的管理,提高了系统的灵活性和效率。 动态分区管理: Doris 支持范围分区,但仅支持日期列的 Range 分区,不支持按照业务量动态调整分区大小。在查询时,必须指定日期才能正确定位到分区,这对查询效率造成了一定影响。此外,业务上有一些数据表数据量非常大,尽管尝试使用自动指定桶数量的配置,但仍存在一定问题。 当前做法:将 Bucket 的数据控制在 1G ~ 5G 之间,执行 SQL 每天自动删除过期分区,并创建未来分区。根据业务量自动创建和调整分区大小;可以通过配置中心在表级维度控制存储时长;使用时间戳字段进行分区,查询时方便指定分区,提高查询效率。 自动降冷策略: Doris 提供了自动降冷的功能,但由于我们的备件库中混用 2.5 和 3.5 的插槽机器数量较少,且大部分 HDD 为 3.5 英寸,因此未使用其提供的方案。 我们利用现有的大量 3.5 英寸备用 HDD 硬盘,新搭建了一批专用的 BE 节点,并将这些节点的 TAG 标签设置为 "cooldown"。随后,我们每天会将历史分区的 TAG 设置为 "cooldown",这样 Doris 就会自动将这些分区的数据调度迁移到成本较低的冷备节点 BE 中,实现了冷数据的高效存储和管理,还可对冷备资源进行归一化管理,方便弹性扩缩容。 示例: yunxin-manager 每天定时查询配置信息,根据配置信息,设置需要降冷分区的replication_allocation 属性到 HDD Group: ALTER TABLE lps_bucket.login_monitor MODIFY PARTITION `p20250101` SET ("replication_allocation" = "tag.location.cooldown:2"); 04 InfluxDB 迁移 建表 SQL 参考: CREATE TABLE `new_mediaServer_transport` ( `partitionTime` date NULL, `appid` bigint NULL, `cid` varchar(65533) NULL, `osType` varchar(65533) NULL, `sampleTime` bigint NULL, `sdkVersion` varchar(65533) NULL, `serverIp` varchar(65533) NULL, `transportId` varchar(65533) NULL, `recvBitrate` bigint NULL, INDEX idx_time (`lps_event_time`) USING INVERTED, INDEX idx_appid (`appid`) USING INVERTED, INDEX idx_cid (`cid`) USING INVERTED, INDEX idx_uid (`uid`) USING INVERTED, INDEX idx_partitionTime (`partitionTime`) USING INVERTED, INDEX idx_transportId (`transportId`) USING INVERTED, INDEX idx_serverIp (`serverIp`) USING INVERTED ) ENGINE=OLAP DUPLICATE KEY(`partitionTime`) PARTITION BY RANGE(`partitionTime`) ( PARTITION p20250614 VALUES [('2025-06-14'), ('2025-06-15')), PARTITION p20250615 VALUES [('2025-06-15'), ('2025-06-16')), PARTITION p20250616 VALUES [('2025-06-16'), ('2025-06-17'))) DISTRIBUTED BY RANDOM BUCKETS AUTO PROPERTIES ( "replication_allocation" = "tag.location.default: 2", "min_load_replica_num" = "1", "compression" = "ZSTD", "compaction_policy" = "time_series" ); InfluxDB 查询 SQL 参考: SELECT mean("recvBitrate") FROM "new_mediaServer_transport" WHERE ("cid" = '$cid' AND "uid" =~ /^$uid$/ AND producing = 'true') AND $timeFilter) GROUP BY time($__interval), "uid", "type", "tag_localHostName" fill(none) Apache Doris 查询 SQL 参考: SELECT avg(recvBitrate) as 'PubTransportRecvBitrate', (CEIL(`lps_event_time`) DIV (1 * 1000)) * 1 as time, uid, type, localHostName as tag_localHostName FROM lps_doris_bucket.new_mediaServer_transport WHERE ( cid = '$cid' AND uid rlike $uid AND producing = 'true') AND lps_event_time > $__unixEpochFrom() * 1000 AND lps_event_time < $__unixEpochTo() * 1000 GROUP BY time, uid, type, tag_localHostName ORDER BY time 应用收益 原集群配置与引入 Doris 后的集群配置对比如下: 成本优化:存储资源降低 70% 引入 Doris 后, 通过统一存储提高了 SSD 存储的利用率,避免了多套存储系统带来的资源浪费。当前 Doris 使用的物理机存储成本降低约 70%(以前的冷备存在 HDFS 中,此处未列出,因此仅计算 Doris SSD 存储部分)。 成本优化:计算资源降低 40% ~ 70% 在离线场景,按照 CPU/核 * s 计算,节省了 40% 的计算资源(当前还在持续改造中)。 在实时场景,使用 Doris 替换 InfluxDB 与 Elasticsearch 集群,总体 CPU 核数降低了约 70%。 效率提升:查询响应 离线任务耗时: 离线任务相对于 Hive,运行时间平均降低 80% 以上,消灭了所有小时级别的任务。 实时场景耗时: 在日志检索场景中,最近 3 小时、1 天 、7 天的日志检索,Doris 查询耗时保持稳定且均低于 4s,最快可在 1s 内响应。而 Elasticsearch 查询耗时呈现出较大的波动,最长耗时高达 75s,即使最短耗时也需要 6-7s。在更低的资源占用下,Doris 的查询效率至少是 Elasticsearch 的 11 倍 。 在高频次实时查询场景中, InfluxDB 出现多次异常波动,查询耗时直线上升,查询稳定性受到严重影响,Doris 的查询性能比 InfluxDB 更稳定,99 次查询均比较平稳、没有明显波动 。 未来规划 未来,我们将深入挖掘 Doris 的价值,探索其在其他场景中的应用: 湖仓一体的应用:进一步拓展联邦查询功能,整合多个数据源,打破数据孤岛。以联邦查询为基础,创建一个统一的分析入口,用户无需在多个数据源间切换,即可实时分析各类数据。 借助大模型搭建能效工具:借助 Doris 和大模型,实现交互式的可视化分析和智能可视化报表生成,搭建内部专业知识库。 深入了解 Doris:紧跟 Doris 社区动态,结合云信业务的使用反哺社区,为 Doris 的健康发展提供助力。

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

druid.io 海量实时OLAP数据仓库 (翻译+总结) (1)——分析框架如hive或者redshift(MPPDB)、ES等

介绍 我是NDPmedia公司的大数据OLAP的资深高级工程师, 专注于OLAP领域, 现将一个成熟的可靠的高性能的海量实时OLAP数据仓库介绍给大家: druid.io NDPmedia在2014年3月就开始使用, 见链接:http://blog.csdn.net/chenyi8888/article/details/37594771 druid是个很新的平台, 2013年底才开源出来, 虽然出现的比较晚, 但druid发展很快, 中国有几个公司开始使用, 2015年druid将会是爆发的一年 最近druid 的华人作者Fangjin从Metamarkets离职, 专门从事druid研发和推广. 以下翻译自http://druid.io/docs/0.7.1.1/, 并添加了自己的注解 什么是Druid Druid 是一个开源的,能在海量时序数据上 (万亿级别数据量, 1000 TB级别数据)上面提供实时分析查询的OLAP数据仓库,Druid提供了廉价的实时数据插入和任意数据探索的能力。 Druid的主要功能 为分析而生 - Druid是为了解决在OLAP工作流中进行探索分析而生的. 它提供了大量的filters, aggregators和 query 类型,并且提供了一个用户添加新功能的框架. 用户可以利用Druid的集群实现例如topN和直方图等功能。 (注: 传统数据库, 查询几千万的数据, 就会出问题, 查不出来) (注: druid就是一个能力超强的数据库, 执行例如SQL: select aColumn, bColumn sum(cColumn) from tableName where aColumn like 'xxx' and bColumn = 5 group by aColumn, bColumn having sum(cColumn) > 5 order by aColumn.) (注: druid对SQL支持有限,现在是实验版本。YeahMobi 重新开发适配了SQL, 屏蔽了下层平台, SQL 语句可以路由到这三个平台 druid, impala, hive) 高交互式 - Druid的低延时数据插入允许数据在生成之后的毫秒范围之内就可以被用户查询到。Druid通过读取和扫描需要的数据来优化查询的延时。 高可用性 - Druid可以被用来实现需要持续提供服务的SaaS应用。即使是在系统升级的过程中,你的数据仍然可以被查询。而且Druid 集群的扩容或者缩减不会带来数据的丢失。 (注: 已经在生产环境之中验证: 添加字段, 集群扩容, 集群缩减) 可扩展性 - 现有的Druid系统可以很轻松的处理每天数十亿条记录和TB级别的数据。Druid本身是被设计来解决PB级别数据的。 为什么要用Druid? Druid的初衷是为了解决在使用Hadoop进行查询时所遇见的高延时问题来提高交互性查询。尤其是当你对数据进行汇总之后并在你汇总之后的数据 上面进行查询时效果更好。将你汇总之后的数据插入Druid,随着你的数据量在不断增长,你仍然可以对Druid的查询能力非常有信心。当前的Druid 安装实例已经可以很好的处理以每小时数TB实时递增的数据量。 (注: 在我们的实践中 druid 查询统计100亿数据, 在5秒内响应。 查询1个月的数据, 基本可以在毫秒内完成。 比hadoop的常用的T+1 Map Reduce 高效多了. 你可以在拥有Hadoop的同时创建一个Druid系统。Druid提供了以一种互动式切片、切块方式来访问数据的能力,它在查询的灵活性和存储格式直接寻找平衡从而来提供更好的查询速度。 如果想了解更多细节,请参考 White Paper 和Design 文档. 什么情况下需要Druid? 当你需要在大数据集上面进行快速的,交互式的查询时 当你需要进行特殊的数据分析,而不只是简单的键值对存储时 当你拥有大量的数据时 (每天新增数百亿的记录、每天新增数十TB的数据) 当你想要分析实时产生的数据时 当你需要一个24x7x365无时无刻不可用的数据存储时 架构概述 druid在一定程度上是受搜索框架的启发, 通过建立不变数据视图和使用便于filter和aggregation的高度优化的格式来提高性能. Druid 集群有一系列不同类型的节点组成, 每种节点将一小部分事情做到极致。 Druid vs… Druid-vs-Impala-or-Shark Druid-vs-Redshift Druid-vs-Vertica Druid-vs-Cassandra Druid-vs-Hadoop Druid-vs-Spark Druid-vs-Elasticsearch 数据框架世界一直在巨大的混乱的变化之中, 这个网页希望帮助潜在的用户评估和确定druid适合用户解决遇到的问题。 如果有错误请通过邮件列表或者其他渠道反馈. 转自:http://www.cnblogs.com/lpthread/p/4519687.html 本文转自张昺华-sky博客园博客,原文链接:http://www.cnblogs.com/bonelee/p/6490891.html,如需转载请自行联系原作者

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

ES索引模板——就是在新建索引时候指定的正则匹配来设置mapping而已,对于自动扩容有用

索引模板 扩容设计»索引模板 Elasticsearch 不要求你在使用一个索引前创建它。对于日志记录类应用,依赖于自动创建索引比手动创建要更加方便。 Logstash 使用事件中的时间戳来生成索引名。 默认每天被索引至不同的索引中,因此一个@timestamp为2014-10-01 00:00:01的事件将被发送至索引logstash-2014.10.01中。 如果那个索引不存在,它将被自动创建。 通常我们想要控制一些新建索引的设置(settings)和映射(mappings)。也许我们想要限制分片数为1,并且禁用_all域。 索引模板可以用于控制何种设置(settings)应当被应用于新创建的索引: PUT /_template/my_logs { "template": "logstash-*", "order": 1, "settings": { "number_of_shards": 1 }, "mappings": { "_default_": { "_all": { "enabled": false } } }, "aliases": { "last_3_months": {} } } 创建一个名为my_logs的模板。 将这个 本文转自张昺华-sky博客园博客,原文链接:http://www.cnblogs.com/bonelee/p/7837282.html ,如需转载请自行联系原作者

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

科大讯飞:成本降低 60%,性能提升 10 倍,从 ES Loki 到 Apache Doris 可观测性存储底座升级

导读:科大讯飞星际日志中心经历了从 Elasticsearch 到 Loki,再到 Apache Doris 的可观测性存储分析底座升级,支持可观测三大支柱 Log Trace Metrics 的存储与分析,有效解决 Elasticsearch 成本高、Loki 查询慢的问题。Doris 能够在降低成本的同时提高查询效率,实现了查询性能提升 10 倍、存储空间缩减至 Elasticsearch 1/6。此外,Doris 提供的半结构化数据类型 VARIANT 能高效存储可扩展的 JSON 数据,具备很高的灵活性,且其性能媲美普通宽表。 作者|科大讯飞软件研发工程师,曲庆伟 指标、日志、链路是服务可观测性的三大支柱。 在保障服务稳定性方面,指标主要用于发现故障和问题,而日志和链路分析则侧重于定位和分析问题。其中,日志实际上在这三大维度中充当了重要桥梁。 讯飞星迹作为科大讯飞推出的一款全领域场景的可观测产品,旨在满足基础设施、应用及业务上的监测需求,提供一体化的解决方案。以实际场景为例,当服务上线后,需要对中间件和操作系统的 CPU 、内存等资源进行监控,这些指标能够实时反映系统的健康状态。当问题出现后,可以通过日志数据分析其来源,再利用日志链路分析迅速定位问题的位置。 为更好的处理及管理日志数据,讯飞星迹推出了星迹日志中心,经过多个版本迭代后, 现已基于 Apache Doris 升级为可观测存储分析底座,实现查询性能提升 10 倍,存储空间缩减至原来 1/6。目前,在集团内多个 BG BU 的项目上稳定运行,帮助业务降本增效,同时通过业务指标分析和用户画像与行为分析助力业务增长。 星迹可观测日志中心建设目标 科大讯飞对星迹可观测日志中心建设目标可以总结为三点:写的多、查的快、易于使用。 写的多: 高吞吐、低延迟:日志数据增长迅猛,每天可产生 10 ~ 100 TB 数据、日志条数可达几十亿~几百亿条。同时,日志写入并发量也在持续提升,可达到 50 万 TPS,对日志平台的写入吞吐要求极高。 低成本存储:日志数据总量庞大且保存时间较长,需实现高压缩率的存储方案,支持冷热存储和数据归档,避免占用过多存储资源。 查的快: 秒级查询响应:即使面对海量数据,也需要提供稳定高性能的查询能力,以应对故障排查等对响应速度有较高要求的使用场景,分钟级的数据延迟往往无法满足业务要求。 易于使用: 运维简单:要求架构极简易用,支持快速部署、扩容及运维,并能够与其他子系统串联打通,支持多云场景。 使用简单:支持工程师熟悉的接口和语言比如 SQL,以及可视化的、便捷的管理界面。 可观测性存储底座架构演进 01 ELK 早期讯飞基于 Elasticsearch 搭建了如图所示的日志处理及分析架构。具体来说,日志采集器将数据上报到 Kafka,接着通过 Logstash 进行 ETL 处理,处理后的数据存储到 Elasticsearch 中,由 Elasticsearch 提供数据的查询及分析服务。 存在的问题: 资源占用高:无论是写入还是查询,CPU 占用率普遍较高。这主要是由于高吞吐量写入时,分词操作和 Segment 合并会导致显著的 CPU 消耗。 存储成本高:受限于 Elasticsearch 的压缩率,海量日志数据存储成本相对较高。 稳定性差:在进行跨天查询或处理大数据量时,系统经常出现 OOM(内存溢出)问题,且故障恢复时间较长。此外,index 加载耗时较长,可能导致写入被拒绝。 02 Loki + Cassandra 第二代架构采用了基于 Grafana Loki 的轻量化架构,通过在各应用物理机上部署采集器来实现日志数据的收集。日志数据通过 Kafka 进入 Logstash 进行 ETL 处理,最终存储到 Loki 中。Loki 的后端存储选用 Cassandra ,其采用 ZSTD 算法压缩,相较于 Elasticsearch,存储成本节约了 5 倍。具体而言,单条日志 300KB,每分钟处理 1.8 万次,经过压缩后,1 天的日志数据仅需 1.4 TB 的存储空间。 存在的问题: CPU 使用率高、查询分析效率低:Loki 的架构与 Prometheus 原理相似。当以标签进行数据查询时,系统先根据 index 定位到 chunk,然后加载 chunk 中的数据,解压后逐条暴力搜索。这种处理方式显然会导致较高的 CPU 使用率,并且经常出现 OOM(内存溢出)问题,查询和分析的效率也并不理想。 创建标签的数量受限: 当标签基数较少时,查询速度尚可。而当标签数量过多或搜索条件复杂时,无法对基数较大的值进行查询加速,而日志场景中的字段基数通常较高。例如,在查询每个链路的 trace ID 时,每条 trace ID 数据几乎都是唯一的,这种高基数的数据不适合使用标签来处理。 03 基于 Apache Doris 的可观测性存储底座 进一步的,为解决上述架构存储的问题,讯飞引入了 Apache Doris。使用 Apache Doris 替换了 Loki 服务。采集后的数据依旧流转至 Kafka 中,消费到日志服务加工处理,按时间和数据大小攒批,设置攒批时长为 3 分钟、数据大小为 200M,满足任一条件则触发数据发送,处理完的数据通过 Stream Load 写入到 Doris 集群中。 Apache Doris 引入后,可以满足文章最初提出的 3 个建设目标: 写的多: 可支撑日均 600 亿条、10TB 的写入流量。与 Elasticsearch 相比,存储成本不及其 1/6,相同的数据 Elasticsearch 存储需要 1T,而依赖于 Doris 的列式存储和 ZSTD 压缩,仅需要 170 G。 查的快: 查询效率提升至少 10 倍,特别是在聚合分析、短语模糊匹配及 TOPN 命中前缀索引等场景下,性能提升效果显著。Doris 即使在没有命中索引的情况下,也可以在分钟内返回结果。与 Elasticsearch 频繁出现的 OOM 相比,效率提升尤为突出。 易于使用: Doris manager 可便捷管理所有 Doris 集群。科大讯飞在交付私有化项目时,通常需要管理至少数十套集群,而 Doris Manager 使这一过程变得便捷轻松。此外,系统还提供 Grafana 和自研的 Web 查询界面,方便用户进行日志检索和分析。 新架构的应用及优化经验 01 Log Trace Metric 选择合适的数据模型 Doris 数据模型分为主键模型、明细模型和聚合模型,这三个模型在可观测日志中心都有不同程度的应用。 主键模型: 主要应用于调用链数据 Trace、配置数据的处理。主键模型能够保证 Key(主键)的唯一性,当用户更新一条数据时,新写入的数据会覆盖具有相同 Key(主键)的旧数据。 CREATE TABLE config ( app_code varchar(50) NOT NULL COMMENT '业务唯一编码', app_name varchar(50) NOT NULL , table_name varchar(200) NOT NULL COMMENT '数据表名', create_time datetime NOT NULL COMMENT '创建时间', update_time datetime NOT NULL COMMENT '更新时间' ) ENGINE = OLAP UNIQUE KEY(app_code) DISTRIBUTED BY HASH(app_code) BUCKETS AUTO PROPERTIES ( "enable_unique_key_merge_on_write" = "true", "light_schema_change" = "true" ); **明细模型:**主要用于日志数据 Log 的处理。数据将按照导入文件中的内容进行存储,且不进行任何聚合,即使 Key 列完全相同的数据也会被保留。在建表语句中指定的 Duplicate Key 仅用于指明数据存储时按哪些列进行排序,以适应业务不断产生的数据。一旦数据产生,就不会再发生变化。 CREATE TABLE log_record ( `log_date` DATETIMEV2(3) COMMENT "日志打印时间", `source` VARCHAR(100) COMMENT "日志来源", `level` VARCHAR(10) COMMENT "日志级别", `host_name` VARCHAR(100) COMMENT "主机名", `msg` STRING COMMENT "日志内容", INDEX idx_msg (`msg`) USING INVERTED ) ENGINE = OLAP DUPLICATE KEY(`log_date`) PARTITION BY RANGE (`log_date`) ( ) DISTRIBUTED BY RANDOM BUCKETS 50 PROPERTIES ( "compaction_policy" = "time_series", "dynamic_partition.enable" = "true", "dynamic_partition.time_unit" = "DAY", "dynamic_partition.create_history_partition" = "true", "dynamic_partition.start" = "-30", "dynamic_partition.end" = "7", "dynamic_partition.prefix" = "p", "compression"="zstd" ); 聚合模型: 主要用于告警指标 Metrics 的聚合。根据 Key 列聚合数据,通过提前聚合大幅提升性能,极大降低聚合查询时所需扫描的数据量和查询的计算量,适合有固定模式的报表类查询场景。 CREATE TABLE alarm_agg ( `id` LARGEINT, `first_trigger_time` DATETIME MIN COMMENT "第一次触发时间", `level` TINYINT MAX COMMENT "告警级别", `msg` STRING REPALCE COMMENT "告警内容", `num` INT SUM COMMENT "总数量" ) AGGREGATE KEY(`id`) DISTRIBUTED BY HASH(`id`) BUCKETS 10; 当相同的 ID 进入时,将根据聚合函数进行处理。例如,对于触发时间,将提取最小值作为第一次触发时间。此外,还会对某些级别和内容进行替换,并对数量进行求和。这一过程可以提前对原始数据进行聚合,从而方便后续查询。通过这种方式,后续查询时能够扫描更少的数据。 数据聚合发生的时段分为三个: 每一批次数据导入的 ETL 阶段。每一批次导入的数据内部进行聚合,聚合后写入 Doris 集群。 BE 进行数据 Compaction 阶段。BE 会对已导入的不同批次的数据进行进一步的聚合。 数据查询阶段。对于查询涉及到的数据,进行对应的聚合。 聚合模型的局限性 对 count(*) 查询不友好, Doris 必须扫描所有的 AGGREGATE KEY 列,并且聚合后,才能得到语意正确的结果。 当聚合列非常多时,count(*) 查询需要扫描大量的数据。 02 可扩展 JSON 数据使用半结构化数据类型 VARIANT { "source":"PC", "time":1722562556534, "userId":"d3079d82", "properties": { "duration":8979, "mark":"标识", "title":"标题", "url":"/statistic/analysis" } } 上述示例展示了一个典型的行为分析日志,其中 properties 是 JSON 格式的可扩展字段,存储了运行时长、标题和 URL 等子字段。然而,不同业务线在 properties 字段中的具体数据可能有所差异。例如,某业务线可能仅依据其标识来记录数据,如果该业务线涉及启动参数或启动状态,字段数量可能会减少,仅包含用户信息,如用户 ID。根据不同的事件类型,整个 properties 字段的内容也会有所不同。因此,在日志场景中,经常需要存储一些动态字段,如 Kubernetes 下的标签以及采集指标数据的标签等。 2-1 基于导入手动指定 json_path 静态 Schema 方案: 对于 properties 这种可扩展的 JSON 字段,最初在使用 Doris 导入时,手动配置json_path 参数以映射到表的字段。为提供用户体验,我们还在页面上添加了可视化配置映射功能。该方案的优缺点都很明显。 优点:能够确保存储和查询的效率,所有字段都映射到 Doris 表中的普通列。相较于 JSON 文本数据,列式存储在存储和查询方面的优势非常明显。 缺点:不够灵活,无法动态扩展字段。每次修改都需手动调整表结构和映射关系,上下游系统需协同处理,这一过程繁琐且低效,容易出错。 2-2 基于 Doris Variant 动态 Schema 方案: 后来,我们将 Apache Doris 升级到了 2.1 版本,该版本提供了 VARIANT 数据类型,支持嵌套的不固定 Schema。VARIANT 数据类型可以存储任何合法的 JSON,可自动从 JSON 中抽取字段并推断其类型,并将这些字段存储为 VARIANT 列的子列。 使用 VARIANT 类型的表结构和查询语句如下所示: CREATE TABLE log_variant ( time BIGINT, source STRING, userId STRING, properties VARIANT ) SELECT * FROM log_variant WHERE time BETWEEN t1 AND t2 AND source = 'PC' AND properties['duration'] > 100000; 拆分成子列并采用列式存储的方式,使得 VARIANT 具备良好的存储和分析性能。在进行聚合/过滤/排序等查询时,只需读取 VARIANT 子列数据(比如 properties['duration']) 即可,不会产生额外的数据读取和解析的开销(比如mark, title, url等字段)。这种方案的性能与静态列相当,而相较于 JSON 字符串,性能提升存在数量级的差异。 在使用 VARIANT 类型的过程中,也遇到一些局限性并已找到解决方法。 早期版本 VARIANT 不能应用在聚合模型中, 升级到 2.1.3 以上的版本可以解决。 当 VARIANT 与 Group Commit 一起使用时,会出现过早反压影响写入性能的问题,升级到 2.1.5 以上的版本可以解决。 在 VARIANT 列上创建索引时,会对所有子列同时创建索引。如果子列数量过多,可能导致索引数量过多,从而影响写入性能。此外,同一 VARIANT 列的分词属性是一致的。例如,如果某列包含十个字段,并使用中文分词,那么这十个字段都将应用相同的中文分词规则。 日期、Decimal 等非标准 JSON 类型,会被默认推断成字符串类型,应尽可能从 VARIANT 中提取出来,用静态类型性能更好。 综合来看,建议需要 VARIANT 的用户使用最新的 2.1.x 版本,目前最新版本为 2.1.6。针对问题 3 和 4,社区已经开发了 VARIANT 内部 Schema 允许用户自定义的功能,预计将在后续版本中正式发布。 03 建表分区分桶优化 Doris 支持以表进行逻辑分区,一个 Table 下面有多个 Partition,每个 Partition 下有多个 Tablet。 如上图所示,Partition 是按照时间进行分区的,假设今天是 8 月 15 日,那么今天写入的数据将写入到右边分区中。该分区的意义在于,将总体数据量进一步拆分,查询时只查符合条件的分区,以此来提升查询的效率。 Tablet 分桶一般有 Hash 和 Random 这两种方式: Hash: 数据写入时将指定 Key 进行 Hash 计算,以确定该分配到哪个 Tablet。这种方法的优势在于,当根据 ID 进行 Hash 时,查询时可以精准地定位到相应的 Tablet。具体来说,先根据时间进行过滤,筛选到特定的数据位置,由于数据是基于 Hash 进行分桶的,这样可以轻松地命中相应的 Tablet。 Random: 在数据写入时,系统会随机选择一个 Tablet,并结合 Single Tablet Load 实现将整个 Batch 数据写入同一个 Tablet 的单个文件。这种方式不仅增加了 IO 批处理的规模、提高了写入的性能,还使得时间相近的数据能够存储在一起,从而优化了存储效率。 分桶数一般设置为磁盘数量的三倍,以确保每个 Tablet 的大小保持在 1-10GB 范围内。 CREATE TABLE log_record ( log_date DATETIMEV2(3), level VARCHAR(10), host VARCHAR(100), message string, INDEX idx_host (host) USING INVERTED), INDEX idx_message (message) USING INVERTED PROPERTIES("parser" = "chinese", "support_phrase" = "true") ) ENGINE = OLAP DUPLICATE KEY(log_date) PARTITION BY RANGE (log_date)() DISTRIBUTED BY RANDOM BUCKETS 50 PROPERTIES ( "compaction_policy" = "time_series", "compression"="zstd“, "dynamic_partition.enable" = "true", "dynamic_partition.time_unit" = "DAY", "dynamic_partition.create_history_partition" = "true", "dynamic_partition.start" = "-30", "dynamic_partition.end" = “3", "dynamic_partition.prefix" = "p" ); 分区分桶优化前后的效果对比非常明显。 例如,下面的根据时间查询前 20 条数据。 SELECT * FROM log_record WHERE log_date >= '2024-08-10 00:00:00.000' AND log_date <= '2024-08-11 00:00:00.000' ORDER BY log_date DESC LIMIT 0, 20 优化前(左)虽然成功返回了 20 条数据,但其扫描的行数和数据量(ScanRowsRead, ScanBytesRead)远大于优化后(右)。优化后的查询以时间为分区,并以时间作为前缀索引,该方式适用于分页场景下,能够更好的提高查询的性能,如红色圈注处,只需在每个节点扫描前 20 行数据,即可完成查询。 04 写入优化 FE 参数优化:尽量让每个节点的 Tablet 均匀 enable_round_robin_create_tablet = true tablet_rebalancer_type = partition BE 参数优化:增加写入 Buffer Size 的大小,write_buffer_size = 1073741824 Stream Load 参数优化:设置单 Tablet 写入,load_to_single_tablet=true 我们会在写入前攒批,压力分摊到写入模块: 写入前攒批与 Group Commit 原理相似(为什么不使用 Group Commit?最初在使用时, Group Commit 还未发布),将 Kafka 消费的数据进行本地存储,也就是说相当于在写入 Doris 之前,数据已经被处理完成。通过这样的方式将写入压力分摊至写入模块,可有效减轻 BE 的压力。 此外,Doris 提供了两种 compaction 策略:size_based 及 time_series,其中 time_series适用于时序场景。 时序数据指每批导入的数据具有时间上的顺序性。由于数据写入到 Tablet 中是有序的,因此在合并时,我们可以直接对 Segment 文件进行合并。然而,在实际应用场景需要注意,可能会出现新的数据插入到原本有序的数据中。比如,第一次写入 10-12 点数据,第二次写入 13-14 点数据,第三次又写入 8-9 点数据。在这种情况下,我们通常会默认采用基于 size 的 Compaction 策略。 优化完成后,在三节点集群测试下,每分钟可写入 600 万条数据,数据流量约 4.5G。整个过程中,磁盘 IO 均值不超过 9%,BE 内存占用均值为 4G,CPU 占用均值 9%。相对于其他系统,Doris 在写入优化方面是非常优秀的,写入吞吐高延迟低,而且资源占用率低。 收益总结 可观测性存储分析底座从 ELK 到 Loki 再到 Apache Doris 的架构升级,带来了多个方面的提升,总结起来还是最初提到的 3 个点:写的多,查的快,易于使用,具体表现在: 写的多: 可支撑日均 10TB,600 亿条日志的写入流量。与 Elasticsearch 相比,存储成本不到其 1/6。 查的快: 查询效率提升至少 10 倍,特别是在聚合分析、短语模糊匹配及 TOPN 命中前缀索引等场景下,性能提升效果显著。 易于使用: Doris Manager 使得管理多个集群变得简单,同时提供 Grafana 和自研的 Web 查询界面,方便用户进行日志检索和分析。 未来展望 未来,还将在 Doris 的基础上进行以下规划: 基于讯飞星火大模型的 AIOps:持续探索智能运维的最佳实践,包括日志异常监测、故障预测和故障诊断等。 用户行为分析:目前采用手动创建物化视图,未来探索如何利用 Doris 的物化视图能力实现自动物化。 存算分离:计划基于 3.0 版本的存算分离,实现中心化的读写分离以及租户的物理隔离。

资源下载

更多资源
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应用均可从中受益。

Sublime Text

Sublime Text

Sublime Text具有漂亮的用户界面和强大的功能,例如代码缩略图,Python的插件,代码段等。还可自定义键绑定,菜单和工具栏。Sublime Text 的主要功能包括:拼写检查,书签,完整的 Python API , Goto 功能,即时项目切换,多选择,多窗口等等。Sublime Text 是一个跨平台的编辑器,同时支持Windows、Linux、Mac OS X等操作系统。

用户登录
用户注册