すべてのプロダクト
Search
ドキュメントセンター

ApsaraMQ for Kafka:Spring Cloudによるメッセージの送受信

最終更新日:Jul 11, 2026

Spring Cloud は、組み込みのサービスディスカバリー、設定管理、およびロードバランシングによってメッセージ駆動型マイクロサービスアプリケーションの構築を簡素化します。 ApsaraMQ for Kafka と統合して、分散システムでメッセージを送受信します。

前提条件

インターネットアクセス (認証と暗号化が必要)

インターネットアクセスでは、メッセージは SASL_SSL プロトコルを介して認証および暗号化されます。クライアントは SSL エンドポイントを介して ApsaraMQ for Kafka に接続します。エンドポイントの詳細については、「エンドポイントの比較」をご参照ください。

この例では、デモパッケージは/home/doc/project/aliware-kafka-demos/kafka-spring-stream-demoにアップロードされています。

  1. Linux システムにログインし、次のコマンドを実行してデモパッケージのディレクトリ/home/doc/project/aliware-kafka-demos/kafka-spring-stream-demoに移動します。

    cd /home/doc/project/aliware-kafka-demos/kafka-spring-stream-demo
  2. 次のコマンドを実行して、設定ファイルのパスに移動します。

    cd sasl-ssl/src/main/resources/
  3. 次のコマンドを実行してapplication.propertiesファイルを編集し、「パラメーターリスト」に基づいてインスタンス情報を設定します。

    vi application.properties
    ## プレースホルダーの値を実際のインスタンス情報に置き換えます。
    kafka.bootstrap-servers=alikafka-pre-cn-zv**********-1.alikafka.aliyuncs.com:9093,alikafka-pre-cn-zv**********-2.alikafka.aliyuncs.com:9093,alikafka-pre-cn-zv**********-3.alikafka.aliyuncs.com:9093
    kafka.consumer.group=test-spring
    kafka.output.topic.name=test-output
    kafka.input.topic.name=test-input
    kafka.ssl.truststore.location=/home/doc/project/aliware-kafka-demos/kafka-spring-stream-demo/sasl-ssl/src/main/resources/kafka.client.truststore.jks
    
    ### 次のバインディングパラメーターは、ApsaraMQ for KafkaをSpring Cloud Streamバインダーに関連付けます。デフォルト値を使用できます。
    spring.cloud.stream.bindings.MyOutput.destination=${kafka.output.topic.name}
    spring.cloud.stream.bindings.MyOutput.contentType=text/plain
    spring.cloud.stream.bindings.MyInput.group=${kafka.consumer.group}
    spring.cloud.stream.bindings.MyInput.destination=${kafka.input.topic.name}
    spring.cloud.stream.bindings.MyInput.contentType=text/plain
    
    ### バインダーは、メッセージミドルウェア用のSpring Cloud抽象化です。次のパラメーターにはデフォルト値を使用できます。
    spring.cloud.stream.kafka.binder.autoCreateTopics=false
    spring.cloud.stream.kafka.binder.brokers=${kafka.bootstrap-servers}
    spring.cloud.stream.kafka.binder.configuration.security.protocol=SASL_SSL
    spring.cloud.stream.kafka.binder.configuration.sasl.mechanism=PLAIN
    spring.cloud.stream.kafka.binder.configuration.ssl.truststore.location=${kafka.ssl.truststore.location}
    spring.cloud.stream.kafka.binder.configuration.ssl.truststore.password=KafkaOnsClient
    ### デモに存在しない場合は、このパラメーターを追加してサーバーホスト名の検証を無効にします。
    ### サーバーホスト名の検証では、SSL証明書のホスト名がサーバーのホスト名と一致するかどうかを確認します。デフォルト値はHTTPSです。
    spring.cloud.stream.kafka.binder.configuration.ssl.endpoint.identification.algorithm=
    表 1. パラメーターリスト

    パラメーター

    説明

    kafka.bootstrap-servers

    ApsaraMQ for Kafka インスタンスのエンドポイントです。このエンドポイントは、ApsaraMQ for Kafka コンソールインスタンスの詳細 ページにある アクセスポイント情報 セクションから取得できます。

    kafka.consumer.group

    メッセージのサブスクライブに使用するグループは、ApsaraMQ for Kafka コンソールGroup の管理 ページで作成できます。 詳細については、「ステップ 3: リソースの作成」をご参照ください。

    kafka.output.topic.name

    送信メッセージのトピック。デモアプリケーションは、固定コンテンツのメッセージをこのトピックに定期的に送信します。トピックは、ApsaraMQ for Kafkaコンソールトピック管理 ページで作成できます。詳細については、「手順3:リソースの作成」をご参照ください。

    kafka.input.topic.name

    受信メッセージのトピック。コンソールからこのトピックにメッセージを送信できます。デモアプリケーションはこれらのメッセージを消費し、ログに出力します。

    kafka.ssl.truststore.location

    SSLルート証明書kafka.client.truststore.jksへのパス。

  4. 次のコマンドを実行してkafka_client_jaas.confファイルを開き、インスタンスのユーザー名とパスワードを設定します。

    vi kafka_client_jaas.conf
    説明
    • インスタンスでアクセス制御リスト (ACL) 機能が無効になっている場合、デフォルトユーザーのユーザー名とパスワードをApsaraMQ for Kafkaコンソールのインスタンスの詳細 ページから取得できます。

    • インスタンスでアクセス制御リスト (ACL) 機能が有効になっている場合は、SASLユーザーがPLAINメカニズムを使用し、メッセージを送受信する権限を持っていることを確認してください。詳細については、「SASLユーザーへの権限付与」をご参照ください。

    KafkaClient {
      org.apache.kafka.common.security.plain.PlainLoginModule required
      username="your-username"
      password="your-password";
    };
  5. /home/doc/project/aliware-kafka-demos/kafka-spring-stream-demo/sasl-sslディレクトリに移動し、次のコマンドを実行してデモを実行します。

    sh run_demo.sh

    プログラムは次の情報を出力します。これは、プログラムがkafka.output.topic.nameで指定されたトピックにメッセージを送信したことを示します。

    Send: hello world !!
    Send: hello world !!
    Send: hello world !!
    Send: hello world !!
  6. ApsaraMQ for Kafkaコンソールにログインし、メッセージが正常に送受信されたことを確認します。

    • kafka.output.topic.nameパラメーターに指定されたトピックが、デモプログラムからメッセージを受信したかどうかを確認します。詳細については、「メッセージのクエリ」をご参照ください。

    • kafka.input.topic.nameで設定されているトピックにメッセージを送信し、デモプログラムのログを確認し、メッセージが出力されることを確認します。詳細については、「メッセージの送信」をご参照ください。

VPC アクセス (認証および暗号化なし)

VPC 内では、メッセージは認証や暗号化なしで PLAINTEXT プロトコルを介して送信されます。クライアントは、デフォルトエンドポイントを介して ApsaraMQ for Kafka に接続します。詳細については、「エンドポイントの比較」をご参照ください。

この例では、デモパッケージは/home/doc/project/aliware-kafka-demos/kafka-spring-stream-demoディレクトリにアップロードされています。

  1. Linux システムにログインし、次のコマンドを実行して、デモパッケージが含まれている/home/doc/project/aliware-kafka-demos/kafka-spring-stream-demoディレクトリに移動します。

    cd /home/doc/project/aliware-kafka-demos/kafka-spring-stream-demo
  2. 次のコマンドを実行して、設定ファイルのパスに移動します。

    cd vpc/src/main/resources/
  3. 次のコマンドを実行してapplication.propertiesファイルを編集し、「パラメーターリスト」に基づいてインスタンス情報を設定します。

    vi application.properties
    ### 実際のインスタンス情報に基づいて、次のパラメーターを変更します。
    kafka.bootstrap-servers=alikafka-pre-cn-zv**********-1-vpc.alikafka.aliyuncs.com:9092,alikafka-pre-cn-zv**********-2-vpc.alikafka.aliyuncs.com:9092,alikafka-pre-cn-zv**********-3-vpc.alikafka.aliyuncs.com:9092
    kafka.consumer.group=test-spring
    kafka.output.topic.name=test-output
    kafka.input.topic.name=test-input
  4. /home/doc/project/aliware-kafka-demos/kafka-spring-stream-demo/vpcディレクトリに移動し、次のコマンドを実行してデモを実行します。

    sh run_demo.sh

    プログラムは次の出力を表示します。

    Send: hello world !!
    Send: hello world !!
    Send: hello world !!
    Send: hello world !!
  5. ApsaraMQ for Kafkaコンソールにログインし、メッセージが正常に送受信されたことを確認します。

    • kafka.output.topic.nameパラメーターに指定されたトピックが、デモプログラムからメッセージを受信したかどうかを確認します。詳細については、「メッセージのクエリ」をご参照ください。

    • kafka.input.topic.nameで設定されているトピックにメッセージを送信し、デモプログラムのログを確認し、メッセージが出力されることを確認します。詳細については、「メッセージの送信」をご参照ください。

関連ドキュメント

Spring Cloud フレームワークの詳細については、「Spring Cloud Stream」をご参照ください。