メッセージキューイング(Message Queue)は、分散システムにおける中核的なコンポーネントです。主な利用シーンとしては、即座に結果を取得する必要はないが、同時接続数の制御が必要な状況が該当します。本稿では、分散型メッセージキューイングシステムであるApache Kafkaとを調整サービスとして使用するZooKeeperを組み合わせたクラスタ構築手順を解説します。
1. メッセージキューイングの基礎
1.1 メッセージキューとは
メッセージ(Message)とはアプリケーション間で送受信されるデータ全般を指します。単純なテキスト文字列から複雑な埋め込みオブジェクトまで、多様な形式を取ることができます。
メッセージキューは、アプリケーション間通信の一形態です。メッセージ送信者はキューにメッセージを送信後、 즉시リ턴を受けられます。メッセージの可靠的配送はメッセージシステム側が保証します。送信側はメッセージをMQにPublishするのみで、誰が取得するか知る必要はなく、受信側はMQからConsumeするのみで、誰がPublishしたか知る必要ありません。この特性により、送信者と受信者は互いの存在を意識する必要がなくなります。
1.2 メッセージキューの特性
(1)バッファリング機能
従来のTCP/UDPといったソケットベースの_req-res_モデルと異なり、メッセージキューは通常、宛先プロセスがメッセージを読み出すか、明示的に削除するまで、メッセージをある種のバッファ領域に蓄積します。
(2)非同期処理
メッセージキューは、アプリケーションに一定程度の非同期性を導入します。送信プロセスはメッセージを送信し、キューに蓄積 tandis que 受信プロセスは処理するメッセージを取り出すといった運用が可能になります。これにより、通信が断続的거나 送信・受信プロセスのいずれかが障害発生した場合でも、アプリケーションの継続的運用が可能となります。
(3)ルーティング機能
メッセージキューは、複数のプロセスが同一キューに対して読み書きできるルーティング機能を提供します。Broadcast혹은 Unicast通信パターンの実現が可能です。
1.3 メッセージキューを採用する理由
(1)疎結合
送信側と受信側の処理プロセスを独立的に拡張・変更することが可能になります。同一のインターフェース制約を守るだけで済みます。
(2)耐障害性
メッセージキューはデータが完全に処理されるまで永続化を行うため、データ損失のリスクを回避できます。多くのメッセージキューが採用する_produce-consume_パターンでは、メッセージがキューから削除される前に、処理システムがそのメッセージが処理完了したことを明示的に通知する必要があります。これにより、データの安全な保存が保証されます。
(3)スケーラビリティ
メッセージキューによる疎結合により、メッセージ投入および処理の頻度を容易に引き上げられます。処理プロセスを追加するだけで対応可能です。
(4)ピーク流量への対応
突発的なアクセス増加においてもシステムは継続稼働必要です、このようなピーク流量は日常的ではありません。ピーク流量対応のために 常時リソースを確保するのはコスト効率が悪いでしょう。メッセージキューを活用すれば、重要なコンポーネントはburst access的压力에도崩溃하지 않고対応可能です。
(5)恢复可能性
システムの一部のコンポーネントが失效しても、全体システムへの影響は限定的입니다。プロセス間の結合度が低減されるためメッセージを処理するプロセ스가停止了としても、キュー内のメッセージはシステム恢复後に正常に處理될 수 있습니다.
(6)順序保証
多くのシナリオでは、データ処理の順序が重要になります。大多数の名前付きメッセージキューは顺序保证机构を実装しており、データは特定顺序で處理될ことを保証します。(Kafkaでは1つのPartition内でのメッセージ顺序が保証されます)
(7)バッファリング
システムを通過するデータフローの速度を制御・最適化できます。Producer와 Consumer간의 처리 속도 차이를 해결합니다.
(8)非同期通信
即座にメッセージを処理する必要がない場合、メッセージキューは非同期処理メカニズムを提供します。必要なだけのメッセージをキューに投入し、後から 원하는 시점에进行处理하면 됩니다.
2. Kafka の基礎
2.1 Kafka の基本概念
Apache Kafkaは、高吞吐量を維持する分散型Publish/Subscribeメッセージシステム입니다。比喩的に説明すると、昨今はデータが爆発的に増加する時代であり、各種ビジネス、ソーシャル、検索、ブラウジング活动中、海量のデータが生成されます。これらのデータを効率的に收集し、リアルタイム分析及を行うことが要求されます。これはProducer가 다양한データを生成하고、Consumer가それを消費・分析するビジネス需求として捉えることができます Producer와 Consumerの間に、通信の桥梁となるメッセージシステムが必要とされます这就是メッセージシステム導入の背景입니다。この,业务需求は跨系統での数据传输需求とも言い換えられます。
KafkaはApacheプロジェクトのオープンソース製品であり、最大の特徴は多様な需求シナリオに対してリアルタイム大量のデータを處理할 수 있다는점에あります。例えば Hadoop プラットフォーム 기반のデータ分析、低延迟实时系统、Storm/Spark 스트리밍処理エンジンなどが考えられます,现在Kafkaは多家的大型企业において多种のデータパイプラインおよびメッセージシステムとして使用されています。
2.2 Kafka の主要コンポーネント
Kafka の核心概念と主要コンポーネントは以下の通りです:
- Broker:Kafkaクラスタは1つ以上のサーバで構成され、各サーバのことをBrokerと呼びます。
- Topic:KafkaクラスタにPublishされる各メッセージは分類を持ち、これをTopic(主题)と呼びます。
- Producer:消息的生产者であり、メッセージをKafka BrokerにPublishします。
- Consumer:消息の消費者であり、Kafka BrokerからデータをPullして.publishされたメッセージを消費します。
- Partition:物理的概念であり、各Topicは1つ以上のPartitionを持ち、各Partitionは順序付けされた队列です。Partition内の各メッセージには一意のID(offset)が割り当てられます。
- Consumer Group:消费者グループであり、各Consumerに消费グループを指定できます。指定しない場合はデフォルトグループに属します。
- Message:通信の基本単位で、各Producerは任意のTopicにメッセージをPublishできます。
2.3 Kafka のシステム構成
典型的なKafkaクラスタは、複数のProducer、複数のBroker、複数のConsumer Group、そしてZooKeeperクラスタで構成されます。KafkaはZooKeeper를 사용하여クラスタ設定を管理し、Leaderを選出し、Consumer Groupに変動が発生した際にrebalanceを行います。Producer는 pushモデルでメッセージをBrokerにpublishし、Consumer는 pullモデルでBrokerからsubscribeおよび消費します。
アーキテクチャ图から确认できるように、典型的なメッセージシステムはProducer、ストレージシステム(Broker)、Consumerの3要素で構成されます。Kafkaは分散型メッセージシステムとして複数のProducerとConsumerをサポートし、Producer可以将消息分散发布到集群中不同节点的不同Partition上,消费者也可以消费集群中多个节点上的多个Partition。書き込み時は複数のProducerが同一Partitionに書き込めますが、読み取り時は1つのPartitionが1つのコンシューマグループ内の1つのConsumerのみ消費可能であり、1つのConsumerは複数のPartitionを消費できます。つまり同一Consumer Group内のConsumer对Partitionは排他的であり、異なるConsumer Group間では共有されます。
Kafkaは消息の永続化存储をサポートしており,永続化データはKafkaのログファイルに保存されます。Producer가メッセージを生成した後、KafkaはメッセージをConsumerに直接 전달하지 않고、まずBrokerで存储します。ディスク書き込み回数を低減するため、Brokerはメッセージをしばらくキャッシュし、メッセージの数またはサイズが的一定閾値に達した際にまとめてディスクに書き込みます。これによりKafkaの実行効率向上とディスクIO呼び出し回数の低減が実現されます。
Kafkaでは各メッセージをPartitionに書き込む际、磁盘に順序書き込みを行います。これは非常に重要な特性です。机械式ディスクではランダム書き込みは効率が悪いですが、順序書き込みは極めて高效입니다。この顺序写入磁盘メカニズムがKafkaの高吞吐量实现の重要な支えとなっています。
2.4 Topic と Partition の関係
KafkaのTopicはPartitionの形で存储됩니다。各Topicに対してPartition数を设定できます。Partition数はTopicを構成するlogファイルの数を决定します。推奨事項としてPartition数は同時に稼働するConsumer数より多い状態が良く、さらにPartition数はクラスタBroker数以内に设定することで消息データが各Brokerに 均分散されます。
Topicに複数のPartitionを設定する理由は、Kafkaはファイル存储ベースの构造이며、複数のPartitionを設定することで消息内容を複数のBrokerに分散存储でき、単一マシンのディスク容量上限に到达することを回避できます。同时、1つのTopicを任意の数のpartitionsに分割することで消息存储・消息消費の効率が向上します。Partition数が越多ればそれだけ多くのConsumer容纳 가능であり、Kafkaの吞吐率が効果的に提升されます。因此、Topicを複数のpartitionsに分割する利点は大量的メッセージを複数バッチに分割して異なるノードに同時書き込みでき、書き込みリクエストをクラスタノード全体に負荷分散できる点上にあります。
存储構造上、各Partitionは物理的に1つのフォルダに対応し、当該フォルダにはそのPartitionのすべてのメッセージとインデックスファイルが存储されます。Partition命名規則は「Topic名称+序号」であり、最初のpartition序号は0から始まり、最大値はpartitions数-1となります。
各Partition(フォルダ)には複数のサイズが等しいsegment(セグメント)データファイルが存在します。各segmentサイズは同一ですが、1件のメッセージサイズは不一定であるため、segmentデータファイル内のメッセージ数は必ずしも等しくありません。segmentデータファイルはindex fileとdata fileの2部分で構成され、両者は一一对应でペアになり、接尾辞分别是".index"と".log"です。
2.5 Producer のessage送信机制
Producerはメッセージとデータの生成者であり、メッセージをBrokerに送信际、Partitionメカニズムに基づいて存储先Partitionを決定します。Partitionメカニズムが適切に设定されていれば、すべてのメッセージを異なるPartitionに 均分散でき、データの负载均衡が実現されます。1つのTopicが1つのファイルに対応する場合、当該ファイルが存在するマシンのI/Oが当該Topicの性能ボトルネックになりますが、Partitionがある場合、異なるメッセージが異なるBrokerの異なるPartitionに並列書き込み可能となり、吞吐率が 크게向上します。
2.6 Consumer のメッセージ消費机制
Kafkaは2つのメッセージPublishパターンをサポートします:キューモード(Queuing)とPublish/Subscribeモードです。キューモードでは、1つの消費グループのみが存在し、当該消費グループには複数のConsumerが含まれ、1つのメッセージは当該消費グループの1つのConsumerのみが消費可能です。Publish/Subscribeモードでは、複数の消費グループが存在でき、各消費グループには1つのConsumerのみが含まれ、同じメッセージが複数の消費グループに消費될 수 있습니다.
KafkaにおけるProducerとConsumerはpush、pullパターンを採用しています。 즉 Producer가 Brokerにメッセージをpushし、Consumer가 Brokerからpullします。pushとpullはメッセージ生成と消費において非同期的に進行します。pullパターンの利点はConsumerがメッセージ消費速度を自己能動的に制御できること、およびバッチでBrokerからデータをpullするか逐条消費するか选择 가능한点上にあります。
3. ZooKeeper の概念
3.1 分散協調サービスとは
ZooKeeperは分散協調技術の一つです。分散協調技術とは、分散環境内で複数のプロセス間の同期制御を行い、共享リソースへのアクセスが有序进行的,防止资源竞争(脑裂)后果的发生为主要目的とした技術です。脑裂とは、Master/Slave切り替え時に切り替えが不完全に或其他原因により、Client와 Slave가误って2つのactivemasterが存在すると認識し、最終的にクラスタ全体が混乱状態に陥る現象を指します。
まず分散システムとは何かを説明します。分散システムとは、異なる地域に分布する複数のサーバが共同で構成し、用户提供服务する应用システムを言います。分散システムにおいて最も重要的是进程的调度です。例えば、3つの地域に分布するサーバで構成される应用システムがあり、第1のマシンにリソースがマウントされており、3つの地域に分布する应用プロセスが 모두当該リソースを競争する必要があるとします。しかし、複数のプロセスが同時にアクセスすることは望ましくありません,这时候就需要一个协调者(ロック)让他们有序地访问共享资源。この協調者が分散システムでいわゆる「ロック」です。例えば「プロセス1」が当該リソースを使用する际、まず当該ロックを獲得する必要があります。「プロセス1」がロックを獲得後、当該リソースに対して独占状态になり、この间他のプロセスは当該リソースにアクセスできません。「プロセス1」がリソースを使用し终えた後、ロックを解放して他のプロセスが获得できるようにします。これにより、この「ロック」メカニズムを通じて分散システム内の複数のプロセスが有序的に共享リソースにアクセスであることが保证されます。这里把这个分散環境下的「ロック」を分散ロックと呼び、この分散ロックが分散協調技術實現の核心的内容です。
目前,在分散協調技術方面做得比較好的实现として、GoogleのChubbyとApacheのZooKeeperがあり、いずれも分散ロックの実現者です。ZooKeeperが提供するロックサービスは分散分野においては 長年にわたり使用されてきた実績があり、その信頼性と可用性は理論と実践の両面で验证済みです。
ZooKeeperは分散应用のために设计された高可用・高性能开源協調サービスであり、基本サービスとして分散ロックサービスをサポートします。同时、データの维护与管理メカニズム诸如:统一命名服务、状态同期服务、クラスタ管理、分散メッセージキュー、分散应用設定項目の管理等を提供します。
3.2 單一故障点問題とZooKeeperによる解決策
(1)單一故障点問題とは
單一故障点とは、Master-Slave構成の分散システムにおいて、Masterノードが任务调度・分配的を担当し、Slaveノードが任务の處理を担当する構成で、Masterノードに障害が発生した場合、システム全体が使用不能となることを言います。この問題を解決するには、クラスタMaster役选择を通じて分散システムにおける單一故障点問題に対応します。
(2)従来の解決方法とその問題点
従来の解決策は、待機ノード(バックアップノード)を采用します。待機ノードは定期的にMasterノードにpingパケットを送信し、Masterノードはpingパケット受領後に待機ノードに応答ACK情報を返送します。待機ノードが応答を受信した場合、現在のMasterノードは正常に稼働中と判断し、服务的継続させます。Masterノードに障害が発生すると、待機ノードは応答を受信できなくなるため、Masterノードの障害発生と判断し、待機ノードがこれを接替して新しいMasterとなります。
この従来の單一故障点解決策は一定の範囲で問題を缓解하지만、潜在的な課題があります network障害の可能性です。例えば次のような状況が発生し得ます:Masterノード自体は障害発生していないが、ACK応答返送時にnetwork障害が発生した場合、待機ノードは応答を受信できず、Masterノードに障害が発生したと判断します。そして待機ノードはMasterノードのサービスを接管し新的Masterとなります。この時点で分散システムには2つのMasterノード(Dual Master構成)が存在することになり、Dual Masterの存在は分散システムの服务混乱を引き起こします。この情況が発生するのを防止하기 위해、ZooKeeperを導入하여問題を解決する必要があります。
3.3 ZooKeeper の動作原理
以下、3つの情形を通じて説明します:
(1)Master起動時
分散システムにZooKeeperを導入した後、複数のMasterノードを設定できます。便宜上2つのMasterノードを設定した場合、假定它们分别是MasterAとMasterBです。2つのMasterノードが起動すると、どちらもZooKeeperにノード情報を登録します。假定MasterAが登録するノード情報はmaster00001であり、MasterBが登録するノード情報はmaster00002とします。登録完了後、选举が执行されます。选举には複数のアルゴリズムが存在しますが、ここでは编号最小を选举アルゴリズムとして採用します。那么编号最小の 노드가选举で勝利し、ロックを獲得してMasterとなります 즉MasterAがロックを獲得してMasterとなり、MasterBはブロッキングされてスタンバイノードとなります。この方式を通じて、ZooKeeperは2つのMasterプロセスのスケジューリングを完了し、Master/Standbyノードの分配と協調を実現します。
(2)Master障害発生時
MasterAに障害が発生した場合、ZooKeeperに登録されたノード情報は自動的に削除されます。ZooKeeperはノードの変更を自動的に感知し、MasterAの障害を検出后会再次发起选举,此时MasterBが选举で勝利し、MasterAに代わって新的Masterとなります。これにより、Master/Standbyノードの重新选举が完了します。
(3)Master恢复時
Masterノードが恢复すると,再次向ZooKeeper注册自身的ノード信息を送信します。ただし这时候注册的ノード信息将是master00003となり、元の情報とは異なります。ZooKeeperはノード変更を感知し再次发起选举,这时候MasterBが再次当选してMasterを継続担任し、MasterAはスタンバイノードを担当します。
ZooKeeperはこのような協調・スケジューリングメカニズムを通じてクラスタの管理と状态同期を行います。
3.4 ZooKeeper クラスタ アーキテクチャ
ZooKeeperは通常クラスタアーキテクチャでサービスを提供します。以下はZooKeeperの基本アーキテクチャ図です:
ZooKeeperクラスタ主要な役立っているServer와 Clientがあり、ServerはさらにLeader、Follower、Observerの3つの役割に分かれます。各役柄の意味は以下の通りです:
- Leader:リーダー役であり、主に投票の発起と決議、システム状态の更新を担当します。
- Follower:フォロワー役であり、クライアントからのリクエストを受け付けて結果を返す他、选举 과정에서投票に参加します。
- Observer:オブザーバー役であり、クライアントからのリクエストを受け取り、書き込みリクエストをLeaderに転送,同時にLeader状态を同期하지만投票には参加しません。Observerの導入目的是扩展系统能力、提高伸缩性を目的としています。
- Client:クライアント役であり、ZooKeeperへのリクエスト发起を担当します。
3.5 ZooKeeper の処理流れ
ZooKeeperデータの更新流れ:ZooKeeperクラスタ内の各Serverは内存中にデータを保存しており、ZooKeeper起動時にクラスタ实例から选举で1つのServerをLeaderに选出します。Leaderがデータ更新等の操作を担当し、绝大多数のServerが内存中のデータを正常に更新했을場合にのみ、数据修改成功と判断されます。
ZooKeeper書き込みの流れ:クライアントClientは最初に1つのServer或者Observerと通信し書き込みリクエスト发起します。その後Serverは書き込みリクエストをLeaderに転送し、Leaderは書き込みリクエストを他のServerに転送します。他のServerは書き込みリクエスト受領後、データを書き込み、Leaderに応答します。Leaderは書き込み成功の応答を大多数から受領した後、データ書き込み成功と判断し最後にClientに応答して書き込み操作完了します。
4. ZooKeeper の Kafka における役割
4.1 Broker登録
Brokerは分散配置されており互いに独立していますが、クラスタ内のBroker全体を管理する登録システムが必要です,这里就使用了ZooKeeper。ZooKeeper上にはBrokerサーバーリストを記録する专用ノードが存在します:
/brokers/ids
各Brokerは起動時にZooKeeperに登録を行います 즉/brokers/ids配下に自身のノードを作成します 예: /brokers/ids/[0...N]のようにります。
Kafkaでは全局的に一意な数字で各Brokerサーバーを識別하며、異なるBrokerは異なるBroker IDを使用して注册する必要があります。ノード作成完了後、各Brokerは自身のIPアドレスとポート情報を 해당ノードに登録します。其中Brokerが作成したノードタイプは临时ノードであり、Brokerが停止すると对应的临时ノードも自動的に削除されます。
4.2 Topic登録
Kafkaでは同一Topicのメッセージは複数の分区に分割されて複数のBrokerに分散配置されます。これらの分区情報とBrokerの対応関係もZooKeeperが维护하며、专用ノードで記録されます 예: /brokers/topicsのようにります。
Kafka内の各Topicは/brokers/topics/[topic]の形式で記録されます 예: /brokers/topics/login와 /brokers/topics/searchなどとなります。Brokerサーバーが起動後、対応するTopicノード(/brokers/topics)配下に自身のBroker IDを登録し、当該Topicの分区総数を書き込みます 예: /brokers/topics/login/3->2はBroker IDが3のBrokerサーバーが"login"というTopicに対して2つの分区でメッセージ存储を提供することを示しています。同样、この分区ノードも临时ノードです。
4.3 Producerの负载分散
同一Topicのメッセージは分区されて複数のBrokerに分散配置されるため、Producerはメッセージをこれらの分散したBrokerに適切に送信する必要があります。那么如何实现Producerの负载分散呢?Kafkaは従来の4层负载分散,支持ZooKeeper方式实现负载分散の両方をサポートしています。
(1)4層负载分散
ProducerのIPアドレスとポートに基づいて関連するBrokerを特定します。通常、1つのProducerは1つのBrokerに対応し、当該Producerが生成するメッセージは全て 해당Brokerに送信されます。この方式は論理的に简单であり、各Producerは他のシステムと追加のTCP接続を確立する必要がなく、Brokerと单个TCP接続を維持するだけで済みます。ただし、真の负载分散做不到であり、実際のシステムでは各Producerのメッセージ生成量と各Brokerのメッセージ保存量は各不相同です。もし一部のProducerのメッセージ生成量が他のProducerより大幅に多い場合、異なるBrokerが受信するメッセージ総数が大きな 차이가产生하며、ProducerはBrokerの新增と削除をリアルタイムで感知できません。
(2)ZooKeeperによる负载分散
各Broker起動時にBroker登録流程が完了するため、Producerは当該ノードの変更を通じてBrokerサーバーリストの変動を動的に感知可以实现动态的负载均衡机制によりProducerはBroker新增・削除をリアルタイムで感知し自動的に负载分散を調整できます。
4.4 Consumerの负载分散
Producerと同様に、Kafka内のConsumerも负载分散を実現して複数のConsumerが対応するBrokerサーバからメッセージを適切に受信する必要があります。各Consumerグループには複数のConsumerが含まれ、各メッセージはグループ内の1つのConsumerのみに送信されます。異なるConsumerグループは各自特定のTopic配下のメッセージを消費し互いに干涉しません。
4.5 Partition と Consumer の対応関係
Consumer Group配下には複数のConsumerが存在します。各Consumer Groupに対して、Kafkaは全局的に一意なGroup IDを割り当てます。Group内の全Consumerは当該IDを共有します。購読するTopic配下の各Partitionは某个group内の1つのconsumerにのみ割り当て可能であり、当然のことながら当該Partitionは他のgroupにも割り当て可能です。
同时、Kafkaは各ConsumerにConsumer IDを割り当て、通常「Hostname:UUID」形式で表されます。
Kafkaでは、各消息分区只能被同组的一个Consumer消费という规则があります,因此需要ZooKeeperに消息分区とConsumerの対応関係を記録する必要があります。各Consumerはメッセージ分区の消費権利を確定した後、Consumer IDをZooKeeperの対応メッセージ分区の临时ノードに書き込みます예: /consumers/[group_id]/owners/[topic]/[broker_id-partition_id]のようにります。
其中、[broker_id-partition_id]は消息分区の識別子であり、ノード内容は当該消息分区上のConsumerのConsumer IDです。
4.6 メッセージ消費進捗Offset記録
Consumerが指定メッセージ分区のメッセージを消費する过程中、一定間隔で分区メッセージの消費進捗OffsetをZooKeeperに記録する必要があります。这样做的目的是、当该Consumerが再起動或其他Consumerが当該メッセージ分区のメッセージ消費を引き継ぐ際、之前的進捗から継続消费可以实现。OffsetはZooKeeper内の専用ノードで記録され、ノード路径は以下の形式です:
/consumers/[group_id]/offsets/[topic]/[broker_id-partition_id]
ノード内容はOffsetの値です。
4.7 Consumer登録
Consumerサーバーが初期起動時にConsumerグループに参加する流程は以下の通りです:
(1)Consumerグループへの登録
各Consumerサーバーが起動する際、ZooKeeperの指定ノード配下に自身に対応するConsumerノードを作成します예: /consumers/[group_id]/ids/[consumer_id]のようにります。ノード作成完了後、Consumerは自身が購読するTopic情報を当該临时ノードに書き込みます。
(2)Consumerグループ内のConsumer変動に対する监听登録
各Consumerは所属Consumerグループ内の他のConsumerサーバーの変動を注目する必要があります 즉/consumers/[group_id]/idsノードに子ノード変動のWatcher监听を登録します。一旦Consumer新增または減少を検出した場合、Consumerの负载分散がトリガーされます。
(3)Brokerサーバー変動に対する监听登録
Consumerは/broker/ids/[0-N]内のノードを监听する必要があります。Brokerサーバーリストに変動があった場合、具体情况に応じてConsumer负载分散が必要かどうかを決定します。
(4)Consumer负载分散の執行
同一Topic配下の異なる分区のメッセージを複数のConsumerに尽可能均等に消費させるため、Consumerと消息分区の割り当て流程を実行します。通常、1つのConsumerグループ内でConsumerサーバー変動またはBrokerサーバー変動が発生した場合、Consumer负载分散が发起されます。
5. 單一ノード Kafka 環境の構築
動作環境
構築対象ホスト:
kafka01:192.168.10.101
5.1 ZooKeeper の導入
[root@kafka01 ~]# yum -y install java [root@kafka01 ~]# tar -xzf apache-zookeeper-3.6.0-bin.tar.gz [root@kafka01 ~]# mv apache-zookeeper-3.6.0-bin /opt/zookeeper [root@kafka01 ~]# cd /opt/zookeeper/conf [root@kafka01 ~]# mv zoo_sample.cfg zoo.cfg [root@kafka01 ~]# vim zoo.cfg dataDir=/opt/zookeeper/data [root@kafka01 ~]# cd /opt/zookeeper/ [root@kafka01 zookeeper]# mkdir -p /opt/zookeeper/data [root@kafka01 zookeeper]# ./bin/zkServer.sh start [root@kafka01 zookeeper]# ./bin/zkServer.sh status
5.2 Kafka の導入
[root@kafka01 ~]# tar -xzf kafka_2.13-2.4.1.tgz [root@kafka01 ~]# mv kafka_2.13-2.4.1 /opt/kafka [root@kafka01 ~]# cd /opt/kafka/ [root@kafka01 kafka]# vim config/server.properties log.dirs=/opt/kafka/logs #60行目付近 [root@kafka01 kafka]# mkdir -p /opt/kafka/logs [root@kafka01 kafka]# bin/kafka-server-start.sh config/server.properties & # 포트 활성화 상태 확인 [root@kafka01 kafka]# netstat -tlnp | grep 2181 [root@kafka01 kafka]# netstat -tlnp | grep 9092
注意:起動時はまずZooKeeperを启动し、关闭時はまずKafkaを終了します。
ZooKeeperの終了:
[root@kafka01 zookeeper]# ./bin/zkServer.sh stop
Kafkaの終了:
[root@kafka01 kafka]# bin/kafka-server-stop.sh # 終了できない場合、プロセスを強制終了 [root@kafka01 kafka]# pkill -f kafka
5.3 動作確認
トピック作成:
[root@kafka01 kafka]# bin/kafka-topics.sh --create \ --zookeeper kafka01:2181 \ --replication-factor 1 \ --partitions 1 \ --topic sample
トピック一覧表示:
[root@kafka01 kafka]# bin/kafka-topics.sh --list --zookeeper kafka01:2181
トピック詳細確認:
[root@kafka01 kafka]# bin/kafka-topics.sh --describe \ --zookeeper kafka01:2181 \ --topic sample
メッセージ生産(プロデューサー):
[root@kafka01 kafka]# bin/kafka-console-producer.sh \ --broker-list kafka01:9092 \ --topic sample
メッセージ消費(コンシューマー):
[root@kafka01 kafka]# bin/kafka-console-consumer.sh \ --bootstrap-server kafka01:9092 \ --topic sample \ --from-beginning
トピック削除:
[root@kafka01 kafka]# bin/kafka-topics.sh --delete \ --zookeeper kafka01:2181 \ --topic sample
6. Kafka クラスタ の構築
動作環境
クラスタ構成ホスト:
kafka01:192.168.10.101 kafka02:192.168.10.102 kafka03:192.168.10.103
6.1 ZooKeeper クラスタ の導入
(1)ZooKeeperのインストール(全ノード共通)
[root@kafka01 ~]# yum -y install java [root@kafka01 ~]# tar -xzf apache-zookeeper-3.6.0-bin.tar.gz [root@kafka01 ~]# mv apache-zookeeper-3.6.0-bin /opt/zookeeper
(2)データ保存用ディレクトリの作成(全ノード共通)
[root@kafka01 ~]# mkdir -p /opt/zookeeper/data
(3)設定ファイルの編集(全ノード共通)
[root@kafka01 ~]# cd /opt/zookeeper/conf [root@kafka01 conf]# mv zoo_sample.cfg zoo.cfg [root@kafka01 conf]# vim zoo.cfg dataDir=/opt/zookeeper/data clientPort=2181 server.1=192.168.10.101:2888:3888 server.2=192.168.10.102:2888:3888 server.3=192.168.10.103:2888:3888
ポート説明:
- 2181:クライアントにサービスを提供するポート
- 3888:Leader选举に使用するポート
- 2888:クラスタ内ノード通信用ポート(LeaderがListen)
(4)ノードIDファイルの作成(ノード별로異なる値を設定)
# kafka01ノード [root@kafka01 conf]# echo '1' > /opt/zookeeper/data/myid # kafka02ノード [root@kafka02 conf]# echo '2' > /opt/zookeeper/data/myid # kafka03ノード [root@kafka03 conf]# echo '3' > /opt/zookeeper/data/myid
(5)全ノードでZooKeeperプロセスを起動
[root@kafka01 ~]# cd /opt/zookeeper/ [root@kafka01 zookeeper]# ./bin/zkServer.sh start [root@kafka01 zookeeper]# ./bin/zkServer.sh status
6.2 Kafka クラスタ の導入
(1)Kafkaのインストール(全ノード共通)
[root@kafka01 ~]# tar -xzf kafka_2.13-2.4.1.tgz [root@kafka01 ~]# mv kafka_2.13-2.4.1 /opt/kafka
(2)設定ファイルの編集
[root@kafka01 ~]# cd /opt/kafka/ [root@kafka01 kafka]# vim config/server.properties broker.id=1 #21行目付近 各ノードで一意の値を設定 listeners=PLAINTEXT://192.168.10.101:9092 #31行目付近 各ノードでIPアドレスを変更 log.dirs=/opt/kafka/logs #60行目付近 num.partitions=1 #65行目付近 Partition数はノード数以内に設定 zookeeper.connect=192.168.10.101:2181,192.168.10.102:2181,192.168.10.103:2181
ポート説明:
- 9092:Kafkaのリスンポート
(3)ログ用ディレクトリの作成(全ノード共通)
[root@kafka01 kafka]# mkdir -p /opt/kafka/logs
(4)全KafkaノードでBrokerプロセスを起動
[root@kafka01 kafka]# bin/kafka-server-start.sh config/server.properties & # 起動に失敗する場合、ログディレクトリ内のデータをクリア后再試行 [root@kafka01 kafka]# rm -rf /opt/kafka/logs/* [root@kafka01 kafka]# bin/kafka-server-start.sh config/server.properties &
6.3 クラスタ動作確認
トピック作成(いずれかのノードで実行):
[root@kafka01 kafka]# bin/kafka-topics.sh --create \ --zookeeper kafka01:2181 \ --replication-factor 1 \ --partitions 1 \ --topic sample
トピック一覧表示(各ノードで表示確認):
[root@kafka01 kafka]# bin/kafka-topics.sh --list --zookeeper kafka01:2181 [root@kafka02 kafka]# bin/kafka-topics.sh --list --zookeeper kafka02:2181 [root@kafka03 kafka]# bin/kafka-topics.sh --list --zookeeper kafka03:2181
メッセージ生産(プロデューサー):
[root@kafka01 kafka]# bin/kafka-console-producer.sh \ --broker-list kafka01:9092 \ --topic sample
メッセージ消費(コンシューマー):
[root@kafka02 kafka]# bin/kafka-console-consumer.sh \ --bootstrap-server kafka02:9092 \ --topic sample \ --from-beginning
トピック削除:
[root@kafka01 kafka]# bin/kafka-topics.sh --delete \ --zookeeper kafka01:2181 \ --topic sample