-
Notifications
You must be signed in to change notification settings - Fork 0
DNET_ApacheKafka
- 戻る(分散処理:データ収集・格納系)
- Apache Sqoop
- Fluentd/Embulk
- Logstash, Beats(Elasticsearchの該当節を参照)
- Apache Flume
- Apache NiFi
- Apache Kafka
-
デファクト・スタンダードの分散メッセージング・システム、
- 複数台のマシンでクラスタを構成
- 分散処理により高いスループットを発揮
-
最近は、
- ストリーム処理の中核と言われるようになって来ているらしい。
- “ETL is Dead; Long Live Stream”なんて言われてるらしい。
多数のサービス間で、従来のMQ的に機能する。
複数のデータ・ソースに様々な処理をする際に、
一度、データハブにキューイングする的な。
ストリームだけでなく、バッチにも対応した、
複数のデータ・ソースからのデータ収集基盤としても利用可能。
- スケールアウト
- 都度フラッシュしない(OS任せ)。
スケールアウト+レプリケーション
-
スケールアウト
- オートスケール
- 対応していない
- 送信・受信側の設定変更必要になるため。
- オートスケール
-
レプリケーション
- Leader / Follower型レプリケーション
- レプリケーションは、メッセージ受信後、直ぐに行われる。
- 都度フラッシュしない問題をレプリケーションで担保。
高信頼性メッセージング(再取得が可能な、確実な送受信)
- データ送信者(Publisher)のメッセージングを受信し、レプリケーションまで完了したら、ACKを返す。
- データ受信者(Subscriber)がメッセージを受信したら、メッセージに処理完了を記録する。
連携可能なプロダクトが多い。
- Publisher / Subscriberメッセージングモデルを採用
- MQTTなどの標準プロトコルではなく、独自プロトコルを使用
の構造と管理
-
論理データ
-
構造
-
管理方法
- メッセージはRecordと言うキーバリュー形式
- Leader Replicaに書き込まれたRecordはFollower Replicaに複製される。
- Partitionへの書き込み / 読み出しはLeader Replicaにのみ行われる。
- Leader Replicaに障害が発生したら、いずれかのFollower Replicaが昇格する。
- 同期しているFollower Replicaを、In Sync Replica(ISR)と呼ぶ。
- ISR数が最小数まで回復すると、書き込みはコミットされ読み出し可能になる。
-
-
物理データ構造
-
構造
- BrokerはReplicaごとにデータディレクトリを作成する。
- RecordをSegmentファイルに保存することで永続化する。
- Segmentファイルの集合をLogと呼ぶ。
-
管理方法
- Brokerには、複数のディスクが接続される。
- データディレクトリは各ディスクにラウンドロビン方式で割り当てられる。
- Segmentファイルのキャッシュは OS制御(WindowsユーザがLinuxに乗り換える際に、知っておくとイイ情報集。の該当節を参照)の処理で高速化可能
-
Partitioning(1つのTopicは複数のPartitionで構成される)の話。
-
Partitioner
- デフォルトのPartitionerのPartitioning方式はハッシュ方式
- partition key が null の場合は、ラウンドロビン方式になる。
- partition key に基づいてメッセージをルーティングできる。
- 別途、独自のPartitionerを実装できる。
- デフォルトのPartitionerのPartitioning方式はハッシュ方式
-
Partition数は、後から増やすことダケはできる。
ただし、Partitioner+partition keyでPartitioningしている場合、
後からPartition数を増やすことが難しくなりがち。
クラスタ内のサーバ・ノードのことをBrokerと呼ぶ。
- Broker同士はクラスタコーディネータであるZooKeeperを使用して連携する。
- Brokerの1つがリーダーとなりBrokerクラスタを管理する。
-
書き込み側のアプリケーション(Publisher)は、
Producerと呼ばれる書き込み用ライブラリを使用してRecordを書き込む。 -
データ送信者(Publisher)は、
-
Recordを送信するためにProducerのSend APIを呼び出す。
-
Producerは、以下を行う。
- 定期的にいずれかのBrokerからメタデータを取得する。
- 各Brokerのホスト名、接続先ポート、Leader Replicaの場所など把握する。
- Produce リクエストでは、通信オーバーヘッド削減のため、ネットワークスレッドが
複数Topicの「Record Batch」を1回のリクエストでまとめて送信する。
-
acksの設定によって、異なる返信タイミングで、リクエスト完了通知を受信する。
- 0:即時
- 1:Leader Replicaへの書き込み完了時
- all:最小ISR数まで複製完了時
-
-
読み出し側のアプリケーション(Subscriber)は、
Consumerと呼ばれる読み出し用ライブラリを使用してRecordを取り出す。 -
データ受信者(Subscriber)は、
-
Recordを取得するためにConsumerのPoll APIを呼び出す。
-
Consumerは、以下を行う。
-
定期的にいずれかのBrokerからメタデータを取得する。
-
各Brokerのホスト名、接続先ポート、Leader Replicaの場所など把握する。
-
Recordをどこまで読みだしたのかを示すOffsetを管理する。
-
Consumer内部のキューに取得対象のRecordがなかった場合
・BrokerにFetchリクエストを投げてRecordを取得
・Fetchリクエストでは、以下の3つを指定
・取得対象のTopic
・Partitionのリスト
・各Partitionで取得したいOffset(Record番号)の範囲
・取得したRecord BatchはConsumer内部のキューに格納される。
-
-
-
Consumer Group
1つ以上のConsumerでConsumer Groupを構成することで、
1TopicのデータをGroup内のConsumerで分散読み出しできる。
分散処理:ストリーム系の該当節を参照。
- リクエストの並列送信数(max.in.flight.requests.per.connection)を1に設定
- 特定のキーに対応するRecordを、特定のPartition(のRecord Batch)に集める
- Consumer Groupの中でもパーティション毎に割り当てが可能
クライアントが使用するBrokerリソースを制御できる。
- ネットワーク帯域幅
- ネットワークI/O
- ディスクI/O
- スレッドのCPU要求レート
-
各キーの最新Recordのみを保持する。
-
一定期間が過ぎたRecordを削除するのではなく、
同じキーの最新のRecord以外を削除する機能。
- JMX (Java Management Extensions) により性能情報などのメトリックを公開。
- JMXに対応した監視ソフトウェアを使用することで、これらのメトリックを収集できる。
-
At-least-once保証
- 最低一回。重複が有り得る。
- Kafkaの既定値。
-
At-most-once保証
- 最大一回。ロストが有り得る。
- 設定&AP実装(冪等性)により可能。
冪等性 → 連番を付与的な。そんなレベル感。
-
Exactly-once保証
- 必ず一回。理想的であるが技術的には難しい。
- トランザクション、エンドツーエンドで保障。
DBなどの外部システムとKafka間でデータを書き込み/読み出しするコネクタを定義
- クライアントライブラリとして機能し、
- ウィンドウ処理などのストリーム系の分散処理を実装可能。
- Kafkaクラスタ間で(Produceによる)ミラーリングしバックアップを行う。
- フォールトトレランス機能としての使用は意図されていない。
-
Apache Kafka
https://kafka.apache.org/ -
Apache Kafka - Wikipedia
https://en.wikipedia.org/wiki/Apache_Kafka -
StormとKafkaによる
リアルタイムデータ処理 - Yahoo! JAPAN Tech Blog
https://techblog.yahoo.co.jp/programming/storm/ -
ITアーキテクトブログ - Medium
Apache Kafkaを使ったアプリ設計で反省している件を正直ベースで話す -
ETLは過去のものか - Apache Kafkaがデータ処理の未来なのか?
https://www.infoq.com/jp/articles/batch-etl-streams-kafka/
-
Kafka
https://qiita.com/kenji-kondo/items/0e9beffb91406f4d7122 -
sigmalist
https://qiita.com/sigmalist-
Apache Kafkaの
-
概要とアーキテクチャ
https://qiita.com/sigmalist/items/5a26ab519cbdf1e07af3 -
Producer/Broker/Consumerのしくみと設定一覧
https://qiita.com/sigmalist/items/3b512e2ab49b07271665 -
推奨構成と性能の見積もり方法
https://qiita.com/sigmalist/items/b3d25e914dc30539ece4
-
-
Apache Kafkaの性能検証
-
(1): 検証環境とパラメータチューニングの内容
https://qiita.com/sigmalist/items/730c7c82e02e5a837b9a -
(2): Producerのチューニング結果
https://qiita.com/sigmalist/items/0cd76d7edb055e89d1e6 -
(3): Brokerのチューニング結果
https://qiita.com/sigmalist/items/0795bde01bcbd78b1c47 -
(4): Producerの再チューニングおよびConsumerのチューニング
https://qiita.com/sigmalist/items/b0bc222f449a9f3cc273 -
(5): システム全体のレイテンシについて
https://qiita.com/sigmalist/items/111e96cdae69bb2d89e6
-
-
-
オープンソースカンファレンス2017 .Enterprise - イベント案内
2017-12-08 (金): めざせ!Kafkaマスター ~Apache Kafkaで最高の性能を出すには~
https://www.ospn.jp/osc2017.enterprise/modules/eguide/e33.html- めざせ!Kafkaマスター ~Apache Kafkaで最高の性能を出すには~
https://www.ospn.jp/osc2017.enterprise/pdf/OSC2017.enterprise_Hitachi_Kafka.pdf
- めざせ!Kafkaマスター ~Apache Kafkaで最高の性能を出すには~
-
最近広がりつつあるストリーム処理を知ろう!
~Apache Kafkaを活用したストリーム処理の基本とユースケース~
セミナープログラム - オープンソースカンファレンス2020 Online/Fall
https://event.ospn.jp/osc2020-online-fall/session/202763- OSC2020 OnlineFall 10-24 A-6 - YouTube
https://www.youtube.com/watch?v=dQyxa-d8zUI
- OSC2020 OnlineFall 10-24 A-6 - YouTube
-
ストリーム処理に広く使われるApache Kafkaの超概要
~一問一答形式で簡単にKafkaをご紹介~
セミナープログラム - オープンソースカンファレンス2020 Online/Fukuoka
https://event.ospn.jp/osc2020-online-fukuoka/session/244772- OSC2020 Online/Fukuoka 2020-11-28 B-5 - YouTube
https://www.youtube.com/watch?v=38LUK4hl3so
- OSC2020 Online/Fukuoka 2020-11-28 B-5 - YouTube
- KafkaのCustom Partitionerを試す - abcdefg.....
http://pppurple.hatenablog.com/entry/2018/11/02/234505 - Apache KafkaのConsumerを、特定のパーティションに手動で割り当てる - CLOVER🍀
https://kazuhira-r.hatenablog.com/entry/20180413/1523631913
移行メモ
- 「Miller Maker」は「MirrorMaker」の誤記のため修正した。
- 「Logsatsh, Beats」→「Logstash, Beats」、 「レプリケーションLeader / Follower型レプリケーション」→ 「Leader / Follower型レプリケーション」(語の重複)、 「これらのメリックを収集」→「メトリックを収集」に修正した。
- マイクロソフト系技術情報 Wiki(techinfoofmicrosofttech.osscons.jp)への URL リンクは、移行済みの メール / 証明書 / Azure Event Hubs に張り替えた。
Tags: 移行, Apache Kafka, 分散メッセージング, Publisher, Subscriber, Broker, Partition, Kafka Streams, Kafka Connect, ストリーム処理
このWikiは「Open棟梁Project」,「OSSコンソーシアム 開発基盤部会」によって運営されています。