メッセージキューの概念
- メッセージキューとは何か
- メッセージ(Message)とはアプリケーション間で送信されるデータのこと
- メッセージキュー(Message Queue)はアプリケーション間の通信方式を解決するソリューションで、メッセージの信頼性のある伝達を保証します。
- メッセージキューが必要な理由
- 疎結合
- 両側の処理プロセスを独立して拡張または修正できるため、同じインターフェース制約を遵守する限り、柔軟性が高まります
- 冗長性
- メッセージキューはデータを永続化し、完全に処理されるまで保存します。この方法により、データ損失のリスクを回避できます。多くのメッセージキューが採用する「挿入-取得-削除」パラダイムでは、メッセージをキューから削除する前に、処理システムが明示的にメッセージが完全に処理されたことを示す必要があり、これによりデータが安全に保存されます
- スケーラビリティ
- メッセージキューが処理プロセスを疎結合にするため、メッセージの投入と処理の頻度を増やすのは簡単です。追加の処理プロセスを用意するだけで済みます
- 柔軟性とピーク処理能力
- アクセス量が急増する場合でも、アプリケーションは機能し続ける必要があります。しかし、このような突発的なトラフィックは稀です。常にピークアクセスを処理できるようにリソースを投入するのは大きな無駄です。メッセージキューを使用すると、重要なコンポーネントが突発的なアクセス圧力に耐えられ、過負荷のリクエストによって完全にクラッシュすることを防ぎます
- 回復性
- システムの一部のコンポーネントが故障しても、全体のシステムに影響を与えません。メッセージキューはプロセス間の結合度を低減するため、メッセージを処理するプロセスがダウンしても、キューに入れられたメッセージはシステムの復旧後に処理できます
- 順序保証
- 多くの使用シナリオでは、データ処理の順序が重要です。ほとんどのメッセージキューは元々ソートされており、特定の順序でデータが処理されることを保証します。(KafkaはPartition内のメッセージの順序性を保証します)
- 非同期通信
- ユーザーがメッセージを即時に処理したくない、あるいは必要ない場合があります。メッセージキューは非同期処理メカニズムを提供し、ユーザーがメッセージをキューに入れ、必要なときに処理できるようにします
メッセージキューの特徴
- ストレージ
- メッセージをバッファ領域に保存し、ターゲットプロセスがメッセージを読み取るか、明示的にキューから削除するまで保持します
- 非同期性
- メッセージキューはバッファリングされたメッセージを通じてアプリケーション内で一定の非同期性を公開し、ソースプロセスがメッセージを送信してキューにメッセージを蓄積し、ターゲットプロセスがメッセージを選択して処理できるようにします
Kafkaの基本概念
- Kafkaとは何か
- Kafkaは高スループットの分散型パブリッシュサブスクライブメッセージシステムです
- KafkaはApache組織のオープンソースシステムです
- 多様なニーズシナリオを満たすために大量のデータをリアルタイムで処理するオープンソースシステムです
- Kafkaの役割と用語
- Broker:Kafkaクラスタは1つ以上のサーバーを含み、各サーバーはbroker(ブローカー)と呼ばれます
- Topic:Kafkaクラスタに公開される各メッセージには分類があり、この分類はTopic(トピック)と呼ばれます
- Producer:メッセージのプロデューサを指します。Kafkaブローカーにメッセージを公開する責任があります
- Consumer:メッセージのコンシューマを指し、Kafkaブローカーからデータをプルし、公開されたメッセージを消費します。
- Partition:物理的な概念で、各Topicは1つ以上のPartitionを含みます。各Partitionのメッセージには順序付きのID(offset)が割り当てられます
- Consumer Group:コンシューマグループで、各Consumerにコンシューマグループを指定できます。指定しない場合、デフォルトのグループが使用されます
- Message:通信の基本単位で、各producerは1つのtopicにいくつかのメッセージを公開できます
ZooKeeperの基本概念
- ZooKeeperは分散型調整技術です。分散型調整技術は主に分散環境における複数のプロセス間の同期制御を解決し、共有リソースに順序正しくアクセスさせ、リソース競合(ブレインスプリット)の結果を防ぎます。
ZooKeeperの動作原理
- マスター起動
- 各ノードはZooKeeperにノード情報を登録し、最小番号アルゴリズムで1つのマスターを選出します。残りのノードはバックアップノードとなり、ZooKeeperが2つのマスターのスケジューリング、マスターとバックアップノードの割り当てとプロトコルを完了します。
- マスター障害
- マスターAが障害を起こした場合、ZooKeeperに登録されたノード情報は自動的に削除され、再選挙が行われます
- マスター修復
- マスターが修復された場合、再びZooKeeperに自身のノード情報を登録しますが、登録されたノード番号は小さくなるため、マスターではなく別のノードが引き続きマスターを担当します
単一ノードでのKafkaデプロイメント
# ZooKeeperのインストール
[root@bogon ~]# yum -y install java
[root@bogon ~]# tar zxvf apache-zookeeper-3.6.0-bin.tar.gz
[root@bogon ~]# mv apache-zookeeper-3.6.0-bin /etc/zookeeper
[root@bogon ~]# cd /etc/zookeeper/conf
[root@bogon conf]# mv zoo_sample.cfg zoo.cfg
[root@bogon conf]# vim zoo.cfg
# 内容を変更
dataDir=/etc/zookeeper/zookeeper-data
[root@bogon conf]# cd /etc/zookeeper/
[root@bogon kafka]# mkdir zookeeper-data
[root@bogon zookeeper]# ./bin/zkServer.sh start
[root@bogon zookeeper]# ./bin/zkServer.sh status
# Kafkaのインストール
[root@bogon ~]# tar zxvf kafka_2.13-2.4.1.tgz
[root@bogon ~]# mv kafka_2.13-2.4.1 /etc/kafka
[root@bogon ~]# cd /etc/kafka/
# 内容を変更
[root@kafka1 kafka]# vim config/server.properties
log.dirs=/etc/kafka/kafka-logs #60行目
[root@bogon kafka]# mkdir /etc/kafka/kafka-logs
[root@bogon kafka]# bin/kafka-server-start.sh config/server.properties &
# 2つのポートの開放状態を確認
[root@bogon kafka]# netstat -anpt | grep 2181
[root@bogon kafka]# netstat -anpt | grep 9092
# テスト
クラスタでのKafkaデプロイメント
# ホストのhostsファイルを修正(すべてのホストで設定)
[root@kafka1 ~]# vim /etc/hosts
192.168.10.101 kafka1
192.168.10.102 kafka2
192.168.10.103 kafka3
# ZooKeeperのデプロイメント
# ZooKeeperのインストール(3つのノードの設定は同じ)
[root@kafka1 ~]# systemctl stop firewalld
[root@kafka1 ~]# setenforce 0
[root@kafka1 ~]# yum -y install java
[root@kafka1 ~]# tar zxvf apache-zookeeper-3.6.0-bin.tar.gz
[root@kafka1 ~]# mv apache-zookeeper-3.6.0-bin /etc/zookeeper
# データ保存ディレクトリの作成(3つのノードの設定は同じ)
[root@kafka1 ~]# cd /etc/zookeeper/
[root@kafka1 zookeeper]# mkdir zookeeper-data
# 設定ファイルの修正(3つのノードの設定は同じ)
[root@kafka1 zookeeper]# cd /etc/zookeeper/conf
[root@kafka1 ~]# mv zoo_sample.cfg zoo.cfg
[root@kafka1 ~]# vim zoo.cfg
dataDir=/etc/zookeeper/zookeeper-data
clientPort=2181
server.1=192.168.10.101:2888:3888 #2181:クライアントにサービスを提供 3888:リーダー選出に使用 #2888:クラスタ内のマシン間通信に使用(リーダーがこのポートを監視)
server.2=192.168.10.102:2888:3888
server.3=192.168.10.103:2888:3888
# ノードIDファイルの作成(server番号に基づいてIDを設定、3つのマシンで異なる)
# ノード1:
[root@kafka1 conf]# echo '1' > /etc/zookeeper/zookeeper-data/myid
# ノード2:
[root@kafka2 conf]# echo '2' > /etc/zookeeper/zookeeper-data/myid
# ノード3:
[root@kafka3 conf]# echo '3' > /etc/zookeeper/zookeeper-data/myid
# 3つのノードでZooKeeperプロセスを起動
[root@kafka1 conf]# cd /etc/zookeeper/
[root@kafka1 zookeeper]# ./bin/zkServer.sh start
[root@kafka1 zookeeper]# ./bin/zkServer.sh status
# Kafkaのインストール(3つのノードの設定は同じ)
[root@kafka1 ~]# tar zxvf kafka_2.13-2.4.1.tgz
[root@kafka1 ~]# mv kafka_2.13-2.4.1 /etc/kafka
# 設定ファイルの修正
[root@kafka1 ~]# cd /etc/kafka/
[root@kafka2 kafka]# vim config/server.properties
broker.id=1 # 他の2つのIDはそれぞれ2と3
listeners=PLAINTEXT://192.168.10.101:9092 # 他のノードはそれぞれのIPアドレスに変更
log.dirs=/etc/kafka/kafka-logs
num.partitions=1 # パーティション数、ノード数以下
zookeeper.connect=192.168.10.101:2181,192.168.10.102:2181,192.168.10.103:218
# ログディレクトリの作成(3つのノードの設定は同じ)
[root@kafka1 kafka]# mkdir /etc/kafka/kafka-logs
# すべてのKafkaノードで起動コマンドを実行し、Kafkaクラスタを生成(3つのノードの設定は同じ)
[root@kafka1 kafka]# ./bin/kafka-server-start.sh config/server.properties &
# テスト
# トピックの作成
[root@bogon kafka]# ./kafka-topics.sh --create --zookeeper kafka1:2181 --replication-factor 1 --partitions 1 --topic test
# トピックの一覧表示
[root@bogon kafka]# ./kafka-topics.sh --list --zookeeper kafka1:2181
# トピックの確認
[root@bogon kafka]# ./kafka-topics.sh --describe --zookeeper kafka1:2181 --topic test
# メッセージの生成
[root@bogon kafka]# ./kafka-console-producer.sh --broker-list kafka1:9092 -topic test
# メッセージの消費(別のターミナルを開き、メッセージを生成しながら消費メッセージを確認)
[root@bogon kafka]# ./kafka-console-consumer.sh --bootstrap-server kafka1:9092 --topic test
# トピックの削除
bin/kafka-topics.sh --delete --zookeeper kafka1:2181 --topic test