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

> Kafka Connect と ClickHouse で HTTP Sink コネクタを使用する

# Confluent HTTP Sink コネクタ

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

HTTP Sink コネクタはデータ型に依存しないため、Kafka のスキーマを必要とせず、Maps や Arrays などの ClickHouse 固有のデータ型にも対応しています。この柔軟性がある一方で、設定はやや複雑になります。

以下では、単一の Kafka トピックからメッセージを取り込み、ClickHouse テーブルに行を挿入するシンプルな導入方法について説明します。

<Note>
  HTTP コネクタは [Confluent Enterprise License](https://docs.confluent.io/kafka-connect-http/current/overview.html#license) の下で配布されています。
</Note>

<div id="quick-start-steps">
  ### クイックスタートの手順
</div>

<Steps>
  <Step title="接続の詳細を収集する" id="1-gather-your-connection-details">
    HTTP(S) で ClickHouse に接続するには、次の情報が必要です。

    | Parameter(s)              | Description                                               |
    | ------------------------- | --------------------------------------------------------- |
    | `HOST` and `PORT`         | 通常、TLS を使用する場合のポートは 8443、TLS を使用しない場合は 8123 です。           |
    | `DATABASE NAME`           | デフォルトでは `default` という名前のデータベースがあります。接続先のデータベース名を使用してください。 |
    | `USERNAME` and `PASSWORD` | デフォルトのユーザー名は `default` です。用途に応じたユーザー名を使用してください。           |

    ClickHouse Cloud サービスの詳細は、ClickHouse Cloud コンソールで確認できます。
    サービスを選択し、**Connect** をクリックします。

    <div className="ch-image-md">
      <Frame>
        <img src="https://mintcdn.com/private-7c7dfe99/CFFsa2agBPbviR4r/images/_snippets/cloud-connect-button.webp?fit=max&auto=format&n=CFFsa2agBPbviR4r&q=85&s=ec0a298a33ca841e947fa5e8bae47362" alt="ClickHouse Cloud サービスの接続ボタン" width="998" height="932" data-path="images/_snippets/cloud-connect-button.webp" />
      </Frame>
    </div>

    **HTTPS** を選択します。接続情報は `curl` コマンドの例として表示されます。

    <div className="ch-image-md">
      <Frame>
        <img src="https://mintcdn.com/private-7c7dfe99/CFFsa2agBPbviR4r/images/_snippets/connection-details-https.webp?fit=max&auto=format&n=CFFsa2agBPbviR4r&q=85&s=cb0fbd98aa2b5b7ca484c9f53395ee07" alt="ClickHouse Cloud HTTPS 接続情報" width="1320" height="1184" data-path="images/_snippets/connection-details-https.webp" />
      </Frame>
    </div>

    セルフマネージド ClickHouse を使用している場合、接続情報は ClickHouse 管理者によって設定されます。
  </Step>

  <Step title="Kafka Connect と HTTP sink コネクタを実行する" id="2-run-kafka-connect-and-the-http-sink-connector">
    2 つの選択肢があります。

    * **セルフマネージド:** Confluent パッケージをダウンロードしてローカルにインストールします。コネクタのインストールについては、[こちら](https://docs.confluent.io/kafka-connect-http/current/overview.html)に記載されている手順に従ってください。
      confluent-hub によるインストール方法を使用する場合、ローカルの設定ファイルが更新されます。

    * **Confluent Cloud:** Kafka のホスティングに Confluent Cloud を使用している場合は、完全マネージド型の HTTP Sink を利用できます。この場合、ClickHouse 環境が Confluent Cloud からアクセス可能である必要があります。

    <Note>
      以下の例では Confluent Cloud を使用しています。
    </Note>
  </Step>

  <Step title="ClickHouse に宛先テーブルを作成する" id="3-create-destination-table-in-clickhouse">
    接続テストの前に、まず ClickHouse Cloud にテスト用のテーブルを作成します。このテーブルが Kafka からのデータを受信します。

    ```sql theme={null}
    CREATE TABLE default.my_table
    (
        `side` String,
        `quantity` Int32,
        `symbol` String,
        `price` Int32,
        `account` String,
        `userid` String
    )
    ORDER BY tuple()
    ```
  </Step>

  <Step title="HTTP Sink の設定" id="4-configure-http-sink">
    Kafka トピックと HTTP Sink コネクタのインスタンスを作成します。

    <Image img="https://mintcdn.com/private-7c7dfe99/pIetLsS_hOGHqoPJ/images/integrations/data-ingestion/kafka/confluent/create_http_sink.webp?fit=max&auto=format&n=pIetLsS_hOGHqoPJ&q=85&s=ade794478790411f406291f68f1f39a1" size="sm" alt="HTTP Sink コネクタの作成方法を示す Confluent Cloud のインターフェイス" border width="718" height="736" data-path="images/integrations/data-ingestion/kafka/confluent/create_http_sink.webp" />

    <br />

    HTTP Sink コネクタを設定します。

    * 作成したトピック名を指定します
    * 認証
      * `HTTP Url` - `INSERT` クエリを指定した ClickHouse Cloud の URL `<protocol>://<clickhouse_host>:<clickhouse_port>?query=INSERT%20INTO%20<database>.<table>%20FORMAT%20JSONEachRow`。**注**: クエリはエンコードする必要があります。
      * `Endpoint Authentication type` - BASIC
      * `Auth username` - ClickHouse のユーザー名
      * `Auth password` - ClickHouse のパスワード

    <Note>
      この HTTP Url は誤りが生じやすいため、問題を避けるにはエスケープを正確に行ってください。
    </Note>

    <Image img="https://mintcdn.com/private-7c7dfe99/pIetLsS_hOGHqoPJ/images/integrations/data-ingestion/kafka/confluent/http_auth.webp?fit=max&auto=format&n=pIetLsS_hOGHqoPJ&q=85&s=cd3fc19e3bd604dd88cff4b487fdddb7" size="lg" alt="HTTP Sink コネクタの認証設定を示す Confluent Cloud のインターフェイス" border width="1944" height="878" data-path="images/integrations/data-ingestion/kafka/confluent/http_auth.webp" />

    <br />

    * 設定
      * `Input Kafka record value format`ソースデータによって異なりますが、多くの場合は JSON または Avro です。以下の設定では `JSON` を前提とします。
      * `advanced configurations` セクション内:
        * `HTTP Request Method` - POST に設定します
        * `Request Body Format` - json
        * `Batch batch size` - ClickHouse の推奨に従い、**少なくとも 1000** に設定します。
        * `Batch json as array` - true
        * `Retry on HTTP codes` - 400-500。必要に応じて調整してください。たとえば、ClickHouse の前段に HTTP プロキシがある場合は変更が必要になることがあります。
        * `Maximum Reties` - デフォルトの (10) で適切ですが、より堅牢に再試行したい場合は調整してもかまいません。

    <Image img="https://mintcdn.com/private-7c7dfe99/pIetLsS_hOGHqoPJ/images/integrations/data-ingestion/kafka/confluent/http_advanced.webp?fit=max&auto=format&n=pIetLsS_hOGHqoPJ&q=85&s=ec905e552defb0fe75a70dde123ecb24" size="sm" alt="HTTP Sink コネクタの詳細設定オプションを示す Confluent Cloud のインターフェイス" border width="786" height="796" data-path="images/integrations/data-ingestion/kafka/confluent/http_advanced.webp" />
  </Step>

  <Step title="接続性のテスト" id="5-testing-the-connectivity">
    HTTP Sink で設定したトピックにメッセージを作成します

    <Image img="https://mintcdn.com/private-7c7dfe99/pIetLsS_hOGHqoPJ/images/integrations/data-ingestion/kafka/confluent/create_message_in_topic.webp?fit=max&auto=format&n=pIetLsS_hOGHqoPJ&q=85&s=60c5ffc663c30145519f583b0cbed9b8" size="md" alt="Kafka トピックでテストメッセージを作成する方法を示す Confluent Cloud のインターフェイス" border width="2138" height="846" data-path="images/integrations/data-ingestion/kafka/confluent/create_message_in_topic.webp" />

    <br />

    続いて、作成したメッセージが ClickHouse インスタンスに書き込まれていることを確認します。
  </Step>
</Steps>

<div id="troubleshooting">
  ### トラブルシューティング
</div>

<div id="http-sink-doesnt-batch-messages">
  #### HTTP Sink がメッセージをバッチ化しない
</div>

[Sink のドキュメント](https://docs.confluent.io/kafka-connectors/http/current/overview.html#http-sink-connector-for-cp)より:

> Kafka ヘッダー値が異なるメッセージを含む場合、HTTP Sink コネクタはリクエストをバッチ化しません。

1. Kafka レコードのキーが同じであることを確認してください。
2. HTTP API の URL にパラメータを追加すると、レコードごとに一意の URL になることがあります。そのため、追加の URL パラメータを使用するとバッチ化は無効になります。

<div id="400-bad-request">
  #### 400 Bad Request
</div>

<div id="cannot_parse_quoted_string">
  ##### CANNOT\_PARSE\_QUOTED\_STRING
</div>

JSONオブジェクトを `String` カラムに挿入する際に、HTTP Sink が次のメッセージを出して失敗する場合:

```response theme={null}
Code: 26. DB::ParsingException: Cannot parse JSON string: expected opening quote: (while reading the value of key key_name): While executing JSONEachRowRowInputFormat: (at row 1). (CANNOT_PARSE_QUOTED_STRING)
```

URL 内で、設定 `input_format_json_read_objects_as_strings=1` を URL エンコードされた文字列 `SETTINGS%20input_format_json_read_objects_as_strings%3D1` として指定します

<div id="load-the-github-dataset-optional">
  ### GitHub データセットを読み込む (任意)
</div>

この例では、GitHub データセットの Array フィールドを保持したまま扱います。サンプルでは、空の github トピックがあることを前提とし、Kafka へのメッセージの挿入に [kcat](https://github.com/edenhill/kcat) を使用します。

<Steps>
  <Step title="設定を準備する" id="1-prepare-configuration">
    インストール形態に応じた Connect のセットアップについては、スタンドアロンと分散クラスターの違いに注意しつつ、[こちらの手順](https://docs.confluent.io/cloud/current/cp-component/connect-cloud-config.html#set-up-a-local-connect-worker-with-cp-install)に従ってください。Confluent Cloud を使用する場合は、分散セットアップが該当します。

    最も重要なパラメータは `http.api.url` です。ClickHouse の [HTTP インターフェイス](/docs/ja/concepts/features/interfaces/http) では、INSERT ステートメントを URL のパラメータとしてエンコードする必要があります。これには、フォーマット (この場合は `JSONEachRow`) と移行先データベースを含める必要があります。フォーマットは Kafka のデータと一致している必要があり、そのデータは HTTP ペイロード内で文字列に変換されます。これらのパラメータは URL エスケープする必要があります。GitHub データセットに対するこのフォーマットの例 (ClickHouse をローカルで実行していることを前提) は、以下のとおりです。

    ```response theme={null}
    <protocol>://<clickhouse_host>:<clickhouse_port>?query=INSERT%20INTO%20<database>.<table>%20FORMAT%20JSONEachRow

    http://localhost:8123?query=INSERT%20INTO%20default.github%20FORMAT%20JSONEachRow
    ```

    ClickHouse で HTTP Sink を使用する際は、以下の追加パラメータが関係します。完全なパラメータ一覧は[こちら](https://docs.confluent.io/kafka-connect-http/current/connector_config.html)で確認できます。

    * `request.method` - **POST** に設定します
    * `retry.on.status.codes` - 任意のエラーコードで再試行するには 400-500 に設定します。データ内で想定されるエラーに応じて調整してください。
    * `request.body.format` - ほとんどの場合、JSON になります。
    * `auth.type` - ClickHouse で認証を使用する場合は BASIC に設定します。現在のところ、ClickHouse と互換性のある他の認証方式はサポートされていません。
    * `ssl.enabled` - SSL を使用する場合は true に設定します。
    * `connection.user` - ClickHouse のユーザー名。
    * `connection.password` - ClickHouse のパスワード。
    * `batch.max.size` - 1 回の batch で送信する行数です。十分に大きな値を設定してください。ClickHouse の[推奨事項](/docs/ja/reference/statements/insert-into#performance-considerations)によると、1000 は最低値と考えるべきです。
    * `tasks.max` - HTTP Sink コネクタは 1 つ以上のタスクの実行をサポートしています。これはパフォーマンス向上に利用できます。batch size とあわせて、パフォーマンス改善の主要な手段となります。
    * `key.converter` - キーの型に応じて設定します。
    * `value.converter` - topic 上のデータ型に基づいて設定します。このデータにスキーマは不要です。ここでのフォーマットは、パラメータ `http.api.url` で指定する FORMAT と一致している必要があります。最も簡単なのは、JSON と org.apache.kafka.connect.json.JsonConverter コンバータを使用する方法です。org.apache.kafka.connect.storage.StringConverter コンバータを使って値を文字列として扱うことも可能ですが、その場合は INSERT ステートメント内で関数を使って値を抽出する必要があります。io.confluent.connect.avro.AvroConverter コンバータを使用する場合、ClickHouse は [Avro format](/docs/ja/reference/formats/Avro/Avro) もサポートしています。

    proxy、再試行、高度な SSL の設定方法を含む設定の完全な一覧は、[こちら](https://docs.confluent.io/kafka-connect-http/current/connector_config.html)で確認できます。

    GitHub サンプルデータ用の設定ファイル例は[こちら](https://github.com/ClickHouse/clickhouse-docs/tree/main/docs/integrations/data-ingestion/kafka/code/connectors/http_sink)にあります。これは、Connect がスタンドアロン モードで実行され、Kafka が Confluent Cloud でホストされていることを前提としています。
  </Step>

  <Step title="ClickHouse テーブルを作成する" id="2-create-the-clickhouse-table">
    テーブルが作成されていることを確認してください。標準的な MergeTree を使用する最小構成の GitHub データセットの例を以下に示します。

    ```sql theme={null}
    CREATE TABLE github
    (
        file_time DateTime,
        event_type Enum('CommitCommentEvent' = 1, 'CreateEvent' = 2, 'DeleteEvent' = 3, 'ForkEvent' = 4,'GollumEvent' = 5, 'IssueCommentEvent' = 6, 'IssuesEvent' = 7, 'MemberEvent' = 8, 'PublicEvent' = 9, 'PullRequestEvent' = 10, 'PullRequestReviewCommentEvent' = 11, 'PushEvent' = 12, 'ReleaseEvent' = 13, 'SponsorshipEvent' = 14, 'WatchEvent' = 15, 'GistEvent' = 16, 'FollowEvent' = 17, 'DownloadEvent' = 18, 'PullRequestReviewEvent' = 19, 'ForkApplyEvent' = 20, 'Event' = 21, 'TeamAddEvent' = 22),
        actor_login LowCardinality(String),
        repo_name LowCardinality(String),
        created_at DateTime,
        updated_at DateTime,
        action Enum('none' = 0, 'created' = 1, 'added' = 2, 'edited' = 3, 'deleted' = 4, 'opened' = 5, 'closed' = 6, 'reopened' = 7, 'assigned' = 8, 'unassigned' = 9, 'labeled' = 10, 'unlabeled' = 11, 'review_requested' = 12, 'review_request_removed' = 13, 'synchronize' = 14, 'started' = 15, 'published' = 16, 'update' = 17, 'create' = 18, 'fork' = 19, 'merged' = 20),
        comment_id UInt64,
        path String,
        ref LowCardinality(String),
        ref_type Enum('none' = 0, 'branch' = 1, 'tag' = 2, 'repository' = 3, 'unknown' = 4),
        creator_user_login LowCardinality(String),
        number UInt32,
        title String,
        labels Array(LowCardinality(String)),
        state Enum('none' = 0, 'open' = 1, 'closed' = 2),
        assignee LowCardinality(String),
        assignees Array(LowCardinality(String)),
        closed_at DateTime,
        merged_at DateTime,
        merge_commit_sha String,
        requested_reviewers Array(LowCardinality(String)),
        merged_by LowCardinality(String),
        review_comments UInt32,
        member_login LowCardinality(String)
    ) ENGINE = MergeTree ORDER BY (event_type, repo_name, created_at)

    ```
  </Step>

  <Step title="Kafka にデータを投入する" id="3-add-data-to-kafka">
    Kafka にメッセージを書き込みます。以下では、[kcat](https://github.com/edenhill/kcat) を使って 10k 件のメッセージを書き込みます。

    ```bash theme={null}
    head -n 10000 github_all_columns.ndjson | kcat -b <host>:<port> -X security.protocol=sasl_ssl -X sasl.mechanisms=PLAIN -X sasl.username=<username>  -X sasl.password=<password> -t github
    ```

    ターゲットテーブル "Github" を簡単に参照すれば、データが挿入されたことを確認できます。

    ```sql theme={null}
    SELECT count() FROM default.github;

    | count\(\) |
    | :--- |
    | 10000 |

    ```
  </Step>
</Steps>
