Skip to main content

はじめに

Apache Kafka は、高性能なデータパイプライン、ストリーミング分析、ミッションクリティカルなアプリケーションに使用される、オープンソースの分散型イベントストリーミングプラットフォームです。

MoEngage と Kafka

この連携により、Connected Sources を介して Kafka トピックのイベントを MoEngage に直接ストリーミングできます。これにより、生のイベントデータが Kafka から MoEngage にほぼリアルタイムで流れます。 この連携では、次のことができます。
  • リアルタイムイベントのストリーミング: Kafka トピックのイベントを 1 秒未満のレイテンシで MoEngage に送信し、キャンペーンの即時トリガーやユーザープロファイルの更新に活用できます。
  • 環境をまたいだスケーリング: Systemd、Docker、Docker Compose、Kubernetes などの複数の環境にデプロイし、本番環境レベルの信頼性を実現できます。

ユースケース

Kafka と MoEngage を連携すると、次のユースケースに対応できます。
  • リアルタイムの購入トリガー: 顧客が購入を完了し、それが Kafka に取り込まれると、注文イベントを即座に MoEngage にストリーミングし、購入商品や注文金額に基づいて、パーソナライズされたお礼メール、商品レビューのリクエスト、クロスセルキャンペーンをトリガーできます。
  • カート放棄のリカバリー: e コマースプラットフォームのカート放棄イベントを Kafka 経由で MoEngage にストリーミングし、パーソナライズされたおすすめ商品や期間限定の割引コードを含むリマインダーメールを自動的にトリガーできます。
  • ユーザーアクティビティのトラッキング: ページビュー、機能の使用状況、コンテンツエンゲージメントなどの大量のユーザーインタラクションを Kafka トピックで取得します。それらを MoEngage にストリーミングして行動プロファイルを構築し、コンテキストに応じたアプリ内メッセージやプッシュ通知をトリガーできます。

連携

この連携は、Kafka のイベントを JSON 形式で MoEngage に直接送信するアーキテクチャに基づいています。 Kafka のイベントが JSON 形式で MoEngage に流れる様子を示すアーキテクチャ図 Kafka Consumer(サンプルの Python スクリプトを以下に示します)を使用して、Kafka トピックのイベントを MoEngage の API エンドポイントにプッシュできます。MoEngage はこれらのイベントを処理し、データをユーザープロファイルに表示します。
前提条件
  • デプロイ環境からアクセス可能なブートストラップサーバーを持つ、アクティブな Kafka クラスター。
  • デプロイ環境に Python 3.7 以上がインストールされていること。
  • Kafka のイベントが標準の JSON スキーマに従った形式になっていること。
  • Kafka のメッセージ構造と認証要件を理解していること。

ステップ 1: MoEngage エンドポイントを取得する

MoEngage Support チームに連絡して、専用の Kafka 連携エンドポイントを取得します。次のステップのために、以下のフィールドの値を控えておいてください。

ステップ 2: Kafka コネクターを設定する

以下のセクションでは、Kafka トピックのデータを専用の MoEngage エンドポイントにプッシュするためのサンプルスクリプトを示します。データ構造や属性のマッピングに応じてスクリプトを変更できます。

ステップ 2.1: 標準のイベント形式

MoEngage は、Kafka イベント用の標準 JSON 形式を提供しています。Kafka メッセージが次の構造に従っていることを確認してください。
Sample Event
フィールドの説明
柔軟なスキーマuser_attributes オブジェクトと event_attributes オブジェクトにはカスタムフィールドを追加できます。MoEngage はこれらのフィールドを自動的に取得して保存します。必須なのは customer_id、updated_at、event_attributes.event_name のみです。

ステップ 2.2: 基本設定

Python の依存関係 Kafka のコンシューム処理と HTTPS 通信に必要な Python パッケージをインストールします。
Shell
必要なパッケージ 環境を設定する 設定を安全に保存するために .env ファイルを作成します。
このファイルはバージョン管理にコミットしないでください。
Environment Configuration
Kafka の環境変数

ステップ 2.3: コンシューマースクリプトを作成する

kafka_consumer.py という名前のファイルを作成します。このスクリプトは Kafka からメッセージをコンシュームし、MoEngage API に送信します。
kafka_consumer.py

ステップ 3: スクリプトをテストする

本番環境にデプロイする前に、連携が正しく動作することを確認します。
テストの前提条件コンシューマースクリプトを実行する前に、Kafka トピックに標準 JSON 形式のテストイベントが含まれていることを確認してください。

3.1. コンシューマースクリプトをローカルで実行する

Shell
想定されるコンソール出力:
Sample Logs

3.2. MoEngage UI でデータを確認する

  1. MoEngage UI で、Settings > Data Management > Events に移動します。
  2. イベント名(event_attributes.event_name の値)を検索します。
  3. Users セクションで、プロファイルが作成または更新されていることを確認します。
  4. Segment > Search Users に移動します。
  5. Customer ID、Email address、または Phone number を使用して顧客を検索します。
  6. User Profile で、イベントが正しい属性とタイムスタンプで表示されていることを確認します。
データ処理の遅延処理キューの影響により、イベントが MoEngage ダッシュボードに表示されるまでに 1〜2 分かかる場合があります。10 分経ってもイベントが表示されない場合は、コンソールログで API エラーを確認し、認証情報を検証してください。

ステップ 4: 本番環境にデプロイする

インフラストラクチャに応じてデプロイ方法を選択します。MoEngage は次のデプロイオプションをサポートしています。

Systemd サービス(Linux)

Linux サーバー上でコンシューマーを永続的なバックグラウンドプロセスとして実行するには、次の手順に従って設定します。1. サービスファイルを作成するnano を使用してサービスファイルを作成します。
必要に応じてパスを調整し、次の設定を追加します。
2. サービスを有効化して起動する次のコマンドを実行して、システムマネージャーの設定を再読み込みし、サービスを有効化します。
3. サービスを監視する次のコマンドを使用して、バックグラウンドプロセスの稼働状況と正常性を確認します。
最適な用途: 単一の Linux サーバー、シンプルな設定スケーラビリティ: 単一インスタンス

レート制限

プラットフォームの安定性を維持するため、MoEngage は Kafka データの取り込みを、ワークスペースあたり最大 500 リクエスト/秒(RPS)に制限しています。ワークスペースがこの制限を超えると、MoEngage は HTTP 429 (Too Many Requests) ステータスコードを返します。 制限に達した際にイベントが失われないよう、コンシューマースクリプトで 429 レスポンスを処理し、バックオフしてから再試行してください。サンプルスクリプトはサーバーエラー時にすでに再試行を行っているため、同じロジックを拡張してレート制限のレスポンスにも対応させてください。

トラブルシューティング

よくある問題と解決策

その他のリソース