> ## 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 (Bring Your Own Cloud) 形態の完全マネージド型サービスとして提供されています。

Streamkap を使用すると、PostgreSQL、MySQL、SQL Server、MongoDB、<a href="https://streamkap.com/connectors" target="_blank">そのほか多数</a>のソースデータベースで発生するあらゆる insert、update、delete を、ミリ秒単位のレイテンシで直接 ClickHouse にストリーミングできます。

そのため、リアルタイム分析ダッシュボード、オペレーショナルアナリティクス、機械学習モデルへのライブデータ供給に最適です。

<div id="key-features">
  ## 主な機能
</div>

* **リアルタイムストリーミング CDC:** Streamkap はデータベースのログから変更を直接取り込み、ClickHouse 内のデータがソースのリアルタイムなレプリカであることを保証します。
  シンプルなストリーム処理: ClickHouse に取り込む前に、データをリアルタイムで変換、エンリッチ、ルーティング、フォーマットし、embeddings を作成できます。複雑さを伴わない Flink を基盤としています

* **完全マネージド型でスケーラブル:** 本番環境に対応した、メンテナンス不要のパイプラインを提供し、Kafka、Flink、Debezium、またはスキーマレジストリのインフラストラクチャを自前で管理する必要をなくします。このプラットフォームは高スループット向けに設計されており、数十億件のイベントを処理できるようリニアにスケールします。

* **自動スキーマ進化:** Streamkap はソースデータベースのスキーマ変更を自動的に検出し、それらを ClickHouse に反映します。手動で介入しなくても、新しいカラムの追加やカラム型の変更に対応できます。

* **ClickHouse 向けに最適化:** このインテグレーションは、ClickHouse の機能を効率的に活用できるように構築されています。既定では、ReplacingMergeTree エンジンを使用して、ソースシステムからの更新や削除をシームレスに処理します。

* **耐障害性の高い配信:** このプラットフォームは at-least-once 配信保証を提供し、ソースと 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 クラスターの接続情報: ホスト名、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. endpoint、port、database name、ユーザー認証情報などの接続情報を入力します。
    5. コネクタを保存します。
  </Step>

  <Step title="ClickHouse の宛先を設定する" id="configure-clickhouse-dest">
    1. **Connectors** セクションで、**Destinations** タブを選択します。
    2. **+ Add** をクリックし、一覧から **ClickHouse** を選択します。
    3. ClickHouse サービスの接続情報を入力します。
       * **ホスト名:** ClickHouse インスタンスのホスト名 (例: `abc123.us-west-2.aws.clickhouse.cloud`)
       * **Port:** 通常は `8443` のセキュアな HTTPS ポート
       * **Username and Password:** ClickHouse ユーザーの認証情報
       * **Database:** ClickHouse の対象データベース名
    4. 宛先を保存します。
  </Step>

  <Step title="パイプラインを作成して実行する" id="run-pipeline">
    1. サイドバーの **Pipelines** に移動し、**+ Create** をクリックします。
    2. 先ほど設定した ログソース と 宛先 を選択します。
    3. ストリーミングするスキーマとテーブルを選択します。
    4. パイプラインに名前を付けて **Save** をクリックします。

    作成すると、パイプラインはアクティブになります。Streamkap はまず既存データのスナップショットを取得し、その後、新しい変更が発生するとストリーミングを開始します。
  </Step>

  <Step title="ClickHouse でデータを確認する" id="verify-data-clickhoouse">
    ClickHouse クラスターに接続し、クエリを実行してターゲットテーブルにデータが取り込まれていることを確認します。

    ```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 では新しい行として書き込まれます。バックグラウンドのマージ処理で、ReplacingMergeTree はこれらの行をまとめ、ソートキーに基づいて最新バージョンだけを残します。

* **削除** は、ReplacingMergeTree の `is_deleted` パラメータに渡されるメタデータフラグによって処理されます。ソースで削除された行はすぐには削除されず、削除済みとしてマークされます。
  * 必要に応じて、削除済みレコードを分析目的で ClickHouse に保持できます

<div id="metadata-columns">
  ### メタデータカラム
</div>

Streamkap は、データの状態を管理するために各テーブルに複数のメタデータカラムを追加します。

| カラム名                      | 説明                                               |
| ------------------------- | ------------------------------------------------ |
| `_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 viewを使用できます。

<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>
