> ## Documentation Index
> Fetch the complete documentation index at: https://private-7c7dfe99-mintlify-67bc7bf8.mintlify.site/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}) => {
  return <Frame>
      <img src={img} alt={alt} />
    </Frame>;
};

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>

<div id="1-gather-your-connection-details">
  #### 1. 接続情報を確認する
</div>

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** をクリックします。

<Image img="https://mintcdn.com/private-7c7dfe99-mintlify-67bc7bf8/bY1iopO46snwcr5S/images/_snippets/cloud-connect-button.png?fit=max&auto=format&n=bY1iopO46snwcr5S&q=85&s=fda686b9b9454ad21df18bb7bac60ff5" size="md" alt="ClickHouse Cloud サービスの接続ボタン" border width="998" height="932" data-path="images/_snippets/cloud-connect-button.png" />

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

<Image img="https://mintcdn.com/private-7c7dfe99-mintlify-67bc7bf8/bY1iopO46snwcr5S/images/_snippets/connection-details-https.png?fit=max&auto=format&n=bY1iopO46snwcr5S&q=85&s=a0ea22430cbcef486eb6f9f596b18a46" size="md" alt="ClickHouse Cloud HTTPS 接続情報" border width="1320" height="1184" data-path="images/_snippets/connection-details-https.png" />

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

<div id="2-run-kafka-connect-and-the-http-sink-connector">
  #### 2. Kafka Connect と HTTP Sink コネクタを実行する
</div>

選択肢は 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>

<div id="3-create-destination-table-in-clickhouse">
  #### 3. ClickHouse に宛先テーブルを作成する
</div>

接続テストの前に、まず 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()
```

<div id="4-configure-http-sink">
  #### 4. HTTP Sink の設定
</div>

Kafka トピックと HTTP Sink コネクタのインスタンスを作成します。

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

<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-mintlify-67bc7bf8/ms4E-Ni592HZjtjR/images/integrations/data-ingestion/kafka/confluent/http_auth.png?fit=max&auto=format&n=ms4E-Ni592HZjtjR&q=85&s=5269df8ce2fd188238ba1eafff7be798" size="lg" alt="HTTP Sink コネクタの認証設定を示す Confluent Cloud のインターフェイス" border width="1944" height="878" data-path="images/integrations/data-ingestion/kafka/confluent/http_auth.png" />

<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-mintlify-67bc7bf8/ms4E-Ni592HZjtjR/images/integrations/data-ingestion/kafka/confluent/http_advanced.png?fit=max&auto=format&n=ms4E-Ni592HZjtjR&q=85&s=d569c5a1141227c703c439fb9ac727d4" size="sm" alt="HTTP Sink コネクタの詳細設定オプションを示す Confluent Cloud のインターフェイス" border width="786" height="796" data-path="images/integrations/data-ingestion/kafka/confluent/http_advanced.png" />

<div id="5-testing-the-connectivity">
  #### 5. 接続のテスト
</div>

HTTP Sink で設定したトピックにメッセージを作成します

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

<br />

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

<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) を使用します。

<div id="1-prepare-configuration">
  ##### 1. 設定を準備する
</div>

インストール形態に応じた 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 インターフェイス](/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 の[推奨事項](/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](/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 でホストされていることを前提としています。

<div id="2-create-the-clickhouse-table">
  ##### 2. ClickHouseテーブルを作成する
</div>

テーブルが作成済みであることを確認してください。標準的な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)

```

<div id="3-add-data-to-kafka">
  ##### 3. Kafka にデータを追加する
</div>

Kafka にメッセージを投入します。以下では、[kcat](https://github.com/edenhill/kcat) を使用して 1 万件のメッセージを投入します。

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

```
