Skip to content

Apache KafkaへMQTTデータをストリームする ​

Apache Kafkaは、高スループットかつリアルタイムのデータ処理を目的とした広く利用されているオープンソースの分散イベントストリーミングプラットフォームです。しかし、Kafkaクライアントは安定したネットワーク接続と高いシステムリソースを必要とするため、エッジIoT通信には適していません。IoTシナリオでは、デバイスは一般的に軽量なMQTTプロトコルを使用して、不安定なネットワーク上でも効率的にデータを送信します。

EMQXはMQTTとKafka/Confluentを統合し、IoTデバイスとバックエンドシステム間のシームレスなデータストリーミングを可能にします。MQTTメッセージはKafkaトピックに取り込まれ、リアルタイム処理、保存、分析に利用される一方で、KafkaトピックのデータはMQTTクライアントに配信され、タイムリーなアクションをトリガーできます。

kafka_bridge

本ページではEMQXとKafkaのデータ統合について紹介し、統合の作成と検証手順を段階的に解説します。

動作概要 ​

Apache Kafkaとのデータ統合はEMQXの組み込み機能であり、MQTTベースのIoTデータをKafkaにストリームして下流処理や分析を可能にします。組み込みのルールエンジンを活用することで、カスタムコードなしにデータのフィルタリング、変換、ルーティングが可能です。

以下の図は、自動車IoTシナリオにおける典型的なEMQX–Kafka統合アーキテクチャを示しています。

kafka_architecture

Apache Kafkaへデータを流入または流出させるには、Kafka Sink(Kafkaへメッセージを送信)またはKafka Source(Kafkaからメッセージを受信)を作成します。以下はKafka Sinkのワークフローです。

  1. メッセージ取り込み: 車両に接続されたIoTデバイスはEMQXにMQTT接続を確立し、定期的に状態データを含むメッセージをパブリッシュします。EMQXがメッセージを受信すると、ルールエンジンでルールマッチングが開始されます。
  2. ルールベース処理: マッチしたルールにより、ペイロードのフィルタリング、変換、強化などが行われます。
  3. Kafkaへのデータ転送: ルールエンジンで定義されたルールがアクションをトリガーし、メッセージをKafkaに転送します。Kafka Sinkを使用してMQTTトピックを事前定義されたKafkaトピックにマッピングし、処理済みメッセージとデータをKafkaトピックに書き込みます。

Kafkaにデータが取り込まれた後は、以下のように複数の方法で消費・処理できます。

  • バックエンドサービスがKafkaトピックからリアルタイムデータストリームを直接消費。
  • Kafka Streamsを利用したリアルタイム集計、相関分析、解析。
  • Kafka Connectを使い、MySQLやElasticsearchなど外部システムへデータ転送し保存・追加処理。

特長と利点 ​

Apache Kafkaとのデータ統合は以下の特長と利点を提供します。

  • 信頼性の高い双方向IoTデータメッセージング: EMQXは不安定なネットワーク環境でもMQTTメッセージをKafkaに確実に転送し、バックエンドからのKafkaメッセージを接続されたIoTクライアントに届けます。
  • ペイロード変換: メッセージはKafkaに転送する前にSQLルールでフィルタリング、強化、変換が可能です。
  • 柔軟なトピックマッピング: MQTTトピックやユーザープロパティをKafkaトピックやヘッダーに柔軟にマッピングでき、1対1、1対多、ワイルドカードベースのマッピングをサポートします。
  • 柔軟なパーティション選択戦略: MQTTトピックやクライアントに基づき、同じKafkaパーティションへメッセージを転送します。
  • 高スループット処理: 同期・非同期のKafka書き込みをサポートし、レイテンシとスループットのバランスを異なるワークロードに応じて調整可能です。
  • ランタイムメトリクス: 各SinkおよびSourceの総メッセージ数、成功/失敗数、現在のレートなどのランタイムメトリクスを表示可能です。
  • 動的設定: ダッシュボードまたは設定ファイルでSinkおよびSourceを動的に設定できます。

これらの機能により、効率的なデータ取り込みと管理を備えたスケーラブルでレジリエントなIoTデータプラットフォームを構築できます。

はじめる前に ​

このセクションでは、EMQXダッシュボードでKafka SinkおよびSourceを作成する前に必要な準備について説明します。

前提条件 ​

Kafkaサーバーのセットアップ ​

ここではmacOSを例にインストールと起動方法を示します。以下のコマンドでKafkaをインストール・起動できます。

bash
wget https://archive.apache.org/dist/kafka/3.3.1/kafka_2.13-3.3.1.tgz

tar -xzf  kafka_2.13-3.3.1.tgz

cd kafka_2.13-3.3.1

# KRaftモードでKafkaを起動
KAFKA_CLUSTER_ID="$(bin/kafka-storage.sh random-uuid)"

bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/kraft/server.properties

bin/kafka-server-start.sh config/kraft/server.properties

詳細な操作手順はKafkaドキュメントのクイックスタートを参照してください。

Kafkaトピックの作成 ​

EMQXでデータ統合を作成する前に、関連するKafkaトピックを作成してください。以下のコマンドでSink用のtesttopic-inとSource用のtesttopic-outの2つのトピックを作成します。

bash
bin/kafka-topics.sh --create --topic testtopic-in --bootstrap-server localhost:9092

bin/kafka-topics.sh --create --topic testtopic-out --bootstrap-server localhost:9092

Kafkaプロデューサーコネクターの作成 ​

Kafka Sinkアクションを追加する前に、EMQXとKafka間の接続を確立するためのKafkaプロデューサーコネクターを作成する必要があります。

  1. EMQXダッシュボードで Integration -> Connector をクリックします。

  2. ページ右上の Create をクリックし、コネクター選択画面で Kafka Producer を選択して Next をクリックします。

  3. 名前と説明を入力します。例:my-kafka。名前はKafka Sinkとコネクターを関連付けるために使用され、クラスター内で一意である必要があります。

  4. Kafka接続に必要なパラメータを設定します。

    • Bootstrap Hosts: 127.0.0.1:9092を入力します。デモではEMQXとKafkaをローカルで実行している前提です。リモート環境の場合は適宜設定を調整してください。

    • Authentication: Kafkaクラスターの認証方式を選択します。以下の方式をサポートしています。

      • None: 認証なし。
      • AWS IAM for MSK: EMQXがAmazon EC2上で稼働し、Amazon MSKクラスターに接続する場合に使用。
      • MSK IAM Roles Anywhere: EC2外の環境からAmazon MSKに接続するためにAWS IAM Roles Anywhereクレデンシャルヘルパーを使用。
      • OAuth: OAuth 2.0ベースの認証を使用し、OAuthまたはOIDCをサポートするKafkaクラスターに接続。
      • Basic Auth: ユーザー名とパスワードによる認証。plain、scram_sha_256、scram_sha_512のいずれかのメカニズムを選択。
      • Kerberos: Kerberos (GSSAPI)認証。Kerberosプリンシパルとキータブファイルを指定。

      詳細は認証方式を参照してください。

    • 暗号化接続を確立する場合は、Enable TLS トグルをオンにします。TLS接続の詳細は外部リソースアクセスのTLSを参照してください。

    • Request Timeout: Kafkaからの応答を待つ最大時間(秒)。デフォルトは30秒。タイムアウト超過時は接続を再確立します。値が小さすぎると、Kafkaはリクエストを受け入れても応答を遅延させ、EMQXが再送することで重複メッセージや過剰な下流データが発生する可能性があります。

    • Advanced Settings(任意): 高度な設定を参照。

  5. Createをクリックする前に、Test ConnectionでKafkaサーバーへの接続が成功するか確認できます。

  6. Createをクリックしてコネクターの作成を完了します。

作成後、コネクターは自動的にKafkaに接続します。次に、このコネクターを基にルールを作成し、Kafkaクラスターへデータを転送します。

認証方式 ​

EMQXでKafkaコネクターを作成する際、Kafkaクラスターのセキュリティ設定に応じて複数の認証方式から選択できます。

  • None: 認証なし。

  • MSK IAM: EMQXがAmazon EC2上で稼働し、Amazon MSKクラスターに接続する場合に使用。

    AWS EC2インスタンスメタデータサービスを利用し、インスタンスに付与されたIAMポリシーに基づく認証トークンを生成します。

    重要なお知らせ

    MSK IAM認証は、EMQXがEC2インスタンス上で稼働しMSKクラスターに接続する場合のみサポートされます。これはEC2インスタンスメタデータサービスに依存しているためです。

    iptablesやnftablesでホストレベルのアウトバウンドフィルタリングを行う場合、169.254.169.254へのアクセスをブロックしないでください。EMQXはMSK IAM認証のためにインスタンスメタデータサービスにアクセスする必要があります。同様の例外はS3、S3 Tables、DynamoDB、KinesisなどEC2メタデータからクレデンシャルを取得するAWSベースの他のコネクターにも適用されます。詳細はルールエンジンポリシーとファイアウォールルールによるSSRF緩和を参照してください。

  • MSK IAM Roles Anywhere: EC2外の環境(オンプレミスなど)からAWS IAM Roles Anywhereクレデンシャルヘルパーを利用してAmazon MSKに接続する場合に使用。

    クレデンシャルヘルパープロセスはserveモードで起動し、EMQXにHTTP APIを公開します。EMQXはこのAPIから一時的なAWSクレデンシャルを取得し、SASL/OAUTHBEARERトークンを生成してMSK IAM認証に使用します。

    必要な設定:

    • Roles Anywhere Endpoint: クレデンシャルヘルパーのAPIエンドポイント。例: http://127.0.0.1:9911
    • AWS Region: MSKクラスターが稼働するAWSリージョン。
  • OAuth: OAuth 2.0ベースの認証で、OAuthまたはOIDCをサポートするKafkaクラスター(Confluent CloudやOAuth有効なセルフマネージドKafkaなど)に接続。

    EMQXはOAuth 2.0クライアントとして動作し、OAuth認可サーバーから定期的にアクセストークンを取得し、SASL/OAUTHBEARER機構でKafkaブローカーに認証します。

    必要な設定:

    • OAuth Grant Type: アクセストークン取得に使用するOAuth 2.0のグラントタイプ(現在はclient_credentialsのみサポート)。
    • OAuth Token Endpoint URI: トークンエンドポイントURI。
    • OAuth Client ID: OAuth認可サーバーに登録されたクライアントID。
    • OAuth Client Secret: クライアントIDに対応するシークレット。
    • OAuth Request Scope: (任意)トークンリクエストに含めるスコープ。
    • SASL Extensions: (高度、任意)認証時にSASL拡張として送信する追加のキー・バリュー。Confluent Cloudなど一部のKafkaプロバイダーでメタデータ(logicalClusterやidentityPoolIdなど)を渡すために必要。

    詳細はConfluent Cloudの公式ドキュメントを参照してください。

  • Basic Auth: ユーザー名とパスワードによる認証。

    必須項目:

    • Mechanism: plain、scram_sha_256、scram_sha_512から選択。
    • Username、Password: 認証情報。
  • Kerberos: Kerberos GSSAPI認証。

    必須項目:

    • Kerberos Principal: 認証に使用するKerberosプリンシパル。
    • Kerberos Keytab File: 非対話認証用のキータブファイルパス。

    重要なお知らせ

    KerberosキータブファイルはすべてのEMQXノードで同じパスに配置し、EMQXサービスユーザーが読み取り権限を持つ必要があります。

Kafka Sinkを使ったルールの作成 ​

このセクションでは、MQTTトピックt/#からのメッセージを処理し、Kafka Sinkを通じてKafkaのtesttopic-inトピックに送信するルールの作成方法を示します。

  1. EMQXダッシュボードで Integration -> Rules をクリックします。

  2. ページ右上の Create をクリックします。

  3. ルールIDを入力します。例:my_rule

  4. SQL Editorに以下のステートメントを入力します。これはトピックt/#のMQTTメッセージをKafkaに転送する例です。

    注意: 独自のSQLを指定する場合は、Sinkで必要なすべてのフィールドをSELECTに含めてください。

    sql
    SELECT
      *
    FROM
      "t/#"

    TIP

    初心者の方はSQL ExamplesやTry It OutをクリックしてSQLルールを学習・テストできます。

    TIP

    EMQX v5.7.2以降、ルールSQLで環境変数を読み取る機能が追加されました。詳細はルールSQLで環境変数を使うを参照してください。

  5. Create Ruleページで + Add Action をクリックし、ルールの出力アクションを定義します。

  6. Type of ActionドロップダウンからKafka Producerを選択します。

    ActionドロップダウンはデフォルトのCreate Actionのままにします。

    既存のSinkを選択することも可能ですが、この例では新規作成します。

  7. Nameと任意でDescriptionを入力します。

  8. Connectorドロップダウンから先ほど作成したmy-kafkaコネクターを選択します。必要に応じて新規作成も可能です。Kafkaプロデューサーコネクターの作成を参照してください。

  9. Sinkのデータ送信方法を設定します。

    • Kafka Topic: メッセージをパブリッシュするKafkaトピック。testtopic-inを入力します。EMQX v5.7.2以降、このフィールドは動的トピック設定もサポートします。変数テンプレートの使用を参照してください。
    • Kafka Headers: Kafkaメッセージに付加する任意のキー・バリューメタデータ。ヘッダー値はオブジェクトとして解決される必要があります。Kafka Header Value Encode Typeドロップダウンでエンコード方法を選択し、Addで複数ヘッダーを追加可能です。
    • Message Key: Kafkaメッセージのキー。パーティション分散やメッセージ順序付けに使用。静的文字列または${.clientid}などのプレースホルダーを含めることができます。
    • Message Value: Kafkaメッセージのペイロード。テンプレートからレンダリングされます。静的文字列または${.}のようなプレースホルダーを使い、ルールコンテキストから動的に生成可能です。テンプレートがNULL(例:参照フィールドが存在しない場合)を返した場合、空文字列ではなくKafkaのNULL値が生成されます。
    • Message Timestamp: Kafkaメッセージのタイムスタンプ。固定値または${timestamp}のようなプレースホルダーで動的に設定可能です。
    • Partition Strategy: プロデューサーがKafkaパーティションにメッセージを分配する方法を選択します。
    • Partitions Limit: プロデューサーがメッセージを送信できる最大パーティション数を制限します。有効にすると、すべてのパーティションではなく指定数のパーティション間でのみメッセージを分配します。
    • Compression: Kafkaメッセージのレコード圧縮・解凍に使用する圧縮アルゴリズムを指定します。
  10. フォールバックアクション(任意): メッセージ配信失敗時の信頼性向上のため、1つ以上のフォールバックアクションを定義できます。詳細はフォールバックアクションを参照してください。

  11. 高度な設定(任意): 高度な設定を参照してください。

  12. CreateをクリックしてSinkの作成を完了します。作成後、ページはCreate Ruleに戻り、新規Sinkがルールアクションに追加されます。

  13. Createをクリックしてルール作成を完了します。

kafka_producer_bridge

これでルールが正常に作成され、Integration -> Rulesページで新規ルールを確認でき、**Actions(Sink)**タブに新規KafkaプロデューサーSinkが表示されます。

また、Integration -> Flow Designerをクリックするとトポロジーを確認でき、トピックt/#のメッセージがルールmy_ruleで解析されKafkaに送信・保存される様子を直感的に把握できます。

Kafkaの動的トピック設定 ​

EMQX v5.7.2以降、KafkaプロデューサーSink設定で環境変数や変数テンプレートを使いKafkaトピックを動的に設定できます。このセクションでは2つのユースケースを紹介します。

環境変数の利用 ​

EMQX v5.7.2は、ルールSQL処理中に環境変数の値を動的に割り当てる機能を追加しました。これはルールエンジンの組み込みSQL関数getenvを使い、EMQXの環境変数を取得し、SQL処理結果に設定します。この機能を応用し、Kafka SinkルールアクションでKafkaトピック設定にルール出力結果のフィールドを参照できます。以下はその例です。

注意

ルールエンジンが読み取る環境変数名は、他のシステム環境変数の漏洩を防ぐため、必ずEMQXVAR_という固定プレフィックスを付ける必要があります。例えばgetenv('KAFKA_TOPIC')で読み取る変数名がKAFKA_TOPICの場合、環境変数名はEMQXVAR_KAFKA_TOPICに設定してください。

  1. Kafkaを起動し、testtopic-inトピックを事前作成します。はじめる前にを参照。

  2. EMQXを起動し環境変数を設定します。zip版インストールの場合、起動時に直接環境変数を指定可能です。例としてKafkaトピックtesttopic-inを環境変数EMQXVAR_KAFKA_TOPICに設定します。

    bash
    EMQXVAR_KAFKA_TOPIC=testtopic-in bin/emqx start
  3. コネクターを作成します。Kafkaプロデューサーコネクターの作成を参照。

  4. Kafka Sinkルールを設定し、SQL Editorに以下を入力します。

    sql
    SELECT
      getenv('KAFKA_TOPIC') as kafka_topic,
      payload
    FROM
      "t/#"

    kafka_dynamic_topic_sql

  5. SQLテストを有効化し、環境変数testtopic-inが正常に取得できることを確認します。

    kafka_dynamic_topic_sql_test

  6. KafkaプロデューサーSinkにアクションを追加します。ルールの右側Action OutputsでAdd Actionをクリック。

    • Connector: 先ほど作成したコネクターtest-kafkaを選択。
    • Kafka Topic: SQLルール出力の変数テンプレート${kafka_topic}形式で設定。

    kafka_dynamic_topic

  7. Kafka Sinkを使ったルールの作成を参照して追加設定を完了し、最後にCreateをクリックしてルール作成を完了します。

  8. Kafkaプロデューサールールのテストの手順に従い、Kafkaにメッセージを送信します。

    bash
    mqttx pub -h 127.0.0.1 -p 1883 -i pub -t t/Connection -q 1 -m 'payload string'

    Kafkaトピックtesttopic-inでメッセージを受信できるはずです。

    bash
    bin/kafka-console-consumer.sh --bootstrap-server 127.0.0.1:9092 \
      --topic testtopic-in
    
    {"payload":"payload string","kafka_topic":"testtopic-in"}
    {"payload":"payload string","kafka_topic":"testtopic-in"}

変数テンプレートの利用 ​

Kafka Topicフィールドに静的なトピック名を設定する代わりに、変数テンプレートを使って動的にトピックを生成できます。これによりメッセージ内容に基づきKafkaトピックを構築し、柔軟なメッセージ処理・振り分けが可能です。例えばdevice-${payload.device}のように指定すると、特定デバイスからのメッセージをdevice-1などのデバイスID付きトピックに簡単に送信できます。

この例では、Kafkaに送信するメッセージのペイロードにdeviceキーが含まれている必要があります。例:

json
{
    "topic": "t/devices/data",
    "payload": {
        "device": "1",
        "temperature": 25.6,
        "humidity": 60.2
    }
}

deviceキーがない場合、トピックのレンダリングに失敗し、メッセージが復旧不能な形でドロップされます。

また、Kafkaにはdevice-1、device-2など、解決されるすべてのトピックを事前作成しておく必要があります。存在しないトピック名に解決された場合もメッセージはドロップされます。

Kafkaプロデューサールールのテスト ​

Kafkaプロデューサールールが期待通りに動作するか、MQTTXを使ってMQTTメッセージをEMQXにパブリッシュするクライアントをシミュレートしてテストできます。

  1. MQTTXでトピックt/1にメッセージを送信します。
bash
mqttx pub -i emqx_c -t t/1 -m '{ "msg": "Hello Kafka" }'
  1. **Actions(Sink)**ページでSink名をクリックし統計情報を確認します。Sinkの稼働状況に新規の受信メッセージ数と送信メッセージ数が1件ずつ増えているはずです。

  2. 以下のコマンドでtesttopic-inトピックにメッセージが書き込まれているか確認します。

    bash
    bin/kafka-console-consumer.sh --bootstrap-server 127.0.0.1:9092  --topic testtopic-in

Kafkaコンシューマーコネクターの作成 ​

Kafka Sourceアクションを追加する前に、EMQXとKafka間の接続を確立するKafkaコンシューマーコネクターを作成する必要があります。

  1. EMQXダッシュボードで Integration -> Connector をクリックします。

  2. ページ右上の Create をクリックします。

  3. Create Connectorページで Kafka Consumer を選択し、Nextをクリックします。

  4. ソースの名前を入力します。英数字の組み合わせで、例:my-kafka-source。

  5. ソースの接続情報を入力します。

    • Bootstrap Hosts: 127.0.0.1:9092を入力します。EMQXとKafkaをローカルで実行している前提です。リモート環境の場合は適宜調整してください。

    • Authentication: Kafkaクラスターの認証方式を選択します。以下をサポートしています。

      • None: 認証なし。
      • AWS IAM for MSK: EC2上のEMQXからAmazon MSKに接続する場合。
      • MSK IAM Roles Anywhere: EC2外の環境からAmazon MSKに接続する場合。
      • OAuth: OAuth 2.0認証。
      • Basic Auth: Mechanism(plain、scram_sha_256、scram_sha_512)とUsername、Passwordを指定。
      • Kerberos: Kerberos PrincipalとKerberos Keytab Fileを指定。

      詳細は認証方式を参照。

    • 暗号化接続を確立する場合はEnable TLSをオンにします。詳細はTLS for External Resource Accessを参照。

    • Advanced Settings(任意): 高度な設定を参照。

  6. Createをクリックする前に、Test ConnectionでKafkaサーバーへの接続を確認できます。

  7. Createをクリックします。関連するルールの作成オプションが表示されます。KafkaコンシューマーSourceを使ったルールの作成を参照してください。

KafkaコンシューマーSourceを使ったルールの作成 ​

このセクションでは、KafkaコンシューマーSourceで転送されたメッセージをEMQXでさらに処理し、MQTTトピックに再パブリッシュするルールの作成方法を示します。

ルールSQLの作成 ​

  1. EMQXダッシュボードで Integration -> Rules をクリックします。

  2. ページ右上の Create をクリックします。

  3. ルールIDを入力します。例:my_rule

  4. Kafkaソース$bridges/kafka_consumer:<sourceName>から変換されたメッセージをEMQXに転送する場合、SQL Editorに以下を入力します。

    注意: 独自SQLを指定する場合、後続の再パブリッシュアクションで必要なすべてのフィールドをSELECTに含めてください。Kafka SourceのSELECT文ではts_type、topic、ts、event、headers、key、metadata、value、timestamp、offset、nodeなどのフィールドが利用可能です。

    sql
    SELECT
      *
    FROM
      "$bridges/kafka_consumer:<sourceName>"

    注意: 初心者はSQL ExamplesやEnable TestをクリックしてSQLルールを学習・テストできます。

KafkaコンシューマーSourceをデータ入力に追加 ​

  1. ルール作成ページ右側のData Inputsタブを選択し、Add Inputをクリックします。

  2. Input TypeドロップダウンからKafka Consumerを選択します。SourceドロップダウンはデフォルトのCreate Sourceのままか、既存のKafka Consumerソースを選択可能です。この例では新規作成してルールに追加します。

  3. ソースの名前と説明を入力します。

  4. Connectorドロップダウンから先ほど作成したmy-kafka-consumerコネクターを選択します。隣のボタンから新規コネクター作成も可能です。Kafkaコンシューマーコネクターの作成を参照してください。

  5. 以下のフィールドを設定します。

    • Kafka Topic: コンシューマーソースが購読するKafkaトピック。
    • Group ID: このソースのコンシューマーグループ識別子。未指定の場合はソース名に基づき自動生成されます。
    • Key Encoding Mode、Value Encoding Mode: Kafkaメッセージのキーと値のエンコードモードを選択。
  6. Offset Reset Policy: コンシューマーがKafkaトピックパーティションのどこから読み始めるかのポリシー。

    • latest: コンシューマー開始時点の最新オフセットから読み、過去のメッセージはスキップ。
    • earliest: パーティションの先頭から読み、過去のメッセージもすべて読み取る。
  7. Advanced Settings(任意): 高度な設定を参照。

  8. Createをクリックする前に、Test ConnectivityでKafkaサーバーへの接続を確認できます。

  9. Createをクリックしてソース作成を完了します。ルール作成ページのData Inputsタブに新規ソースが表示されます。

再パブリッシュアクションの追加 ​

  1. Action Outputsタブを選択し、+ Add Actionをクリックしてルールがトリガーするアクションを定義します。

  2. Type of ActionドロップダウンからRepublishを選択します。

  3. TopicおよびPayloadフィールドに再パブリッシュするメッセージのトピックとペイロードを入力します。例として t/1と${.}を入力します。

    • Topicフィールドには${}を使い動的にMQTTトピックを指定可能です。例:t/${key}(${}内のパラメータはSQLのSELECT文に含める必要があります)。
  4. Addをクリックしてアクションをルールに追加します。

  5. ルール作成ページに戻り、Saveをクリックします。

Kafka_consumer_rule

Kafka Sourceルールのテスト ​

Kafkaソースとルールが期待通りに動作するか、MQTTXを使ってEMQXのトピックをサブスクライブするクライアントをシミュレートし、KafkaプロデューサーでKafkaトピックにデータを生成してテストできます。EMQXがKafkaのデータをクライアントがサブスクライブするトピックに再パブリッシュするか確認します。

  1. MQTTXでトピックt/1をサブスクライブします。

    bash
    mqttx sub -t t/1 -v
  2. 新しいコマンドラインを開き、以下のコマンドでKafkaプロデューサーを起動します。

    bash
    bin/kafka-console-producer --bootstrap-server 127.0.0.1:9092 --topic testtopic-out

    メッセージ入力待ちになります。

  3. {"msg": "Hello EMQX"}を入力し、testtopic-outトピックにメッセージを生成します。

  4. MQTTXのサブスクリプションで、Kafkaからの以下のメッセージがトピックt/1で受信されることを確認します。

    json
    {
        "value": "{\"msg\": \"Hello EMQX\"}",
        "ts_type": "create",
        "ts": 1679665968238,
        "topic": "testtopic-out",
        "offset": 2,
        "key": "key",
        "headers": {
            "header_key": "header_value"
        }
    }

高度な設定 ​

このセクションでは、データ統合のパフォーマンス最適化や特定シナリオに応じたカスタマイズに役立つ高度な設定オプションを説明します。コネクター、Sink、Source作成時にAdvanced Settingsを展開し、ビジネス要件に応じて以下の設定を行えます。

フィールド名説明推奨値
Allow Auto Topic Creation(プロデューサーコネクターのみ)有効にすると、クライアントがメタデータ取得要求時に存在しないKafkaトピックを自動作成します。disabled
Min Metadata Refresh IntervalKafkaブローカーやトピックのメタデータ更新間隔の最小時間。小さすぎるとKafkaサーバーに不要な負荷がかかる可能性があります。3秒
Metadata Request TimeoutKafkaからメタデータを要求する際の最大待機時間。5秒
Connect TimeoutTCP接続確立の最大待機時間。認証時間も含みます。5秒
Max Wait Time (Source)Kafkaブローカーからのフェッチ応答を待つ最大時間。1秒
Fetch Bytes (Source)1回のフェッチ要求でKafkaから取得するバイト数。設定値がメッセージサイズ未満だとフェッチ性能に悪影響を与える可能性があります。896 KB
Max Batch Bytes (Sink)Kafkaバッチ内で収集可能なメッセージの最大バイト数。Kafkaブローカーのデフォルトは1MBですが、EMQXはエンコードオーバーヘッドを考慮しやや小さめに設定。単一メッセージが上限を超える場合は別バッチで送信されます。896 KB
Offset Commit Interval (Source)コンシューマーグループごとにオフセットコミット要求を送る間隔。5秒
Required Acks (Sink)Kafkaパーティションリーダーがフォロワーから待つ必要があるアックの種類。
all_isr: 全インシンクレプリカからのアックを要求。
leader_only: パーティションリーダーのみからのアックを要求。
none: Kafkaからのアック不要。
all_isr
Partition Count Refresh Interval (Source)Kafkaプロデューサーがパーティション数増加を検知する間隔。増加検知後、指定のpartition_strategyに基づき新パーティションにメッセージを分配。60秒
Max Inflight (Sink)Kafkaプロデューサーがアック受信前に送信可能な最大バッチ数(パーティション単位)。値が大きいほどスループット向上。ただし1より大きいとメッセージの順序入れ替わりリスクあり。10
Query Mode (Source)非同期または同期クエリモードを選択し、メッセージ伝送を最適化。非同期モードではKafka書き込みがMQTTパブリッシュ処理をブロックしませんが、クライアントがKafka到着前にメッセージを受信する可能性があります。Async
Synchronous Query Timeout (Sink)同期モード時の最大待機時間。メッセージ伝送完了を保証し長時間待機を防止。同期モード時のみ有効。5秒
Buffer Mode (Sink)メッセージ送信前のバッファリング方法。メモリバッファリングは送信速度向上に寄与。
memory: メモリにバッファ。EMQXノード再起動でメッセージは失われる。
disk: ディスクにバッファ。再起動後もメッセージ保持。
hybrid: 初期はメモリバッファ。一定容量超過時に順次ディスクにオフロード。メモリモード同様、再起動でメッセージは失われる。
memory
Per-partition Buffer Limit (Sink)Kafkaパーティションごとの最大バッファサイズ(バイト)。上限到達時は古いメッセージを破棄しバッファ領域を確保。メモリ使用量と性能のバランス調整に有効。2 GB
Segment File Bytes (Sink)バッファモードがdiskまたはhybridの場合に適用。メッセージ保存用分割ファイルのサイズを制御し、ディスクストレージの最適化に影響。100 MB
Memory Overload Protection (Sink)バッファモードがmemoryの場合に適用。メモリ圧迫時に古いメッセージを自動破棄し、システム安定性を確保。Linuxのみ有効。Enabled
Socket Send / Receive Buffer Sizeソケットバッファサイズを管理しネットワーク伝送性能を最適化。1024 KB
TCP KeepaliveKafkaブリッジ接続のTCPキープアライブ設定。長時間の非アクティブ状態による接続切断を防止。Idle, Interval, Probesの3つの数値をカンマ区切りで指定。例: 240,30,5は240秒アイドル後にキープアライブ開始、30秒間隔で最大5回プローブ送信。none
Max Batch Age (Sink)プロデューサーバッファ内のメッセージが送信されずに保持可能な最大時間。超過するとメッセージは破棄され、dropped.expiredメトリクスにカウント。デフォルトはinfinityで期限切れなし。バッファオーバーフロー時は期限切れに関わらず破棄される可能性あり。infinity
Max Retries (Sink)Kafkaがリトライ可能なエラー(例:パーティションリーダー変更)を返した際の最大リトライ回数。初回試行とリトライがすべて失敗するとバッチは破棄され、failedメトリクスにカウント。接続喪失による再送はリトライ回数にカウントされず、max_batch_ageで制限。デフォルトは無制限。infinity
Reconnect Delay (Sink)接続喪失後、プロデューサーがKafkaに再接続を試みるまでの遅延時間。切断中もメッセージはバッファに蓄積されるがバッファ制限やmax_batch_ageの影響を受ける。デフォルトは2秒。2秒
Max Linger Timeパーティションごとのプロデューサーがより大きなバッチを作るために待機する最大時間。すべてのバッファモードに適用。デフォルト0は待機なしでレイテンシ最適化。多少の遅延を許容できる場合は設定するとリクエスト数削減に寄与。ディスクバッファ時はバッチ書き込み前の待機時間。最低5ms推奨。0ミリ秒
Max Linger Bytesパーティションごとのプロデューサーがバッチ送信前に蓄積する最大バイト数。10 MB
Health Check Intervalコネクターの稼働状態をチェックする間隔。15秒

さらに詳しく ​

EMQXはApache Kafkaとのデータ統合に関する豊富な学習リソースを提供しています。以下のリンクから詳細を学べます。

ブログ:

ベンチマークレポート:

動画: