2 万亿行数据不停机迁移:ClickHouse + ClickPipes 的一次真实大规模实践
本文字数:6463;估计阅读时间:17 分钟
作者:Lionel Palacin
Meetup活动
ClickHouse 深圳第三届 Meetup 讲师招募中,欢迎讲师在文末扫码报名!
ClickPy 是一个免费的分析型 Python 下载统计平台,最近达成了一个令人瞩目的里程碑:其主数据集已经包含超过 2 万亿行记录。数据集中的每一行都代表一次 Python 包下载,历史数据最早可以追溯到 2011 年。
SELECT count() FROM pypi.pypi┌───────count()─┐│ 2214017475506 │ -- 2.21 trillion└───────────────┘
ClickHouse 是支撑 ClickPy 的分析型数据库,使 Python 爱好者能够查看自己喜爱的 Python 包的流行度,或洞察 Python 生态系统中的新趋势。
在几乎无需持续运维投入的情况下达到如此规模,充分体现了 ClickHouse 作为高吞吐量分析数据长期存储方案的优势。数据行数突破 2 万亿,也促使我们重新审视现有的数据摄取管道,并对其进行加固和优化。
在这一规模下对数据摄取管道进行演进并非易事。系统必须始终保持可用,数据摄取不能中断,同时所有变更都需要非常谨慎地引入。在此过程中,我们也对历史数据进行了更深入的检查,从而发现了一些数据层面的不一致问题。
本文接下来的内容将介绍我们如何重新设计数据摄取管道,以及在整个过程中保持 ClickPy 持续运行的前提下,如何修复这些历史数据问题。
最初的摄取管道构建于 ClickPipes 尚未出现的时期,因此依赖了大量自定义脚本。我们曾在另一篇博客文章中详细介绍过这套原始的数据摄取流程。
PyPI 的下载统计数据以 BigQuery 表的形式对外公开,我们运行了一个自动化的 BigQuery 作业,每天将数据导出到 Google Cloud Storage 存储桶中。随后,一个自定义脚本会每日执行一次 ClickHouse insert 语句,将数据直接从 GCS 存储桶摄取到 ClickPy 的主表 pypi 中。
这种方案在过去运行良好,但在当前规模下仍有明显的改进空间。最大的优化点在于,用 ClickPipes 替换图中标注为 ClickLoad 的自定义摄取脚本。
迁移到 ClickPipes 后,可以消除由基于 cron 的脚本所带来的多项运维限制:
重试、退避以及失败处理机制由系统内建提供,而无需自行实现
摄取管道的状态和进度清晰可见,更容易发现停滞或不完整的摄取任务
将摄取管道的日常维护交由 ClickHouse 团队负责
调整摄取逻辑时涉及的组件更少,自定义代码量也更低
在引入 ClickPipes 之后,数据摄取成为系统中的核心组成部分,而不再是一个需要额外监控和维护的外部流程。
用 ClickPipes 替换自定义脚本的最大挑战在于当前的运行规模。我们无法承受任何数据损坏的风险。主表 pypi 已包含超过 2 万亿行数据,同时我们还基于 pypi 表构建了多个物化视图,这使得从头重新摄取全部数据成为不可行的方案。
为了在降低风险的前提下,用 ClickPipes 替换原有的自定义摄取脚本,我们采用了分阶段推进的策略。
第一步是将新的摄取管道与生产环境完全隔离。我们把 pypi 数据库中的所有表结构和物化视图克隆到了一个名为 pypi_clickpipes 的独立数据库中,从而可以在不影响现有查询和仪表盘的情况下,对摄取、转换和聚合逻辑进行验证。
-- Create the databaseCREATE DATABASE pypi_clickpipes;-- Clone the table schemasCREATE TABLE pypi_clickpipes.pypi AS pypi.pypi;CREATE TABLE pypi_clickpipes.pypi_downloads AS pypi.pypi_downloads;CREATE TABLE pypi_clickpipes.pypi_downloads_by_version AS pypi.pypi_downloads_by_version;CREATE TABLE pypi_clickpipes.pypi_downloads_max_min AS pypi.pypi_downloads_max_min;CREATE TABLE pypi_clickpipes.pypi_downloads_per_day AS pypi.pypi_downloads_per_day;CREATE TABLE pypi_clickpipes.pypi_downloads_per_day_by_installer AS pypi.pypi_downloads_per_day_by_installer;CREATE TABLE pypi_clickpipes.pypi_downloads_per_day_by_version AS pypi.pypi_downloads_per_day_by_version;CREATE TABLE pypi_clickpipes.pypi_downloads_per_day_by_version_by_country AS pypi.pypi_downloads_per_day_by_version_by_country;CREATE TABLE pypi_clickpipes.pypi_downloads_per_day_by_version_by_file_type AS pypi.pypi_downloads_per_day_by_version_by_file_type;CREATE TABLE pypi_clickpipes.pypi_downloads_per_day_by_version_by_installer_by_type AS pypi.pypi_downloads_per_day_by_version_by_installer_by_type;CREATE TABLE pypi_clickpipes.pypi_downloads_per_day_by_version_by_installer_by_type_by_country AS pypi.pypi_downloads_per_day_by_version_by_installer_by_type_by_country;CREATE TABLE pypi_clickpipes.pypi_downloads_per_day_by_version_by_python AS pypi.pypi_downloads_per_day_by_version_by_python;CREATE TABLE pypi_clickpipes.pypi_downloads_per_day_by_version_by_python_by_country AS pypi.pypi_downloads_per_day_by_version_by_python_by_country;CREATE TABLE pypi_clickpipes.pypi_downloads_per_day_by_version_by_system AS pypi.pypi_downloads_per_day_by_version_by_system;CREATE TABLE pypi_clickpipes.pypi_downloads_per_day_by_version_by_system_by_country;CREATE TABLE pypi_clickpipes.pypi_downloads_per_month AS pypi.pypi_downloads_per_day_by_version_by_system_by_country;
在克隆表准备就绪后,我们配置 ClickPipes 将数据从 GCS 摄取到新的目标表 pypi_raw。ClickPipes 会在首次运行时自动创建该表并推断表结构。这个表仅作为一个短暂存在的暂存层。
在具体配置中,我们刻意为目标表选择了 Null 引擎。由于并不需要在 ClickHouse 中长期保存原始行数据,实际的数据转换由物化视图完成,并直接写入包含主数据集的 pypi 表。
新的物化视图承载了此前由自定义脚本实现的全部转换逻辑,包括字段规范化、类型转换以及模式对齐。将这些逻辑内聚到 ClickHouse 中,使其在后续维护和演进时更加清晰可控。
CREATE MATERIALIZED VIEW pypi_clickpipes.pypi_mv TO pypi.pypi(`date` Date,`country_code` LowCardinality(String),`project` String,`type` LowCardinality(String),`installer` LowCardinality(String),`python_minor` LowCardinality(String),`system` LowCardinality(String),`version` String,`ci` Enum8('false' = 0, 'true' = 1, 'unknown' = 2),`filename` String,`libc` Tuple(lib LowCardinality(String),version LowCardinality(String)))AS SELECTtoDate(timestamp) AS date,ifNull(country_code, '') AS country_code,ifNull(project, '') AS project,ifNull(tupleElement(file, 'type'), '') AS type,ifNull(tupleElement(installer, 'name'), '') AS installer,arrayStringConcat(arraySlice(splitByChar('.', ifNull(python, '')), 1, 2), '.') AS python_minor,ifNull(tupleElement(system, 'name'), '') AS system,ifNull(tupleElement(file, 'version'), '') AS version,CAST(ifNull(ci, 0), 'Enum8(\'false\' = 0, \'true\' = 1, \'unknown\' = 2)') AS ci,ifNull(tupleElement(file, 'filename'), '') AS filename,tuple(ifNull(tupleElement(tupleElement(distro, 'libc'), 'lib'), ''), ifNull(tupleElement(tupleElement(distro, 'libc'), 'version'), '')) AS libcFROM pypi_clickpipes.pypi_raw
下面展示了使用 ClickPipes 从 BigQuery 流向 ClickHouse 的整体数据流示意。
在重构数据摄取管道的过程中,我们也顺势为数据集新增了一些字段,用于响应社区近期提出的功能需求 [1,2]。
当暂存表和物化视图就位后,ClickPipes 的配置过程相对简单。ClickPipes 支持持续摄取,并能够跟踪已处理的文件,但并不支持从指定的历史时间点开始摄取。为了避免重新导入完整的历史数据,我们调整了 BigQuery 的导出作业,将新增数据写入一个新的 GCS 存储桶。随后将 ClickPipes 指向该新位置,使新的摄取管道从迁移时刻开始生效。
我们让新的摄取管道并行运行了数天,并持续对比生产数据库 pypi 与暂存数据库 pypi_clickpipes 中的每日行数。当数据完全一致后,切换过程就非常顺利:我们停用了运行自定义脚本的 cron 作业,并将物化视图的写入目标更新为 pypi.pypi。
完成切换后,我们清理了 pypi_clickpipes 数据库中的克隆表,仅保留 pypi_raw、pypi_mv 以及由 ClickPipes 自动创建和管理的表。这样形成了清晰的职责划分:pypi 数据库专注于为 ClickPy 提供查询服务,而 pypi_clickpipes 则专门用于数据摄取与转换。
在验证新摄取管道的过程中,我们发现 BigQuery(数据源)与 ClickHouse 中的数据存在差异。通过对比每日行数,我们注意到 ClickHouse 中某些历史日期的数据并不完整。
在如此规模的数据系统中,这类问题很容易被忽视。查询依然能够返回结果,仪表盘也能正常展示,但数据本身却已经悄然出现偏差。一旦确认问题存在,随之而来的挑战是:如何在不中断现有摄取流程、也不重建完整数据集的前提下修复这些历史数据。
修复历史数据最直接的方式,是删除受影响日期的数据,并从源头重新摄取。然而,列式数据库天生更适合高效写入和分析查询,而非对已有数据进行修改,ClickHouse 亦是如此。好在 ClickHouse 近期引入了轻量级的删除与更新能力,使得即便是在包含数万亿行数据的表上,也能安全地删除大范围的数据。
针对单日数据的修复,我们采用了如下流程:
从所有按天聚合的表中删除对应日期的数据。
从主表 pypi 中删除该日期的数据。
临时删除不按天聚合的物化视图。
将该日期的源数据重新摄取到 pypi 表中。
重新创建不按天聚合的物化视图。
重建由这些物化视图生成的表。
区分按天聚合与非按天聚合的表至关重要。对于按天分组的表(例如 pypi_downloads_per_day),我们可以直接删除某一天的数据,并依赖物化视图在重新摄取时自动恢复。
而对于不按天分组的表(例如按月聚合的表),则无法精确定位属于某一天的数据行。在这些情况下,必须在重新摄取前先移除相关物化视图,以避免产生重复数据;待历史数据回填完成后,再重新创建并重建这些视图。
用于上述修复流程的完整脚本可以在这里查看(https://github.com/ClickHouse/clickpy/blob/main/scripts/day-fix.sh)。
ClickPy 达到 2 万亿行不仅是一个数字上的成就,它既体现了 Python 生态系统日益增长的活跃度,也验证了 ClickHouse 在这一规模下支撑分析型负载的能力。我们将 ClickPy 作为一个免费的开源服务提供,是因为我们享受在大规模数据集之上构建应用,希望回馈开源社区,并且得益于 ClickHouse,即使在数万亿行的数据规模下也能实现高性价比的运行。
本文介绍的工作,包括将摄取管道迁移到 ClickPipes 以及修复历史数据,都是我们持续努力的一部分,旨在随着 ClickPy 的不断成长,保持平台的健康性和长期价值。社区的参与在塑造 ClickPy 的过程中至关重要,我们鼓励大家提交功能需求和问题反馈,因为它们会直接推动平台的改进。近期的一个例子便是通过 Metabase 导出图表的功能支持。
在完成这些改造之后,ClickPy 的数据更加准确,运维更加简单,也已经为未来的持续增长做好了准备。
我们正为深圳活动招募讲师,如果你有独特的技术见解、实践经验或 ClickHouse 使用故事,非常欢迎你加入我们,成为这次活动的讲师,与大家分享你的经验。
/END/
试用阿里云 ClickHouse企业版
轻松节省30%云资源成本?阿里云数据库ClickHouse 云原生架构全新升级,首次购买ClickHouse企业版计算和存储资源组合,首月消费不超过99.58元(包含最大16CCU+450G OSS用量)了解详情:https://t.aliyun.com/Kz5Z0q9G
征稿启示
面向社区长期正文,文章内容包括但不限于关于 ClickHouse 的技术研究、项目实践和创新做法等。建议行文风格干货输出&图文并茂。质量合格的文章将会发布在本公众号,优秀者也有机会推荐到 ClickHouse 官网。请将文章稿件的 WORD 版本发邮件至:[email protected]