Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

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

Kafka と ClickHouse で Vector を使用する

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

ClickHouse 向け Vector の入門ガイドでは、ログのユースケースとファイルからのイベント読み取りに重点を置いています。ここでは、Kafka トピックに保持されたイベントを含むGithub sample datasetを使用します。

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

Vector はデータのtransformationにも対応しています。これはこのガイドの対象外です。データセットに対してこの機能が必要な場合は、Vector のドキュメントを参照してください。

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

ライセンス

Vector は MPL-2.0 ライセンス のもとで配布されています

接続情報を確認する

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

ClickHouse Cloud サービスの接続ボタン

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

ClickHouse Cloud HTTPS 接続情報

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

手順

  1. Kafka の github トピックを作成し、GitHub データセットを投入します。
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 行で構成されています。

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

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);
  1. Vector をダウンロードしてインストールしますkafka.toml の設定ファイルを作成し、Kafka と ClickHouse の各インスタンスに合わせて値を調整します。
[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 ガイドで説明されているように、ClickHouse 側で処理する必要があります。
  • Vector はメッセージに 複数のフィールドを追加します。この例では、構成パラメータ skip_unknown_fields = true を使って、ClickHouse sink でこれらのフィールドを無視しています。これにより、ターゲットテーブルのスキーマに含まれないフィールドは無視されます。offset などのメタフィールドも含めたい場合は、必要に応じてスキーマを調整してください。
  • inputs パラメータを使って、sink がイベントログソースを参照している点に注目してください。
  • ClickHouse sink の動作については、こちら の説明も確認してください。最適なスループットを得るには、buffer.max_eventsbatch.timeout_secsbatch.max_bytes の各パラメータを調整するとよいでしょう。ClickHouse の推奨事項によれば、1 回のバッチに含めるイベント数は最低でも 1000 を目安にしてください。継続的に高スループットが見込まれるユースケースでは、buffer.max_events パラメータを増やすことを検討してください。スループットの変動が大きい場合は、batch.timeout_secs パラメータの調整が必要になることがあります。
  • auto_offset_reset = "smallest" パラメータを指定すると、Kafka ログソースはトピックの先頭から読み取りを開始します。これにより、手順 (1) で公開したメッセージを確実に消費できます。必要な動作が異なる場合もあります。詳しくは こちら を参照してください。
  1. Vector を起動します
vector --config ./kafka.toml

既定では、ClickHouse への挿入を開始する前に ヘルスチェック が必要です。これにより、接続を確立でき、スキーマを読み取れることを確認できます。問題が発生した場合に役立つ追加のログを取得するには、先頭に VECTOR_LOG=debug を付けます。

  1. データが挿入されたことを確認します。
SELECT count() AS count FROM github;
件数
200000
Navigation