跳转到主要内容
Streamkap 是一个实时数据集成平台,专注于流式 CDC (变更数据捕获) 和流处理。它基于 Apache Kafka、Apache Flink 和 Debezium 构建,具备高吞吐和可扩展性,并以 SaaS 或 BYOC (自备 Cloud) 部署模式提供全托管服务。 Streamkap 支持将 PostgreSQL、MySQL、SQL Server、MongoDB 等源数据库中的每一次插入、更新和删除直接实时流式传输到 ClickHouse,延迟可低至毫秒级。支持的数据源还有更多 这使其非常适合为实时分析仪表盘和运营分析提供支持,并为机器学习模型提供实时数据。

主要特性

  • 实时流式 CDC: Streamkap 直接从数据库日志中捕获变更,确保 ClickHouse 中的数据是源端数据的实时副本。 简化的流处理:在数据写入 ClickHouse 之前,实时完成转换、增强、路由、格式化以及创建嵌入向量。由 Flink 驱动,却无需承担其复杂性
  • 全托管且可扩展: 它提供可用于生产环境的零维护管道,无需自行管理 Kafka、Flink、Debezium 或 Schema Registry 基础设施。该平台专为高吞吐而设计,并且可以线性扩展以处理数十亿事件。
  • 自动 schema 演进: Streamkap 会自动检测源数据库中的 schema 变更,并将其同步到 ClickHouse。它可以在无需人工干预的情况下处理新增列或列类型变更。
  • 针对 ClickHouse 优化: 该集成专为高效利用 ClickHouse 特性而构建。默认情况下,它使用 ReplacingMergeTree 引擎,无缝处理来自源系统的更新和删除。
  • 可靠交付: 该平台提供至少一次投递保证,确保源端与 ClickHouse 之间的数据一致性。对于 upsert 操作,它会基于主键执行去重。

入门

本指南简要介绍如何设置 Streamkap 管道,将数据导入 ClickHouse。

前置条件

  • 一个 Streamkap 账户
  • 你的 ClickHouse 集群连接信息:主机名、端口、用户名和密码。
  • 一个已配置为允许 CDC 的源数据库 (例如 PostgreSQL、SQL Server) 。你可以在 Streamkap 文档中找到详细的设置指南。

第 1 步:在 Streamkap 中配置源

  1. 登录你的 Streamkap 账户。
  2. 在侧边栏中,前往 Connectors,然后选择 Sources 选项卡。
  3. 点击 + Add,然后选择源 database 类型 (例如 SQL Server RDS) 。
  4. 填写连接信息,包括端点、端口、database 名称以及用户凭据。
  5. 保存该 connector。

第 2 步:配置 ClickHouse 目标端

  1. Connectors 部分,选择 Destinations 选项卡。
  2. 点击 + Add,然后从列表中选择 ClickHouse
  3. 输入你的 ClickHouse 服务连接信息:
    • Hostname: 你的 ClickHouse 实例主机名 (例如 abc123.us-west-2.aws.clickhouse.cloud)
    • Port: 安全的 HTTPS 端口,通常为 8443
    • Username and Password: 你的 ClickHouse 用户名和密码
    • Database: ClickHouse 中目标数据库的名称
  4. 保存目标端。

第 3 步:创建并运行管道

  1. 在侧边栏中前往 Pipelines,然后点击 + 创建
  2. 选择您刚刚配置的 源 和目标端。
  3. 选择您要进行流式传输的 schema 和表。
  4. 为管道命名,然后点击 保存
创建后,管道将立即激活。Streamkap 会先对现有数据进行一次快照,然后开始流式传输后续产生的新变更。

第 4 步:在 ClickHouse 中验证数据

连接到你的 ClickHouse 集群,并运行查询以查看写入目标表的数据。

与 ClickHouse 的工作方式

Streamkap 的集成专为在 ClickHouse 中高效处理 CDC (变更数据捕获) 数据而设计。

表引擎与数据处理

默认情况下,Streamkap 使用 upsert 摄取模式。当它在 ClickHouse 中创建表时,会使用 ReplacingMergeTree 引擎。该引擎非常适合处理 CDC 事件:
  • 源表的主键会在 ReplacingMergeTree 表定义中用作 ORDER BY 键。
  • 源端的更新会作为新行写入 ClickHouse。在后台 merge 过程中,ReplacingMergeTree 会合并这些行,并仅保留基于排序键的最新版本。
  • 删除通过一个元数据标志来处理,该标志会传递给 ReplacingMergeTree 的 is_deleted 参数。源端已删除的行不会立即移除,而是会被标记为已删除。
    • 也可以选择将已删除的记录保留在 ClickHouse 中,以用于分析

元数据列

Streamkap 会为每个表添加几个元数据列,用于管理数据状态:

查询最新数据

由于 ReplacingMergeTree 会在后台处理更新和删除,因此在合并完成前,简单的 SELECT * 查询可能会显示历史行或已删除的行。要获取数据的最新状态,必须过滤掉已删除的记录,并只选择每一行的最新版本。 你可以使用 FINAL 修饰符来实现这一点。这样做很方便,但可能会影响查询性能:
为了提升大型表的查询性能,尤其是在无需读取所有列且只进行一次性分析查询时,可以使用 argMax 函数为每个主键手动选出最新记录:
对于生产环境以及需要并发处理终端用户周期性查询的场景,可以使用 Materialized Views 对数据进行建模,以更好地适应下游访问模式。

延伸阅读

最后修改于 2026年7月2日