Details and best practices of RocketMQ's consumer types
RocketMQ 5.0 では、クライアントタイプ、特にコンシューマータイプの概念が重視されています。RocketMQ には、PushConsumer、SimpleConsumer、PullConsumer の 3 種類のコンシューマーがあり、それぞれ異なるビジネスシナリオに対応しています。
コンシューマータイプの概要
この記事では、各コンシューマータイプについて詳しく説明します。各メッセージタイプを紹介する前に、RocketMQ のコンシューマーに共通するワークフローを整理しておきます。コンシューマーでは、クライアントがサーバーに積極的にリクエストを送信し、ロングポーリングを維持することでメッセージを受信します。メッセージ到着の適時性を確保するため、クライアントはサーバーに継続的にリクエストを送信し続ける必要があります(リクエストの開始がクライアント主導かどうかは、コンシューマータイプによって異なります)。条件を満たす新しいメッセージがサーバーに到着すると、クライアントはそのメッセージを受信します。最後に、サーバーはクライアントの処理結果に応じて、メッセージの処理結果を記録します。
さらに、PushConsumer と SimpleConsumer には ConsumerGroup という概念があります。これは、同じサブスクリプション関係を持つコンシューマーグループの共通 ID に相当します。サーバーは ConsumerGroup ごとに消費プログレスを記録します。同じ ConsumerGroup 内のメッセージコンシューマーは、サブスクリプショングループの条件を満たすすべてのメッセージを協調して消費し、個別に消費するわけではありません。PullConsumer と比較して、PushConsumer と SimpleConsumer はビジネス統合シナリオに適しています。消費状態とプログレスがサーバー側で管理されるため、比較的軽量でシンプルな実装となります。
簡潔にまとめると以下の通りです。
・PushConsumer:完全にホストされたコンシューマータイプです。ユーザーはメッセージリスナーを登録するだけで、対応するサブスクリプション関係に一致するメッセージが自動的に消費メソッドを呼び出します。ビジネス統合で最も一般的に使用されるコンシューマータイプです。
・SimpleConsumer:メッセージの受信とプログレス同期を分離したコンシューマータイプです。ユーザーはサーバーからメッセージを個別に受信し、確認できます。PushConsumer と同様に消費プログレスはサーバーが管理しますが、ユーザーが消費レートを独自に制御する必要があるビジネスシナリオに適しています。
・PullConsumer:ストリーム処理フレームワークが管理するコンシューマータイプです。ユーザーはキュー(Topic の最小論理単位)ごとにメッセージを受信し、消費オフセットの自動コミットまたは手動コミットを選択できます。
PushConsumer
PushConsumer は、現在 RocketMQ で最も広く使用されているコンシューマーです。ユーザーはサブスクリプション関係を確認した後、対応するリスナーを登録するだけで済みます。プロデューサーがサブスクリプション関係に一致するメッセージを送信すると、コンシューマーのリスナーインターフェイスが即座に呼び出されます。このとき、ユーザーはリスナー内に対応するビジネスロジックを実装する必要があります。
ユーザーは自身のビジネス処理結果に応じて、ConsumeResult.SUCCESS または ConsumeResult.FAILURE を返す必要があります。ConsumeResult.SUCCESS を返した場合、メッセージは正常に消費されたと見なされます。ConsumeResult.FAILURE を返した場合、サーバーは消費失敗と判断し、メッセージのバックオフ再試行を行います。バックオフ再試行とは、メッセージが正常に消費されるまで、登録された MessageListener に複数回配信されることであり、2 回の配信間の時間間隔はバックオフルールに従います。
特に、各 ConsumerGroup には最大消費回数が設定されています。現在のメッセージ消費がこの設定を超えると、メッセージは再配信されず、代わりにデッドレターキューに送信されます。この消費回数は、メッセージが MessageListener に配信されるたびに自動的に増加します。たとえば、メッセージの最大消費回数が 1 の場合、消費成功または消費失敗の戻り値に関わらず、メッセージは 1 回だけ消費されます。
アプリケーションシナリオとベストプラクティス
PushConsumer はほぼ完全にホストされたコンシューマーです。ここでの「ホスト」とは、ユーザーがメッセージの受信について気にする必要がなく、メッセージの消費処理のみに関心を払えばよいことを意味します。その他のロジックはすべて PushConsumer の実装にカプセル化されています。ユーザーは受信した各メッセージに応じて異なる消費結果を返すだけでよいため、最も人気のあるコンシューマータイプです。
ほとんどのシナリオでは、ユーザーは消費ロジックを素早く処理して消費成功を返すべきであり、消費ロジックを長時間ブロックすべきではありません。重い消費ロジックの場合は、まず消費ステータスを報告してから、メッセージを非同期で処理することが推奨されます。
実際、PushConsumer の実装では、メッセージ消費の適時性を確保するために、クライアントが事前にメッセージをプルして後続の消費に備えています。そのため、クライアント内にはプルされたメッセージサイズのキャッシュが存在します。キャッシュされたメッセージが多すぎてクライアントのメモリリークが発生しないよう、クライアントパラメータはユーザーが自身で設定できるよう予約されています。
SimpleConsumer では、ユーザーは SimpleConsumer#receive インターフェイスを通じて自身でメッセージをプルし、ビジネスロジックの処理結果に応じてプルしたメッセージを異なる方法で処理する必要があります。SimpleConsumer#receive もロングポーリングでサーバーからメッセージを受信します。具体的なロングポーリング時間は SimpleConsumerBuilder#setAwaitDuration を使用して設定できます。
SimpleConsumer では、SimpleConsumer#receive を通じてタイムウィンドウ(このインターフェイスで受信したメッセージの非表示タイムウィンドウ)を設定する必要があります。タイムウィンドウはユーザーがメッセージを受信した時点からカウントされます。この期間中、メッセージはコンシューマーに再配信されませんが、タイムウィンドウを超えるとメッセージは再配信されます。このプロセスで、メッセージの消費回数も増加します。PushConsumer と同様に、消費回数が ConsumerGroup の最大回数を超えると、再配信されなくなります。
PushConsumer と比較して、SimpleConsumer のユーザーはメッセージ受信のリズムを独自に制御できます。SimpleConsumer#receive は現在のサブスクリプション関係の条件を満たすメッセージをサーバーからプルします。実際、SimpleConsumer の各メッセージ受信リクエストは、具体的な Topic パーティションごとに個別に発行されます。実際の Topic パーティション数は多い場合があるため、メッセージ受信の適時性を確保するには、自身のビジネス処理能力に合わせて SimpleConsumer#receive の同時実行を適度に上げることが推奨されます。
メッセージ受信後、ユーザーはメッセージに対して ack または changeInvisibleDuration のいずれかを選択できます。前者はサーバーにメッセージを確認する意味で、PushConsumer の消費成功に相当します。後者は現在のメッセージの可視化時間を遅延させることを意味し、サーバーは現在の期間が経過した後にクライアントにメッセージを配信します。ここでのメッセージ再配信も ConsumerGroup の最大消費回数の制限に従う点に注意してください。つまり、メッセージの消費回数が最大消費回数を超えると(消費回数はメッセージが可視時間になるたびに自動的に増加)、メッセージは再配信されず、代わりにデッドレターキューに入ります。例:
・ack を実行すると、メッセージ消費が正常に確認され、サーバーが消費プログレスを同期します。
・changeInvisibleDuration:
1) メッセージが現在の ConsumerGroup の最大消費回数を超えている場合、メッセージはその後デッドレターキューに配信されます。
2) メッセージが現在の ConsumerGroup の最大消費回数を超えておらず、リクエストが最後のメッセージ可視時間より前に発行された場合、変更は成功します。それ以外の場合は変更は失敗します。
アプリケーションシナリオとベストプラクティス
PushConsumer では、メッセージが MessageListener に渡されて処理されます。SimpleConsumer では、ユーザーが同時に複数のメッセージを取得できます。各バッチの最大メッセージ数は SimpleConsumer#receive のパラメータによって決まります。一部の IO 集約型アプリケーションでは、より便利な選択肢となります。この場合、ユーザーは毎回バッチでメッセージを取得し、集中的に処理することで消費速度を向上できます。
PullConsumer
PullConsumer も RocketMQ がこれまでサポートしてきたコンシューマータイプです。RocketMQ 5.0 の新しい PullConsumer API はまだ開発中です。ご期待ください。以下の PullConsumer の説明では、4.0 の既存の LitePullConsumer を使用します。これが現在推奨される方法です。
概要
RocketMQ では、メッセージはキューを通じて送受信されます。Topic は複数のキューで構成されます。メッセージはキューの形式で 1 つずつ格納されます。同じキュー内のメッセージには異なるオフセットがあり、オフセットのサイズはメッセージがサーバーに到達した時間とともに増加します。本質的に、サーバー上の異なる ConsumerGroup の消費プログレスは、キュー内のオフセット情報です。クライアントは消費プログレスをサーバーに同期しますが、これは本質的にメッセージのオフセット同期です。
PullConsumer では、キューの概念がユーザーに完全に公開されています。ユーザーは関心のある Topic にルートリスナーを設定してキューの変化を検知し、現在のコンシューマーにキューを割り当てることができます。LitePullConsumer#poll を使用すると、割り当て済みのキューからメッセージの取得を試みます。LitePullConsumer#setAutoCommit が設定されている場合、メッセージがクライアントに到達するとオフセットが自動コミットされます。それ以外の場合は、LitePullConsumer#commitSync インターフェイスを使用して手動コミットする必要があります。
アプリケーションシナリオとベストプラクティス
PullConsumer では、ユーザーがメッセージオフセットを管理する完全な権限を持ち、消費プログレスを自身で管理できます。これが PushConsumer および SimpleConsumer との最も本質的な違いであり、消費レートと消費プログレスの両方を独立して制御する必要があるストリームコンピューティングシナリオで広く使用されている理由でもあります。多くの場合、PullConsumer は特定のストリームコンピューティングフレームワークと統合して使用されます。
コンシューマータイプの概要
この記事では、各コンシューマータイプについて詳しく説明します。各メッセージタイプを紹介する前に、RocketMQ のコンシューマーに共通するワークフローを整理しておきます。コンシューマーでは、クライアントがサーバーに積極的にリクエストを送信し、ロングポーリングを維持することでメッセージを受信します。メッセージ到着の適時性を確保するため、クライアントはサーバーに継続的にリクエストを送信し続ける必要があります(リクエストの開始がクライアント主導かどうかは、コンシューマータイプによって異なります)。条件を満たす新しいメッセージがサーバーに到着すると、クライアントはそのメッセージを受信します。最後に、サーバーはクライアントの処理結果に応じて、メッセージの処理結果を記録します。
さらに、PushConsumer と SimpleConsumer には ConsumerGroup という概念があります。これは、同じサブスクリプション関係を持つコンシューマーグループの共通 ID に相当します。サーバーは ConsumerGroup ごとに消費プログレスを記録します。同じ ConsumerGroup 内のメッセージコンシューマーは、サブスクリプショングループの条件を満たすすべてのメッセージを協調して消費し、個別に消費するわけではありません。PullConsumer と比較して、PushConsumer と SimpleConsumer はビジネス統合シナリオに適しています。消費状態とプログレスがサーバー側で管理されるため、比較的軽量でシンプルな実装となります。
簡潔にまとめると以下の通りです。
・PushConsumer:完全にホストされたコンシューマータイプです。ユーザーはメッセージリスナーを登録するだけで、対応するサブスクリプション関係に一致するメッセージが自動的に消費メソッドを呼び出します。ビジネス統合で最も一般的に使用されるコンシューマータイプです。
・SimpleConsumer:メッセージの受信とプログレス同期を分離したコンシューマータイプです。ユーザーはサーバーからメッセージを個別に受信し、確認できます。PushConsumer と同様に消費プログレスはサーバーが管理しますが、ユーザーが消費レートを独自に制御する必要があるビジネスシナリオに適しています。
・PullConsumer:ストリーム処理フレームワークが管理するコンシューマータイプです。ユーザーはキュー(Topic の最小論理単位)ごとにメッセージを受信し、消費オフセットの自動コミットまたは手動コミットを選択できます。
PushConsumer
PushConsumer は、現在 RocketMQ で最も広く使用されているコンシューマーです。ユーザーはサブスクリプション関係を確認した後、対応するリスナーを登録するだけで済みます。プロデューサーがサブスクリプション関係に一致するメッセージを送信すると、コンシューマーのリスナーインターフェイスが即座に呼び出されます。このとき、ユーザーはリスナー内に対応するビジネスロジックを実装する必要があります。
ユーザーは自身のビジネス処理結果に応じて、ConsumeResult.SUCCESS または ConsumeResult.FAILURE を返す必要があります。ConsumeResult.SUCCESS を返した場合、メッセージは正常に消費されたと見なされます。ConsumeResult.FAILURE を返した場合、サーバーは消費失敗と判断し、メッセージのバックオフ再試行を行います。バックオフ再試行とは、メッセージが正常に消費されるまで、登録された MessageListener に複数回配信されることであり、2 回の配信間の時間間隔はバックオフルールに従います。
特に、各 ConsumerGroup には最大消費回数が設定されています。現在のメッセージ消費がこの設定を超えると、メッセージは再配信されず、代わりにデッドレターキューに送信されます。この消費回数は、メッセージが MessageListener に配信されるたびに自動的に増加します。たとえば、メッセージの最大消費回数が 1 の場合、消費成功または消費失敗の戻り値に関わらず、メッセージは 1 回だけ消費されます。
アプリケーションシナリオとベストプラクティス
PushConsumer はほぼ完全にホストされたコンシューマーです。ここでの「ホスト」とは、ユーザーがメッセージの受信について気にする必要がなく、メッセージの消費処理のみに関心を払えばよいことを意味します。その他のロジックはすべて PushConsumer の実装にカプセル化されています。ユーザーは受信した各メッセージに応じて異なる消費結果を返すだけでよいため、最も人気のあるコンシューマータイプです。
ほとんどのシナリオでは、ユーザーは消費ロジックを素早く処理して消費成功を返すべきであり、消費ロジックを長時間ブロックすべきではありません。重い消費ロジックの場合は、まず消費ステータスを報告してから、メッセージを非同期で処理することが推奨されます。
実際、PushConsumer の実装では、メッセージ消費の適時性を確保するために、クライアントが事前にメッセージをプルして後続の消費に備えています。そのため、クライアント内にはプルされたメッセージサイズのキャッシュが存在します。キャッシュされたメッセージが多すぎてクライアントのメモリリークが発生しないよう、クライアントパラメータはユーザーが自身で設定できるよう予約されています。
SimpleConsumer では、ユーザーは SimpleConsumer#receive インターフェイスを通じて自身でメッセージをプルし、ビジネスロジックの処理結果に応じてプルしたメッセージを異なる方法で処理する必要があります。SimpleConsumer#receive もロングポーリングでサーバーからメッセージを受信します。具体的なロングポーリング時間は SimpleConsumerBuilder#setAwaitDuration を使用して設定できます。
SimpleConsumer では、SimpleConsumer#receive を通じてタイムウィンドウ(このインターフェイスで受信したメッセージの非表示タイムウィンドウ)を設定する必要があります。タイムウィンドウはユーザーがメッセージを受信した時点からカウントされます。この期間中、メッセージはコンシューマーに再配信されませんが、タイムウィンドウを超えるとメッセージは再配信されます。このプロセスで、メッセージの消費回数も増加します。PushConsumer と同様に、消費回数が ConsumerGroup の最大回数を超えると、再配信されなくなります。
PushConsumer と比較して、SimpleConsumer のユーザーはメッセージ受信のリズムを独自に制御できます。SimpleConsumer#receive は現在のサブスクリプション関係の条件を満たすメッセージをサーバーからプルします。実際、SimpleConsumer の各メッセージ受信リクエストは、具体的な Topic パーティションごとに個別に発行されます。実際の Topic パーティション数は多い場合があるため、メッセージ受信の適時性を確保するには、自身のビジネス処理能力に合わせて SimpleConsumer#receive の同時実行を適度に上げることが推奨されます。
メッセージ受信後、ユーザーはメッセージに対して ack または changeInvisibleDuration のいずれかを選択できます。前者はサーバーにメッセージを確認する意味で、PushConsumer の消費成功に相当します。後者は現在のメッセージの可視化時間を遅延させることを意味し、サーバーは現在の期間が経過した後にクライアントにメッセージを配信します。ここでのメッセージ再配信も ConsumerGroup の最大消費回数の制限に従う点に注意してください。つまり、メッセージの消費回数が最大消費回数を超えると(消費回数はメッセージが可視時間になるたびに自動的に増加)、メッセージは再配信されず、代わりにデッドレターキューに入ります。例:
・ack を実行すると、メッセージ消費が正常に確認され、サーバーが消費プログレスを同期します。
・changeInvisibleDuration:
1) メッセージが現在の ConsumerGroup の最大消費回数を超えている場合、メッセージはその後デッドレターキューに配信されます。
2) メッセージが現在の ConsumerGroup の最大消費回数を超えておらず、リクエストが最後のメッセージ可視時間より前に発行された場合、変更は成功します。それ以外の場合は変更は失敗します。
アプリケーションシナリオとベストプラクティス
PushConsumer では、メッセージが MessageListener に渡されて処理されます。SimpleConsumer では、ユーザーが同時に複数のメッセージを取得できます。各バッチの最大メッセージ数は SimpleConsumer#receive のパラメータによって決まります。一部の IO 集約型アプリケーションでは、より便利な選択肢となります。この場合、ユーザーは毎回バッチでメッセージを取得し、集中的に処理することで消費速度を向上できます。
PullConsumer
PullConsumer も RocketMQ がこれまでサポートしてきたコンシューマータイプです。RocketMQ 5.0 の新しい PullConsumer API はまだ開発中です。ご期待ください。以下の PullConsumer の説明では、4.0 の既存の LitePullConsumer を使用します。これが現在推奨される方法です。
概要
RocketMQ では、メッセージはキューを通じて送受信されます。Topic は複数のキューで構成されます。メッセージはキューの形式で 1 つずつ格納されます。同じキュー内のメッセージには異なるオフセットがあり、オフセットのサイズはメッセージがサーバーに到達した時間とともに増加します。本質的に、サーバー上の異なる ConsumerGroup の消費プログレスは、キュー内のオフセット情報です。クライアントは消費プログレスをサーバーに同期しますが、これは本質的にメッセージのオフセット同期です。
PullConsumer では、キューの概念がユーザーに完全に公開されています。ユーザーは関心のある Topic にルートリスナーを設定してキューの変化を検知し、現在のコンシューマーにキューを割り当てることができます。LitePullConsumer#poll を使用すると、割り当て済みのキューからメッセージの取得を試みます。LitePullConsumer#setAutoCommit が設定されている場合、メッセージがクライアントに到達するとオフセットが自動コミットされます。それ以外の場合は、LitePullConsumer#commitSync インターフェイスを使用して手動コミットする必要があります。
アプリケーションシナリオとベストプラクティス
PullConsumer では、ユーザーがメッセージオフセットを管理する完全な権限を持ち、消費プログレスを自身で管理できます。これが PushConsumer および SimpleConsumer との最も本質的な違いであり、消費レートと消費プログレスの両方を独立して制御する必要があるストリームコンピューティングシナリオで広く使用されている理由でもあります。多くの場合、PullConsumer は特定のストリームコンピューティングフレームワークと統合して使用されます。
Related Articles
-
A detailed explanation of Hadoop core architecture HDFS
Knowledge Base Team
-
What Does IOT Mean
Knowledge Base Team
-
6 Optional Technologies for Data Storage
Knowledge Base Team
-
What Is Blockchain Technology
Knowledge Base Team
Explore More Special Offers
-
Short Message Service(SMS) & Mail Service
50,000 email package starts as low as USD 1.99, 120 short messages start at only USD 1.00
