为什么要将数据从 SQL Server 流式传输到 ClickHouse?
- 不会拖慢生产应用的内部报表
- 需要响应迅速且始终保持最新的面向客户的仪表盘
- 事件流处理,例如让用户活动日志持续保持最新状态,以便进行分析
开始前需要准备的内容
前置条件
- 一个正在运行的 SQL Server 实例
- 本教程使用的是 AWS RDS for SQL Server,但任何现代 SQL Server 实例都可以。从零开始设置 AWS SQL Server。
- 一个 ClickHouse 实例
- 可使用自托管或 Cloud。从零开始设置 ClickHouse。
- Streamkap
- 该工具将作为数据流式管道的核心。
连接信息
- SQL Server 服务器地址、端口、用户名和密码。建议为 Streamkap 创建单独的用户和角色,以访问你的 SQL Server 数据库。查看我们的配置文档。
- ClickHouse 服务器地址、端口、用户名和密码。ClickHouse 中的 IP 访问列表决定了哪些服务可以连接到你的 ClickHouse 数据库。请按照此处的说明操作。
- 你想要流式传输的表——暂时先从一张表开始
1
在 Streamkap 中创建 SQL Server 数据源
开始吧!我们先设置 source 连接,这样 Streamkap 才知道要从哪里拉取变更。操作步骤如下:
- 打开 Streamkap,进入 Sources 部分。
- 创建一个新的 source。
- 给它起一个便于识别的名称 (例如:sqlserver-demo-source) 。
- 填写 SQL Server 连接详情:
- Host (例如:your-db-instance.rds.amazonaws.com)
- Port (SQL Server 的默认端口是 3306)
- 用户名和密码
- 数据库名称
后台会发生什么
完成设置后,Streamkap 会连接到你的 SQL Server 并检测表。此次演示中,我们会选择一张已经有一些数据持续流入的表,例如 events 或 transactions。2
在 Streamkap 中添加 ClickHouse 目标端
现在,我们来配置接收这些数据的目标端。和源端一样,我们将使用 ClickHouse 连接详情来创建目标端。
步骤:
- 前往 Streamkap 中的 destinations 部分。
- 添加一个新的目标端——选择 ClickHouse 作为目标端类型。
- 输入你的 ClickHouse 信息:
- Host
- Port (默认为 9000)
- 用户名和密码
- 数据库名称
Upsert 模式:这是什么?
这是一个重要步骤:我们要使用 ClickHouse 的 “upsert” 模式——它在底层使用 ClickHouse 的 ReplacingMergeTree engine。这样可以高效合并传入记录,并通过 ClickHouse 所说的“分片合并”在摄取后处理更新。- 这样可以确保当 SQL Server 端发生变化时,目标端表不会堆满重复数据。
处理 Schema 演进
ClickHouse 和 SQL Server 有时并不会有完全相同的列——尤其是在应用已经上线、开发人员还在持续动态添加列的情况下。- 好消息是:Streamkap 可以处理基础的 schema 演进。这意味着,如果你在 SQL Server 中新增了一列,它也会出现在 ClickHouse 这一侧。
3
在 Streamkap 中配置管道
源端和目标端都已设置好,现在到了最有意思的部分——开始流式传输数据!
管道设置
- 前往 Streamkap 中的 Pipelines 选项卡。
- 创建一个新管道。
- 选择你的 SQL Server 源 (sqlserver-demo-source) 。
- 选择你的 ClickHouse 目标端 (clickhouse-tutorial-destination) 。
- 选择你要流式传输的表——这里假设是 events。
- 配置 Change Data Capture (CDC) 。
- 本次我们将流式传输新数据 (刚开始可以先跳过回填,重点关注 CDC 事件) 。
需要回填吗?
你可能会问:需要回填旧数据吗?在很多分析场景中,你可能只想从现在开始流式传输变更;不过之后你也随时可以回过头来加载历史数据。除非你有明确需求,否则目前选择“不要回填”即可。4
查看数据流
现在你的管道已配置完成并开始运行!以下是具体过程:在高负载场景下,可能会有一定延迟,但大多数用例都能实现接近实时的流式传输。
- 当新数据进入 SQL Server 中的源表时,Streamkap 管道会捕获变更并将其发送到 ClickHouse。
- ClickHouse (借助 ReplacingMergeTree 和分片合并) 会摄取这些行,并合并更新。
- schema 也能保持同步——在 SQL Server 中新增列后,这些列也会出现在 ClickHouse 中。
底层原理:Streamkap 实际在做什么?
- Streamkap 会监控 SQL Server 的二进制日志 (也就是用于复制的那种日志) 。
- 只要你的表中有行被插入、更新或删除,Streamkap 就会捕获这个事件。
- 它会将该事件转换为 ClickHouse 能理解的内容并发送过去——把变更立即应用到你的分析数据库中。
高级选项
Upsert 与 Insert 模式
- Insert Mode:每个新行都会被添加——即使它其实是一次更新,也会产生重复数据。
- Upsert Mode:对现有行的更新会覆盖原有内容——更适合让分析数据保持最新且整洁。
处理 schema 变更
- 给业务表新增一列? Streamkap 会自动识别,并在 ClickHouse 端同步添加这一列。
- 删除一列? 这取决于具体设置,你可能需要执行迁移——不过大多数新增都能顺利处理。
生产环境监控:随时掌握管道运行状态
检查管道运行状况
- 查看管道滞后情况 (你的数据有多新鲜?)
- 监控行数和吞吐量
- 在出现异常时接收告警
需要重点关注的常见指标
- 滞后:ClickHouse 比 SQL Server 落后多少?
- 吞吐量:每秒处理的行数
- 错误率:应接近零
开始使用:查询 ClickHouse
后续步骤与深入研究
- 设置带过滤的流 (仅同步部分表/列)
- 将多个数据源流式传输到同一个分析型数据库
- 将其与 S3/数据湖结合用于冷存储
- 在更改表时自动执行 schema 迁移
- 使用 SSL 和防火墙规则保护你的管道
常见问题与故障排查
总结
- Upsert 与 Insert 的区别,以及两者的技术细节
- 端到端延迟:多久才能获得最终的分析视图?
- 性能调优与吞吐量
- 基于这套技术栈构建的真实场景仪表盘