> ## 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 および ClickHouse で Vector を使用する

# Kafka および ClickHouse で Vector を使用する

<div id="using-vector-with-kafka-and-clickhouse">
  ## Kafka と ClickHouse で Vector を使用する
</div>

Vector は、Kafka から読み取ったイベントを ClickHouse に送信できる、ベンダー非依存のデータパイプラインです。

ClickHouse 向け Vector の[入門ガイド](/docs/ja/integrations/connectors/data-ingestion/etl-tools/vector-to-clickhouse)では、ログのユースケースとファイルからのイベント読み取りに重点を置いています。ここでは、Kafka トピックに保持されたイベントを含む[Github sample dataset](https://datasets-documentation.s3.eu-west-3.amazonaws.com/kafka/github_all_columns.ndjson)を使用します。

Vector は、プッシュまたはプルのモデルでデータを取得するために[ログソース](https://vector.dev/docs/introduction/concepts/#sources)を使用します。一方、[Sinks](https://vector.dev/docs/introduction/concepts/#sinks)はイベントの宛先となります。そのため、ここでは Kafka ログソース と ClickHouse sink を使用します。なお、Kafka は Sink としてサポートされていますが、ClickHouse ログソース は利用できません。したがって、ClickHouse から Kafka にデータを転送したい場合、Vector は適していません。

Vector はデータの[transformation](https://vector.dev/docs/reference/configuration/transforms/)にも対応しています。これはこのガイドの対象外です。データセットに対してこの機能が必要な場合は、Vector のドキュメントを参照してください。

なお、現在の ClickHouse sink の実装では HTTP インターフェイスを使用しています。現時点では、ClickHouse sink は JSON スキーマの使用をサポートしていません。データは、プレーンな JSON フォーマットまたは String として Kafka に公開する必要があります。

<div id="license">
  ### ライセンス
</div>

Vector は [MPL-2.0 ライセンス](https://github.com/vectordotdev/vector/blob/master/LICENSE) のもとで配布されています

<div id="gather-your-connection-details">
  ### 接続情報を確認する
</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** をクリックします。

<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 管理者によって設定されます。

<div id="steps">
  ### 手順
</div>

1. Kafka の `github` トピックを作成し、[GitHub データセット](https://datasets-documentation.s3.eu-west-3.amazonaws.com/kafka/github_all_columns.ndjson)を投入します。

```bash theme={null}
cat /opt/data/github/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
```

このデータセットは、`ClickHouse/ClickHouse` リポジトリを対象とした 200,000 行で構成されています。

2. ターゲットテーブルが作成されていることを確認します。以下ではデフォルトのデータベースを使用します。

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

```

3. [Vector をダウンロードしてインストールします](https://vector.dev/docs/setup/quickstart/)。`kafka.toml` の設定ファイルを作成し、Kafka と ClickHouse の各インスタンスに合わせて値を調整します。

```toml theme={null}
[sources.github]
type = "kafka"
auto_offset_reset = "smallest"
bootstrap_servers = "<kafka_host>:<kafka_port>"
group_id = "vector"
topics = [ "github" ]
tls.enabled = true
sasl.enabled = true
sasl.mechanism = "PLAIN"
sasl.username = "<username>"
sasl.password = "<password>"
decoding.codec = "json"

[sinks.clickhouse]
type = "clickhouse"
inputs = ["github"]
endpoint = "http://localhost:8123"
database = "default"
table = "github"
skip_unknown_fields = true
auth.strategy = "basic"
auth.user = "username"
auth.password = "password"
buffer.max_events = 10000
batch.timeout_secs = 1
```

この設定と Vector の動作に関して、いくつか重要な注意点があります。

* この例は Confluent Cloud でテストされています。そのため、`sasl.*` および `ssl.enabled` のセキュリティオプションは、セルフマネージド環境には適さない可能性があります。
* 構成パラメータ `bootstrap_servers` には、プロトコルのプレフィックスは不要です。例: `pkc-2396y.us-east-1.aws.confluent.cloud:9092`
* ログソースパラメータ `decoding.codec = "json"` を指定すると、メッセージは 1 つの JSONオブジェクトとして ClickHouse sink に渡されます。メッセージを文字列として扱い、デフォルト値の `bytes` を使用する場合、メッセージの内容は `message` フィールドに追加されます。ほとんどの場合、これは [Vector getting started](/docs/ja/integrations/connectors/data-ingestion/etl-tools/vector-to-clickhouse#4-parse-the-logs) ガイドで説明されているように、ClickHouse 側で処理する必要があります。
* Vector はメッセージに [複数のフィールドを追加します](https://vector.dev/docs/reference/configuration/sources/kafka/#output-data)。この例では、構成パラメータ `skip_unknown_fields = true` を使って、ClickHouse sink でこれらのフィールドを無視しています。これにより、ターゲットテーブルのスキーマに含まれないフィールドは無視されます。`offset` などのメタフィールドも含めたい場合は、必要に応じてスキーマを調整してください。
* `inputs` パラメータを使って、sink がイベントログソースを参照している点に注目してください。
* ClickHouse sink の動作については、[こちら](https://vector.dev/docs/reference/configuration/sinks/clickhouse/#buffers-and-batches) の説明も確認してください。最適なスループットを得るには、`buffer.max_events`、`batch.timeout_secs`、`batch.max_bytes` の各パラメータを調整するとよいでしょう。ClickHouse の[推奨事項](/docs/ja/reference/statements/insert-into#performance-considerations)によれば、1 回のバッチに含めるイベント数は最低でも 1000 を目安にしてください。継続的に高スループットが見込まれるユースケースでは、`buffer.max_events` パラメータを増やすことを検討してください。スループットの変動が大きい場合は、`batch.timeout_secs` パラメータの調整が必要になることがあります。
* `auto_offset_reset = "smallest"` パラメータを指定すると、Kafka ログソースはトピックの先頭から読み取りを開始します。これにより、手順 (1) で公開したメッセージを確実に消費できます。必要な動作が異なる場合もあります。詳しくは [こちら](https://vector.dev/docs/reference/configuration/sources/kafka/#auto_offset_reset) を参照してください。

4. Vector を起動します

```bash theme={null}
vector --config ./kafka.toml
```

既定では、ClickHouse への挿入を開始する前に [ヘルスチェック](https://vector.dev/docs/reference/configuration/sinks/clickhouse/#healthcheck) が必要です。これにより、接続を確立でき、スキーマを読み取れることを確認できます。問題が発生した場合に役立つ追加のログを取得するには、先頭に `VECTOR_LOG=debug` を付けます。

5. データが挿入されたことを確認します。

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

| 件数     |
| :----- |
| 200000 |
