> ## Documentation Index
> Fetch the complete documentation index at: https://clickhouse.com/docs/llms.txt
> Use this file to discover all available pages before exploring further.

> 使用 Airbyte 数据管道将流数据传输到 ClickHouse

# 将 Streamkap 连接到 ClickHouse

export const PartnerBadge = () => {
  return <div className="PartnerBadge">
            <div className="PartnerBadgeIcon">
                <svg width="16" height="16" viewBox="0 0 16 16" fill="none" xmlns="http://www.w3.org/2000/svg">
                    <polyline points="12.5 9.5 10 12 6 11 2.5 8.5" stroke="currentColor" strokeLinecap="round" strokeLinejoin="round" strokeWidth="1" />
                    <polyline points="4.54 4.41 8 3.5 11.46 4.41" stroke="currentColor" strokeLinecap="round" strokeLinejoin="round" strokeWidth="1" />
                    <path d="M2.15,3.78 L0.55,6.95 A0.5,0.5 0,0,0 0.77,7.62 L2.5,8.5 L4.54,4.41 L2.82,3.55 A0.5,0.5 0,0,0 2.15,3.78 Z" stroke="currentColor" strokeLinecap="round" strokeLinejoin="round" strokeWidth="1" />
                    <path d="M13.5,8.5 L15.23,7.62 A0.5,0.5 0,0,0 15.45,6.95 L13.85,3.78 A0.5,0.5 0,0,0 13.18,3.55 L11.46,4.41 Z" stroke="currentColor" strokeLinecap="round" strokeLinejoin="round" strokeWidth="1" />
                    <path d="M11.5,4.5 L9,4.5 L6.15,7.27 A0.5,0.5 0,0,0 6.24,8.05 C7.33,8.74 8.81,8.72 10,7.5 L12.5,9.5 L13.5,8.5" stroke="currentColor" strokeLinecap="round" strokeLinejoin="round" strokeWidth="1" />
                    <polyline points="7.75 13.5 5.15 12.85 3.5 11.67" stroke="currentColor" strokeLinecap="round" strokeLinejoin="round" strokeWidth="1" />
                </svg>
            </div>
            合作伙伴集成
        </div>;
};

export const Image = ({img, alt, size = "lg"}) => {
  const normalizedSize = ["sm", "md", "lg"].includes(size) ? size : "lg";
  return <div className={`ch-image-${normalizedSize}`}>
      <Frame>
        <img src={img} alt={alt} />
      </Frame>
    </div>;
};

<PartnerBadge />

<a href="https://streamkap.com/" target="_blank">Streamkap</a> 是一个实时数据集成平台，专注于流式 CDC (变更数据捕获) 和流处理。它基于 Apache Kafka、Apache Flink 和 Debezium 构建，具备高吞吐和可扩展性，并以 SaaS 或 BYOC (自备 Cloud) 部署模式提供全托管服务。

Streamkap 支持将 PostgreSQL、MySQL、SQL Server、MongoDB 等源数据库中的每一次插入、更新和删除直接实时流式传输到 ClickHouse，延迟可低至毫秒级。支持的数据源还有<a href="https://streamkap.com/connectors" target="_blank">更多</a>。

这使其非常适合为实时分析仪表盘和运营分析提供支持，并为机器学习模型提供实时数据。

<div id="key-features">
  ## 主要特性
</div>

* **实时流式 CDC：** Streamkap 直接从数据库日志中捕获变更，确保 ClickHouse 中的数据是源端数据的实时副本。
  简化的流处理：在数据写入 ClickHouse 之前，实时完成转换、增强、路由、格式化以及创建嵌入向量。由 Flink 驱动，却无需承担其复杂性

* **全托管且可扩展：** 它提供可用于生产环境的零维护管道，无需自行管理 Kafka、Flink、Debezium 或 Schema Registry 基础设施。该平台专为高吞吐而设计，并且可以线性扩展以处理数十亿事件。

* **自动 schema 演进：** Streamkap 会自动检测源数据库中的 schema 变更，并将其同步到 ClickHouse。它可以在无需人工干预的情况下处理新增列或列类型变更。

* **针对 ClickHouse 优化：** 该集成专为高效利用 ClickHouse 特性而构建。默认情况下，它使用 ReplacingMergeTree 引擎，无缝处理来自源系统的更新和删除。

* **可靠交付：** 该平台提供至少一次投递保证，确保源端与 ClickHouse 之间的数据一致性。对于 upsert 操作，它会基于主键执行去重。

<div id="started">
  ## 入门
</div>

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

<div id="prerequisites">
  ### 前置条件
</div>

* 一个 <a href="https://app.streamkap.com/account/sign-up" target="_blank">Streamkap 账户</a>。
* 你的 ClickHouse cluster 连接信息：Hostname、Port、Username 和 Password。
* 一个已配置为允许 CDC 的源数据库 (例如 PostgreSQL、SQL Server) 。你可以在 Streamkap 文档中找到详细的设置指南。

<Steps>
  <Step title="在 Streamkap 中配置源" id="configure-clickhouse-source">
    1. 登录你的 Streamkap 账户。
    2. 在侧边栏中，前往 **Connectors** 并选择 **Sources** 选项卡。
    3. 点击 **+ Add**，然后选择你的源数据库类型 (例如 SQL Server RDS) 。
    4. 填写连接信息，包括端点、端口、数据库名称和用户凭据。
    5. 保存 connector。
  </Step>

  <Step title="配置 ClickHouse 目标端" id="configure-clickhouse-dest">
    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. 保存目标端。
  </Step>

  <Step title="创建并运行管道" id="run-pipeline">
    1. 在侧边栏中前往 **Pipelines**，然后点击 **+ Create**。
    2. 选择你刚刚配置的源和目标端。
    3. 选择你希望进行流式传输的 schema 和表。
    4. 为你的管道命名，然后点击 **Save**。

    创建后，管道将变为活动状态。Streamkap 会先对现有数据执行一次 snapshot，然后开始流式传输后续发生的所有新变更。
  </Step>

  <Step title="在 ClickHouse 中验证数据" id="verify-data-clickhoouse">
    连接到你的 ClickHouse cluster，并运行查询以查看写入目标表的数据。

    ```sql theme={null}
    SELECT * FROM your_table_name LIMIT 10;
    ```
  </Step>
</Steps>

<div id="how-it-works-with-clickhouse">
  ## 与 ClickHouse 的工作方式
</div>

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

<div id="table-engine-data-handling">
  ### 表引擎与数据处理
</div>

默认情况下，Streamkap 使用 upsert 摄取模式。当它在 ClickHouse 中创建表时，会使用 ReplacingMergeTree 引擎。该引擎非常适合处理 CDC 事件：

* 源表的主键会在 ReplacingMergeTree 表定义中用作 ORDER BY 键。

* 源端的**更新**会作为新行写入 ClickHouse。在后台 merge 过程中，ReplacingMergeTree 会合并这些行，并仅保留基于排序键的最新版本。

* **删除**通过一个元数据标志来处理，该标志会传递给 ReplacingMergeTree 的 `is_deleted` 参数。源端已删除的行不会立即移除，而是会被标记为已删除。
  * 也可以选择将已删除的记录保留在 ClickHouse 中，以用于分析

<div id="metadata-columns">
  ### 元数据列
</div>

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

| Column Name               | Description                          |
| ------------------------- | ------------------------------------ |
| `_STREAMKAP_SOURCE_TS_MS` | 源数据库中事件的时间戳 (以毫秒为单位) 。               |
| `_STREAMKAP_TS_MS`        | Streamkap 处理该事件时的时间戳 (以毫秒为单位) 。      |
| `__DELETED`               | 布尔标志 (`true`/`false`) ，表示该行是否已在源端删除。 |
| `_STREAMKAP_OFFSET`       | 来自 Streamkap 内部日志的偏移量，可用于排序和调试。      |

<div id="query-latest-data">
  ### 查询最新数据
</div>

由于 ReplacingMergeTree 会在后台处理更新和删除，因此在合并完成前，简单的 SELECT \* 查询可能会显示历史行或已删除的行。要获取数据的最新状态，必须过滤掉已删除的记录，并只选择每一行的最新版本。

你可以使用 FINAL 修饰符来实现这一点。这样做很方便，但可能会影响查询性能：

```sql theme={null}
-- 使用 FINAL 获取正确的当前状态
SELECT * FROM your_table_name FINAL WHERE __DELETED = 'false';
SELECT * FROM your_table_name FINAL LIMIT 10;
SELECT * FROM your_table_name FINAL WHERE <filter by keys in ORDER BY clause>;
SELECT count(*) FROM your_table_name FINAL;
```

为了提升大型表的查询性能，尤其是在无需读取所有列且只进行一次性分析查询时，可以使用 argMax 函数为每个主键手动选出最新记录：

```sql theme={null}
SELECT key,
       argMax(col1, version) AS col1,
       argMax(col2, version) AS col2
FROM t
WHERE <您的过滤条件>
GROUP BY key;
```

对于生产环境以及需要并发处理终端用户周期性查询的场景，可以使用 Materialized Views 对数据进行建模，以更好地适应下游访问模式。

<div id="further-reading">
  ## 延伸阅读
</div>

* <a href="https://streamkap.com/" target="_blank">Streamkap 网站</a>
* <a href="https://docs.streamkap.com/clickhouse" target="_blank">适用于 ClickHouse 的 Streamkap 文档</a>
* <a href="https://streamkap.com/blog/streaming-with-change-data-capture-to-clickhouse" target="_blank">博客：通过变更数据捕获将数据流式传输到 ClickHouse</a>
* <a href="https://streamkap.com/blog/streaming-with-change-data-capture-to-clickhouse" target="_blank">ClickHouse 文档：ReplacingMergeTree</a>
