搜狐技术产品

关键词命中统计ClickHouse改造

Image

 Image 

本文字数:8154字

预计阅读时间:21分钟

目录

  • 1.背景
    • 1.1. 业务需求
    • 1.2. 问题难点
  • 2. Clickhouse解决方案
    • 2.1 Clickhouse简介
    • 2.2 Clickhouse如何解决传统方案存在问题
  • 3. 使用ClickHouse 进行关键词命中统计
    • 3.1 创建集群
    • 3.2 建表
    • 3.3 与传统方案对比
  • 4. 总结与展望

1.背景

1.1 业务需求

关键词统计需求是审核系统中业务人员,需要对审核中所有关键词每天的命中情况进行统计,并且要求根据触发次数、数据删除率进行筛选和排序Image

1.2 问题难点

  1. 关键词数据量大,关键词表本身有 将近70W的数据,每天都进行统计,也就是每个月有将近2100W 数据 。
  2. 需求上删除率是指删除次数/触发次数,是一个计算结果,需要把每个词的 在指定时间内的所有触发情况、删除情况全部计算后,再进行筛选,排序。

1.2.1 传统Mysql方案

由于数据量巨大,我们需要先将mysql进行分表,按照当前数据量,按月进行分表,每个表最多2100W 数据。

  1. 根据用户的时间条件确认数据范围在哪些个表中。
  2. 根据其他筛选条件 对数据进行筛选
  3. 聚合数据,统计触发数量和命中删除率,并对聚合后的数据进行筛选
  4. 按照排序条件,分别在每个表中进行排序取前N条(业务要求,此N值可配置, 当前为5000条),最后将结果合并排序,取最终的前N条。
  5. 纪录统计结果并进行展示。

1.2.2 Mysql-java 方案

上述方案中,第3、4 两步发生在mysql中,如果数据量大会导致当前mysql节点无法响应其他请求,另外,由于关键词触发数量、删除率的筛选和排序,是对数据聚合之后再进行筛选,所以加索引的效果并不明显。

故而需要将3、4步骤 放到jvm内存中进行处理

配置内存中处理的窗口数据量

  1. 从数据库中滚动获取一定量数据,将符合筛选条件的数据保留下来。
  2. 将上一步保存下来数据 同内存中已暂存的N条数据进行比对,保留符合排序条件的TopN,重复上述两步,直到将当前表的数据全部处理完。

1.2.3 传统方案存在问题

  • 数据查询用到的筛选条件多而灵活: 为每种查询条件都配置索引的话,索引量较大,影响写入效率。
  • 数据聚合后再筛选排序: 基于触发数量、删除率的的筛选和排序 是对数据进行聚合之后再进行筛选和排序的,传统数据库的索引对这种数据影响不大。
  • 整体数据量较大: 全部放到数据库中进行筛选和排序,会严重影响数据库本身的并发能力。

2. ClickHouse 解决方案

2.1 Clickhouse 简介

ClickHouse 是俄罗斯Yandex在2016年年开源的一个⾼高性能分析型SQL数 据库,主要⾯面向OLAP场景。开源之后,凭借优异的查询性能,受到业界的青睐。

优点:

1)为了高效的使用CPU,数据不仅仅按列存储,同时还按向量进行处理;

2)数据压缩空间大,减少io;处理单查询高吞吐量每台服务器每秒最多数十亿行;

3)索引非B树结构,不需要满足最左原则;只要过滤条件在索引列中包含即可;即使在使用的数据不在索引中,由于各种并行处理机制ClickHouse全表扫描的速度也很快;

4)写入速度非常快,50-200M/s,对于大量的数据更新非常适用;

2.2 Clickhouse 如何解决传统方案存在问题

2.2.1 列式存储 - 解决灵活筛选条件、快速响应

与行存将每一行的数据连续存储不同,列存储将每一列的数据连续存储。示例图如下:

Image

相比于行式存储,列式存储在分析场景下有着许多优良的特性。

1)在行存模式下,数据按行连续存储,所有列的数据都存储在一个block中,不参与计算的列在IO时也要全部读出,读取操作被严重放大。而列存模式下,只需要读取参与计算的列即可,极大的减低了IO cost,加速了查询。

2)同一列中的数据属于同一类型,压缩效果显著。列存往往有着高达10十倍甚至更高的压缩比,节省了大量的存储空间,降低了存储成本。

3)更高的压缩比意味着更小的data size,从磁盘中读取相应数据耗时更短。

4)高压缩比,意味着同等大小的内存能够存放更多数据,系统cache效果更好。

使用列式存储数据结构,无论筛选哪一列数据,都可以有非常快速的查询响应,整体查询速度也会非常快。

下面截取了测试环境下 某一个节点的数据压缩情况

Image

可以看到数据压缩比非常高,100G+的数据只用了不到6G的存储容量,可以说非常节约存储,同时如此高的压缩比下,读取大量数据只需要相对较少的IO次数,并且系统缓存效果也会非常好。

2.2.2 Clickhouse 多核并行与向量化执行 - 解决聚合后再排序筛选

2.2.2.1 多核并行

ClickHouse将数据划分为多个partition,每个partition再进一步划分为多个index granularity,然后通过多个CPU核心分别处理其中的一部分来实现并行数据处理。

在这种设计下,单条Query就能利用整机所有CPU。极致的并行处理能力,极大的降低了查询延时。

2.2.2.2 向量化执行与SIMD(单指令多数据运算)

ClickHouse不仅将数据按列存储,而且按列进行计算。传统OLTP数据库通常采用按行计算,原因是事务处理中以点查为主,SQL计算量小,实现这些技术的收益不够明显。但是在分析场景下,单个SQL所涉及计算量可能极大,将每行作为一个基本单元进行处理会带来严重的性能损耗:

1)对每一行数据都要调用相应的函数,函数调用开销占比高;

2)存储层按列存储数据,在内存中也按列组织,但是计算层按行处理,无法充分利用CPU cache的预读能力,造成CPU Cache miss严重;

3)按行处理,无法利用高效的SIMD指令;

ClickHouse实现了向量执行引擎(Vectorized execution engine),对内存中的列式数据,一个batch调用一次SIMD指令(而非每一行调用一次),不仅减少了函数调用次数、降低了cache miss,而且可以充分发挥SIMD指令的并行能力,大幅缩短了计算耗时。向量执行引擎,通常能够带来数倍的性能提升。

借助 ClickHouse 多核并行,每个核心又针对同一列上的数据进行快速大量的筛选计算并排序,从指令集这一层进行优化,大大缩短了计算的耗时。

举个例子:传统计算方法进行筛选,如果筛选列 column1 没有索引,数据库会读取每一行数据,并对数据进行进行筛选条件命中,命中后,保留筛选成功的数据。

而在 clickhouse中,首先column1 这一列的数据,会分段存在很多个文件中,每一个cpu核心处理其中一组文件,其次,使用向量执行引擎,每次执行一个指令是针对寄存器内很多个数据进行执行的,相比于传统方案中,每次读取一行数据到内存,这一行数据再放到寄存器中让cpu参与计算,获取计算结果这样的方式,会有数倍的性能提升。

2.2.3 Clickhouse 数据分片与分布式计算 - 处理大数据量问题

2.2.3.1 数据分片

Clickhouse 支持数据分布式存储,在分布式模式下,ClickHouse会将数据分为多个分片,并且分布到不同节点上。不同的分片策略在应对不同的SQL Pattern时,各有优势。ClickHouse提供了丰富的sharding策略,让业务可以根据实际需求选用。

1) random随机分片:写入数据会被随机分发到分布式集群中的某个节点上。

2) constant固定分片:写入数据会被分发到固定一个节点上。

3)column value分片:按照某一列的值进行hash分片。

4)自定义表达式分片:指定任意合法表达式,根据表达式被计算后的值进行hash分片。

2.2.3.2 分布式计算

除了优秀的单机并行处理能力,ClickHouse还提供了可线性拓展的分布式计算能力。ClickHouse会自动将查询拆解为多个task下发到集群中,然后进行多机并行处理,最后把结果汇聚到一起。

在存在多副本的情况下,ClickHouse提供了多种query下发策略:

  • 随机下发:在多个replica中随机选择一个;
  • 最近hostname原则:选择与当前下发机器最相近的hostname节点,进行query下发。在特定的网络拓扑下,可以降低网络延时。而且能够确保query下发到固定的replica机器,充分利用系统cache。
  • in order:按照特定顺序逐个尝试下发,当前一个replica不可用时,顺延到下一个replica。
  • first or random:在In Order模式下,当第一个replica不可用时,所有workload都会积压到第二个Replica,导致负载不均衡。first or random解决了这个问题:当第一个replica不可用时,随机选择一个其他replica,从而保证其余replica间负载均衡。另外在跨region复制场景下,通过设置第一个replica为本region内的副本,可以显著降低网络延时。

借助 ClickHouse 分布式存储与分布式计算能力,将海量数据分散到多个节点中,避免一个节点数据量过大,造成瓶颈。同时,将一个查询分配到多台机器上并行执行,大大降低响应时间,再者,一份数据多个副本,多个副本间互不影响,整体上可以同时在多个副本上进行查询,提高了服务整体的并发能力。

2.2.4 OLAP及其他数据库产品技术选型

除了Clickhouse之外,我们也对比了市面上常见的其他两种数据库,TiDB、DorisDB,与Clickhouse进行横向对比,结合我们现有业务,分析如下:

数据库选型TiDBDorisDBClickhouse
优势支持高频词update操作集群功能丰富多表关联性能良好支持update 操作支持多种表引擎支持多种复合结构
缺点复杂结构只支持Set、Json,Json以Binary形式存储,不能索引。当前不支持update操作,需要使用insert overwrite 进行模拟操作只能在duplicate table中创建Array类型非标准SQL,有学习成本多表join性能不好建立分布式表时,需要先建立本地表,再建立相应视图,过程繁琐。

通过横向对比优缺点,在我们当前业务场景下最终先择了 clickhouse

原因如下:

  1. 查询性能上 doris、clickhouse明显优于TiDB,本身业务上存在数据更新,但不会特别频繁。
  2. Doris 对复合数据类型支持不够,例如:Array,在关键词业务中,会给关键词打上一些产品标签,这些产品标签也是会随着业务变更,进行增删的。同时不同产品标签下,关键词作用也不同,需要进行查询。

3 使用ClickHouse 进行关键词命中统计

下面 介绍一下实际使用中,审核后台如何利用ClickHouse 进行关键词命中统计。

关键词命中统计业务上比较简单。

每天凌晨通过前一天的审核记录统计,前一天每个关键词的命中情况。

查询时,直接使用每天关键词命中情况进行聚合统计。

3.1 创建集群

创建cluster 命名 cluster_3shards_2replicas包含6个clickhouse节点, 3个分片,2个副本

Image

占用三台物理机 别名 ch19-117-253, ch19-117-254, ch19-117-255 每台机器上有两个clickhouse实例,分别占用9000 和 9001端口,同一个分片的不同副本分配到不同机器上。

3.2 建表

涉及到的表结构有两个,一个是关键词表,记录关键词属性。一个是关键词命中情况表,记录每一天关键词的命中情况。

3.2.1 为什么是两张表?

这里说明下,为什么没有将统计表设计成一张大宽表,而是分成两个表。

关键词表在当前业务中,从统计上讲相对频繁,业务上在一些特定时期(例如:6-4时期、国庆时期、佩洛西访台等特殊事件),关键词力度、关键词类型会有变化,简单来说,有些词,在特定时期就会从普通警示词提升为禁词。当词性变化时,如果需要对历史数据进行统计查看,就需要洗数据,这样宽表得不偿失,所以我这里把变化的关键词,不变的关键词命中情况,分成了两张表进行记录。当有禁词发生变化时,通过mq消息,只修改关键词表。

3.2.2 创建关键词表

使用clickhouse 的分布式DDL创建语句可以直接在集群中,为每个节点都创建一份本地表

create table monitor_banword_local on cluster cluster_3shards_2replicas (  
    id Int32 COMMENT '关键词ID',  
    type UInt8 COMMENT '关键词类型',  
    name String COMMENT '关键词',  
    product Array(UInt8) COMMENT '适用产品',  
    parent_id UInt8 COMMENT '级别',  
    is_core UInt8 COMMENT '是否核心词:0:否;1:是',  
    status UInt8  
) ENGINE = MergeTree  
PARTITION BY type  
ORDER BY type  

3.2.3 创建关键词命中情况统计表

CREATE TABLE banword_hit_statistics_mt_2021_local on cluster cluster_3shards_2replicas  
(  
    `id` Int32 COMMENT '统计ID',  
    `hit_banword_id` Int32 COMMENT '命中关键词ID',  
    `trigger_times` Int32 COMMENT '触发次数',  
    `delete_times` Int32 COMMENT '删除次数',  
    `delete_rate` Decimal(9, 6) COMMENT '删除率',  
    `snapshot_date` Date COMMENT '数据快照时间',  
    `update_time` Date COMMENT '更新时间'  
)  
ENGINE = MergeTree  
PARTITION BY toYYYYMM(snapshot_date)  
ORDER BY (snapshot_date, hit_banword_id)  

3.2.4 为节点上的本地表创建分布式视图

CREATE TABLE monitor_banword_all ON CLUSTER cluster_3shards_2replicas (  
      id Int32 COMMENT '关键词ID',  
      type UInt8 COMMENT '关键词类型',  
      name String COMMENT '关键词',  
      product Array(UInt8) COMMENT '适用产品', 
      parent_id UInt8 COMMENT '级别',  
      is_core UInt8 COMMENT '是否核心词:0:否;1:是',  
      status UInt8  
) ENGINE = Distributed(cluster_3shards_2replicas, default, monitor_banword_local,intHash64(type));  
CREATE TABLE banword_hit_statistics_mt_2021_all ON CLUSTER cluster_3shards_2replicas  
(  
    `id` Int32 COMMENT '统计ID',  
    `hit_banword_id` Int32 COMMENT '命中关键词ID',  
    `trigger_times` Int32 COMMENT '触发次数',  
    `delete_times` Int32 COMMENT '删除次数',  
    `delete_rate` Decimal(9, 6) COMMENT '删除率',  
    `snapshot_date` Date COMMENT '数据快照时间',  
    `update_time` Date COMMENT '更新时间'  
) ENGINE = Distributed(cluster_3shards_2replicas, default, banword_hit_statistics_mt_2021_local, toMonth(snapshot_date) ); 

3.2.5 查询示例

统计 2022-01-01 至 2022-07-01 之间 禁词类型为屏蔽性禁词,触发次数大于1的数据,按照删除率倒叙排序取前5000条数据。

select   
   hit_banword_id,   
   sum(trigger_times) as triggerTimes,   
   sum(delete_times) as deleteTimes,  
   if (deleteTimes>0, divide(deleteTimes, triggerTimes), 0) as deleteRate    
from default.banword_hit_statistics_mt_2021_all  
where   
   hit_banword_id GLOBAL IN (  
                select id from monitor_banword_all  
                where parent_id = 2  
                )  
   and snapshot_date between toDate('2022-01-01') and toDate('2022-07-01')  
group by hit_banword_id  
having   
   triggerTimes > 1  
order by  deleteRate  desc  
limit 5000  
Image

Image

可以看到总数据量在1.2亿的情况下,进行筛选、分组聚合、聚合后再筛选、排序这样的查询用例下查询时间只要 1.837s

3.2.6 从传统方案到ClickHouse迁移踩过的坑

  • 历史数据迁移到ClickHouse时,批量写入数据占用过多分区。

报错:Too many partitions for single INSERT block (more than 100)

原因:clickhouse 使用max_partitions_per_insert_block  参数用来限制单个插入Block中,包含的最大分区数量,默认值为100。设置为0时,表示不限制

修改方法:

1.创建会话时,SET max_partitions_per_insert_block=1000,

2.在users.xml中,<profile>模块中添加

<max_partitions_per_insert_block>1000</max_partitions_per_insert_block>

  • SQL语句支持较为简单

举例:

  1. 如果建立jdbc连接时没有选择数据库,查询语句不会默认使用default库,而是需要在查询时,查询default.{表名}
  2. 使用insert into table 插入数据时,数据列需要使用``引用起来,注意这个是反引号,一般在数字1的左面
  3. 一般情况下使用clickhouse都是集群,所以对分布式表进行 删除数据、更新数据时,需要使用特定语法 alter {分布式表名} on cluster {集群名} delete where /update where

3.3 与传统方案对比

分别统计了1天、一个月、半年相同查询条件进行统计对比


mysql 方案Mysql java 方案clickhouse 方案
一天 70w数据5秒12分钟350毫秒
一个月 2100w数据2分钟3.7个小时2秒以内
半年1.2亿数据时间过长,超过了20分钟。时间过长,超过12个小时影响了统计的时效2秒以内

可以看到mysql方案中计算放到mysql中在数据量较小时,还是可以接受,但是一旦数据过大,会影响到整个数据库的响应,触发系统监控。另外在整个过程中,无法控制,无法查看进度,存在阻塞业务风险。

Mysql-java 方案中,整体计算流程是在java中进行的,所以可以控制进度,随时查看统计任务进度,并且对数据库压力也小,都是简单筛选查询。但问题是统计时间过长。

Clickhouse 方案,即使1.2亿数据也能在2秒内响应,统计时间上非常理想。

4 总结与展望

Clickhouse 从设计和实现上都更加适合统计业务的场景,提供了非常好的查询性能同时借助于其列式存储结构与高压缩比的设计,单机处理数据量也非常庞大。另外,调用上,提供jdbc链接方式和完善的SQL支持,上手十分简单。

但 clickhouse 本身也有一些不足:

  1. 并发能力不足,每次查询都使用多核并行运算,数据文件io也是经过压缩,所以查询时消耗CPU,并发不足
  2. 不支持事务

未来,除了关键词命中统计外,还将继续发掘类似的结构化数据的统计需求,诸如: AI 文字检测统计、历史人工审核记录统计分析等,更好的为业务方提供数据支持和多维度的数据比对结果。