ClickHouseInc

ClickHouse 原生采样:千万级数据聚合,秒级响应不是梦!

图片

本文字数:9861;估计阅读时间:25 分钟

作者:Mark Needham

Image

有时,对所有数据运行聚合查询的速度可能不够快。ClickHouse 原生的随机采样(random sampling)功能允许你转而对数据的一个随机子集或样本执行查询——只要配置得当,查询结果依然能保持惊人的准确性。

查看媒体链接:https://youtu.be/RHkfT7Uj38k

Image
设置

我们将使用 英国房价数据集(https://clickhouse.com/docs/getting-started/example-datasets/uk-price-paid),其中包含超过 3000 万条房产交易记录。首先,我们来创建这张表:

CREATE TABLE uk_price_paid

(

    price UInt32,

    date Date,

    postcode1 LowCardinality(String),

    postcode2 LowCardinality(String),

    type Enum8('terraced' = 1, 'semi-detached' = 2, 'detached' = 3, 'flat' = 4, 'other' = 0),

    is_new UInt8,

    duration Enum8('freehold' = 1, 'leasehold' = 2, 'unknown' = 0),

    addr1 String,

    addr2 String,

    street LowCardinality(String),

    locality LowCardinality(String),

    town LowCardinality(String),

    district LowCardinality(String),

    county LowCardinality(String)

)

ENGINE = MergeTree

ORDER BY (postcode1, postcode2, addr1, addr2);

然后导入数据:

INSERT INTO uk_price_paid

SELECT

    toUInt32(price_string) AS price,

    parseDateTimeBestEffortUS(time) AS date,

    splitByChar(' ', postcode)[1] AS postcode1,

    splitByChar(' ', postcode)[2] AS postcode2,

    transform(a, ['T', 'S', 'D', 'F', 'O'], ['terraced', 'semi-detached', 'detached', 'flat', 'other']) AS type,

    b = 'Y' AS is_new,

    transform(c, ['F', 'L', 'U'], ['freehold', 'leasehold', 'unknown']) AS duration,

    addr1,

    addr2,

    street,

    locality,

    town,

    district,

    county

FROM url(

    'http://prod1.publicdata.landregistry.gov.uk.s3-website-eu-west-1.amazonaws.com/pp-complete.csv',

    'CSV',

    'uuid_string String,

    price_string String,

    time String,

    postcode String,

    a String,

    b String,

    c String,

    addr1 String,

    addr2 String,

    street String,

    locality String,

    town String,

    district String,

    county String,

    d String,

    e String'

) SETTINGS max_http_get_redirects=10;

接下来,我们可以编写以下查询来统计表中的所有行数:

SELECT count()

FROM uk_price_paid;

查询结果如下:

┌──count()─┐

│ 30452463 │ -- 30.45 million

└──────────┘

1 row in set. Elapsed: 0.002 sec.

Image
选择采样键

现在,假设我们想对这份数据的样本运行聚合查询,因为直接对所有数据运行这些查询效率不高。为此,我们需要在创建表时指定一个采样键(sample key)。

采样键是一个表达式,用于确定在数据采样时哪些行将被包含。采样键应为无符号整数(例如哈希值),具有高基数(high cardinality),并且在其取值范围内均匀分布。sipHash64 就是一个这样的函数,它是一个返回 64 位值的加密哈希函数。

我们可以通过检查 sipHash64 对 postcode1 和 postcode2 的分布均匀性来对此进行验证。以下查询使用 sipHash64 对这两个邮编字段进行哈希,然后将哈希值向下取整并分配到 10 个桶中,每个桶预期应包含约 10% 的数据。

WITH pow(2, 64) AS maxUInt64, 10 AS numBuckets

SELECT

    floor(sipHash64(postcode1, postcode2) / (maxUInt64 / numBuckets)) AS bucket,

    count() AS rows,

    round(count() / (SELECT count() FROM uk_price_paid) * 100, 2) AS pct

FROM uk_price_paid

GROUP BY bucket

ORDER BY bucket;

运行该查询,我们将看到以下结果:

┌─bucket─┬────rows─┬───pct─┐

│      0 │ 3106578 │  10.2 │

│      1 │ 3025622 │  9.94 │

│      2 │ 3043133 │  9.99 │

│      3 │ 3033438 │  9.96 │

│      4 │ 3063534 │ 10.06 │

│      5 │ 3048884 │ 10.01 │

│      6 │ 3005960 │  9.87 │

│      7 │ 3053378 │ 10.03 │

│      8 │ 3029517 │  9.95 │

│      9 │ 3042419 │  9.99 │

└────────┴─────────┴───────┘

每个桶都包含大约 10% 的数据——虽然略有波动,但总体分布非常接近 10%。因此,postcode1 和 postcode2 将是采样键的良好选择。

另一方面,如果我们转而对 county 字段进行哈希,结果会如何呢?

WITH pow(2, 64) AS maxUInt64, 10 AS numBuckets

SELECT

    floor(sipHash64(county) / (maxUInt64 / numBuckets)) AS bucket,

    count() AS rows,

    round(count() / (SELECT count() FROM uk_price_paid) * 100, 2) AS pct

FROM uk_price_paid

GROUP BY bucket

ORDER BY bucket;

查询结果如下:

┌─bucket─┬────rows─┬───pct─┐

│      0 │ 2853723 │  9.37 │

│      1 │ 2543953 │  8.35 │

│      2 │ 4568680 │    15 │

│      3 │ 1402093 │   4.6 │

│      4 │ 5547516 │ 18.22 │

│      5 │ 1402865 │  4.61 │

│      6 │ 4854564 │ 15.94 │

│      7 │ 2820334 │  9.26 │

│      8 │ 2649867 │   8.7 │

│      9 │ 1808868 │  5.94 │

└────────┴─────────┴───────┘

这次,每个桶中的数据量范围从 4.6% 到 18%,这表明数据分布不够均匀。因此,county 对于我们的采样键来说不是一个理想的选择。其原因在于 county 是一个低基数(low cardinality)列,仅有 132 个唯一值。因此,我们计算了 132 个哈希值,并将它们分散到 10 个桶中,平均每个桶约包含 13 个县的数据。每个桶中的行数取决于相应县的房产销售数量。

以下查询展示了按县划分的房产销售分布不均的程度:

SELECT county,

       round((100 * count()) / sum(count()) OVER (), 2) AS pct

FROM uk_price_paid

GROUP BY county

ORDER BY pct DESC

LIMIT 10

查询结果如下:

┌─county─────────────┬──pct─┐

│ GREATER LONDON     │ 12.7 │

│ GREATER MANCHESTER │ 4.45 │

│ WEST MIDLANDS      │  3.8 │

│ WEST YORKSHIRE     │  3.8 │

│ KENT               │ 2.84 │

│ ESSEX              │ 2.77 │

│ HAMPSHIRE          │  2.6 │

│ LANCASHIRE         │ 2.29 │

│ SURREY             │ 2.23 │

│ MERSEYSIDE         │ 2.13 │

└────────────────────┴──────┘

大伦敦区贡献了近13%的总销售额,因此无论它被分到哪个数据桶,都会造成严重的数据倾斜。

相比之下,有132万个不同的 postcode1/postcode2 组合,这足以让哈希函数(hash function)发挥作用,实现数据的均匀分布。

Image
创建表

接下来,我们使用采样键创建一个新版本的表。SAMPLE BY 键必须是主键(primary key)的一部分(即 ORDER BY 表达式的一部分)。我们会在 ORDER BY 语句的开头添加 sipHash64(postcode1, postcode2),以便最大限度地利用采样功能。

CREATE TABLE uk_price_paid_sample

(

    price UInt32,

    date Date,

    postcode1 LowCardinality(String),

    postcode2 LowCardinality(String),

    type Enum8('terraced' = 1, 'semi-detached' = 2, 'detached' = 3, 'flat' = 4, 'other' = 0),

    is_new UInt8,

    duration Enum8('freehold' = 1, 'leasehold' = 2, 'unknown' = 0),

    addr1 String,

    addr2 String,

    street LowCardinality(String),

    locality LowCardinality(String),

    town LowCardinality(String),

    district LowCardinality(String),

    county LowCardinality(String)

)

ENGINE = MergeTree

ORDER BY (sipHash64(postcode1, postcode2), postcode1, postcode2, addr1, addr2)

SAMPLE BY sipHash64(postcode1, postcode2);

现在,我们将原始表中的数据插入到新表中:

INSERT INTO uk_price_paid_sample

SELECT * 

FROM uk_price_paid;

Image
表内数据采样

表配置完毕后,是时候进行数据采样了!数据采样是确定性的(deterministic),这意味着相同的 SELECT .. SAMPLE 查询将始终返回相同的结果。SAMPLE 子句(clause)有两种工作模式。

第一种模式是按比例(by fraction),我们提供一个介于0到1之间的值,查询将针对该比例的数据执行。

SELECT count() 

FROM uk_price_paid_sample 

SAMPLE 0.1;

此查询提取哈希值范围的前10%,实际上等同于所有满足 sipHash64(postcode1, postcode2) < 2^64 * 0.1 条件的行。执行该查询的结果如下所示:

┌─count()─┐

│ 3106578 │ -- 3.11 million

└─────────┘

1 row in set. Elapsed: 0.043 sec.

该表共有3045万行数据,因此10%约为304.5万行,这与我们实际返回的311万行数据相当接近。

第二种模式是按行数(by row count),我们指定要返回的最小行数:

SELECT count() 

FROM uk_price_paid_sample 

SAMPLE 100000;

这将返回至少100,000行。ClickHouse 会计算一个阈值 (100000 / estimated_table_rows) * 2^64,并选择哈希值低于该阈值的行。执行该查询的输出如下所示:

┌─count()─┐

│  103347 │

└─────────┘

1 row in set. Elapsed: 0.009 sec.

在这两种情况下,由于哈希值是主索引(primary index)的一部分,ClickHouse 可以快速跳过哈希值高于阈值的 granules。

Image
_sample_factor 虚拟列

在使用采样时,我们还可以利用 _sample_factor 虚拟列,它表示每个采样行在完整数据集中所代表的原始行数。如果没有进行采样,_sample_factor 的值为1。对于10%的采样,其值为10。而对于按行数采样,该值大约为304.5。

如果我们想估算总行数,可以将 _sample_factor 的值进行求和,如下面的查询所示:

SELECT 'None' AS sampleType, count(), sum(_sample_factor)

FROM uk_price_paid_sample

UNION ALL

SELECT 'Fraction (0.1)' AS sampleType, count(), sum(_sample_factor)

FROM uk_price_paid_sample SAMPLE 0.1

UNION ALL

SELECT 'Number (100,000)' AS sampleType, count(), sum(_sample_factor)

FROM uk_price_paid_sample SAMPLE 100000;

查询结果如下所示:

┌─sampleType───────┬──count()─┬─sum(_sample_factor)─┐

│ None             │ 30452463 │            30452463 │

│ Fraction (0.1)   │  3106578 │            31065780 │

│ Number (100,000) │   103347 │  31471706.936610293 │

└──────────────────┴──────────┴─────────────────────┘

这两种采样技术都将总行数高估了 2-3%,误差在可接受范围内。

同样的原理也适用于其他聚合操作。如果我们对 10% 的样本运行 sum(price),我们将大致得到实际总数的 10%:

SELECT sum(price) 

FROM uk_price_paid_sample 

SAMPLE 0.1;

查询结果如下所示:

┌───sum(price)─┐

│ 741594301263 │ -- 741.59 billion

└──────────────┘

为了获得全数据集的估算结果,我们需要乘以 _sample_factor:

SELECT sum(price * _sample_factor) 

FROM uk_price_paid_sample 

SAMPLE 0.1;

查询结果如下所示:

┌─sum(multiply⋯le_factor))─┐

│            7415943012630 │ -- 7.42 trillion

└──────────────────────────┘

avg 和 count 聚合在样本上无需调整即可安全使用——对子集进行平均或计数会得到与完整数据集相同的结果。但任何依赖于总体总和的聚合操作(例如 sum),都需要 _sample_factor 进行缩放以还原到全数据集的量级。

需要注意的是,采样操作总是会存在一定的误差——这是我们为换取更快的查询时间而做出的权衡。

Image
采样一个真实分析查询

现在,我们将其应用到一个真实的分析查询中。对于给定年份的每个县 (county)、城镇 (town) 和房产类型 (property type),我们希望获取其总销售额、平均价格和中位数。我们限制每个县和城镇只显示一行结果,否则将只会得到大伦敦 (Greater London) 的数据。

SELECT county, town, type,

       toYear(date) AS year,

       sum(_sample_factor) AS sales,

       round(avg(price)) AS avgPrice,

       round(quantile(0.5)(price)) AS median

FROM uk_price_paid_sample

GROUP BY county, town, type, year

ORDER BY sales DESC

LIMIT 1 BY county, town

LIMIT 5;

查询结果如下所示:

┌─county─────────────┬─town───────┬─type─────┬─year─┬─sales─┬─avgPrice─┬─median─┐

│ GREATER LONDON     │ LONDON     │ flat     │ 2006 │ 65567 │   296793 │ 242500 │

│ GREATER MANCHESTER │ MANCHESTER │ terraced │ 2004 │  9896 │    80212 │  71000 │

│ WEST MIDLANDS      │ BIRMINGHAM │ terraced │ 2002 │  8266 │    73946 │  67500 │

│ MERSEYSIDE         │ LIVERPOOL  │ terraced │ 2003 │  6507 │    52895 │  42500 │

│ WEST YORKSHIRE     │ LEEDS      │ terraced │ 2002 │  5921 │    64886 │  55000 │

└────────────────────┴────────────┴──────────┴──────┴───────┴──────────┴────────┘

我对此进行了多次运行,其耗时如下所示:

5 rows in set. Elapsed: 0.687 sec. Processed 30.45 million rows, 290.09 MB (44.36 million rows/s., 422.55 MB/s.)

Peak memory usage: 481.82 MiB.

5 rows in set. Elapsed: 0.721 sec. Processed 30.45 million rows, 290.09 MB (42.22 million rows/s., 402.20 MB/s.)

Peak memory usage: 496.21 MiB.

5 rows in set. Elapsed: 0.748 sec. Processed 30.45 million rows, 290.09 MB (40.69 million rows/s., 387.62 MB/s.)

Peak memory usage: 500.59 MiB.

现在,让我们看看对 10% 的数据进行采样会发生什么:

SELECT county, town, type,

       toYear(date) AS year,

       sum(_sample_factor) AS sales,

       round(avg(price)) AS avgPrice,

       round(quantile(0.5)(price)) AS median

FROM uk_price_paid_sample SAMPLE 0.1

GROUP BY county, town, type, year

ORDER BY sales DESC

LIMIT 1 BY county, town

LIMIT 5;

查询结果如下所示:

┌─county─────────────┬─town───────┬─type─────┬─year─┬─sales─┬─avgPrice─┬─median─┐

│ GREATER LONDON     │ LONDON     │ flat     │ 2006 │ 70500 │   287644 │ 240000 │

│ GREATER MANCHESTER │ MANCHESTER │ terraced │ 2004 │ 10380 │    75479 │  66000 │

│ WEST MIDLANDS      │ BIRMINGHAM │ terraced │ 2002 │  8480 │    73869 │  68000 │

│ MERSEYSIDE         │ LIVERPOOL  │ terraced │ 2004 │  6640 │    78769 │  68000 │

│ WEST YORKSHIRE     │ LEEDS      │ terraced │ 2002 │  6320 │    64598 │  54050 │

└────────────────────┴────────────┴──────────┴──────┴───────┴──────────┴────────┘

同样,我们将其运行三次:

5 rows in set. Elapsed: 0.169 sec. Processed 2.91 million rows, 38.31 MB (17.23 million rows/s., 226.97 MB/s.)

Peak memory usage: 238.06 MiB.

5 rows in set. Elapsed: 0.141 sec. Processed 3.15 million rows, 41.49 MB (22.30 million rows/s., 294.09 MB/s.)

Peak memory usage: 177.00 MiB.

5 rows in set. Elapsed: 0.185 sec. Processed 2.12 million rows, 27.82 MB (11.45 million rows/s., 150.14 MB/s.)

Peak memory usage: 179.40 MiB.

取每次运行的最佳时间,采样查询耗时 141 毫秒,而查询整个数据集耗时 687 毫秒,性能提升了约 80%。

采样查询的查询计划如下所示:

┌─explain──────────────────────────────────────────────────────┐

 1. │ Expression (Project names)                                   │

 2. │   Limit                                                      │

 3. │     LimitBy                                                  │

 4. │       Expression ((Before LIMIT BY + (Before ORDER BY + Proj⋯│

 5. │         Sorting (Sorting for ORDER BY)                       │

 6. │           Expression ((Before ORDER BY + Projection))        │

 7. │             Aggregating                                      │

 8. │               Expression ((Before GROUP BY + Change column n⋯│

 9. │                 ReadFromMergeTree (default.uk_price_paid_sam⋯│

10. │                 Indexes:                                     │

11. │                   PrimaryKey                                 │

12. │                     Keys:                                    │

13. │                       sipHash64(postcode1, postcode2)        │

14. │                     Condition: (sipHash64(postcode1, postcod⋯│

15. │                     Parts: 9/9                               │

16. │                     Granules: 384/3719                       │

17. │                     Search Algorithm: binary search          │

18. │                   Ranges: 9                                  │

    └──────────────────────────────────────────────────────────────┘

从第 16 行可以看出,该查询仅扫描了 3719 个数据粒 (granules) 中的 384 个,这正是查询时间缩短的原因。

从结果来看,销售额数据大约有 7-8% 的误差,平均价格更为接近,中位数也有轻微偏差。对于探索性工作而言,这可能是一个合理的权衡。

Image
总结

当需要快速、近似的答案且能容忍少量误差时,采样技术非常有用。如果本文您只需记住三点:选择一个高基数 (high-cardinality) 列作为采样键 (sampling key);将采样键置于 ORDER BY 子句的首位以实现最快的查询;以及使用 _sample_factor 对求和聚合进行缩放。

/END/

征稿启示

面向社区长期正文,文章内容包括但不限于关于 ClickHouse 的技术研究、项目实践和创新做法等。建议行文风格干货输出&图文并茂。质量合格的文章将会发布在本公众号,优秀者也有机会推荐到 ClickHouse 官网。请将文章稿件的 WORD 版本发邮件至:[email protected]

关于我们

ClickHouse 是全球速度最快,资源利用最高效的在线分析列式数据库管理系统。现在,ClickHouse可以作为一个安全可扩展的无服务器应用在云中提供服务。通过云服务,ClickHouse使得任何人都能轻松获取高效的实时分析处理能力。2023年,ClickHouse正式进入中国,请访问clickhouse.com以获取更多信息。

图片