Tous les produits
Search
Centre de documentation

ApsaraMQ for Kafka:Send and receive messages with the Ruby SDK

Dernière mise à jour :Aug 12, 2026

Connectez une application Ruby à ApsaraMQ for Kafka pour produire et consommer des messages à l’aide du gem ruby-kafka.

Prérequis

Avant de commencer, assurez-vous de disposer des éléments suivants :

  • Ruby installé sur votre serveur

  • Une instance ApsaraMQ for Kafka avec un topic et un groupe de consommateurs créés

  • (Endpoint SSL uniquement) Le certificat racine SSL téléchargé et enregistré sous le nom cert.pem

Installer le gem ruby-kafka

gem install ruby-kafka -v 0.6.8

Récupérer les paramètres de connexion

Récupérez les valeurs suivantes dans la console ApsaraMQ for Kafka, puis remplacez les espaces réservés correspondants dans le code du producteur et du consommateur.

Endpoint SSL

Espace réservé Emplacement
<your-broker-host1:port,your-broker-host2:port> Instance Details > Endpoint Information -- copiez l’endpoint SSL
<your-topic-name> Page Topics
<your-sasl-username> Instance Details > Configuration Information > Username
<your-sasl-password> Instance Details > Configuration Information > Password
<your-consumer-group-id> Page Groups

Endpoint par défaut

Espace réservé Emplacement
<your-broker-host1:port,your-broker-host2:port> Instance Details > Endpoint Information -- copiez l’endpoint par défaut
<your-topic-name> Page Topics
<your-consumer-group-id> Page Groups
Remarque
  • Si la fonctionnalité ACL n’est pas activée sur votre instance, récupérez le nom d’utilisateur et le mot de passe SASL depuis les paramètres Username et Password de la section Configuration Information, sur la page Instance Details de la console ApsaraMQ for Kafka.

  • Si la fonctionnalité ACL est activée, vérifiez que l’utilisateur SASL dispose des autorisations requises. Pour plus d’informations, consultez Accorder des autorisations aux utilisateurs SASL.

Produire des messages

Créez un fichier nommé producer.ruby en utilisant l’un des exemples ci-dessous, selon votre type d’endpoint.

Endpoint SSL

# frozen_string_literal: true

$LOAD_PATH.unshift(File.expand_path("../../lib", __FILE__))

require "kafka"

logger = Logger.new($stdout)
logger.level = Logger::INFO

# Replace the placeholders with your actual values.
brokers = "<your-broker-host1:port,your-broker-host2:port>"
topic = "<your-topic-name>"
username = "<your-sasl-username>"
password = "<your-sasl-password>"

kafka = Kafka.new(
    seed_brokers: brokers,
    client_id: "sasl-producer",
    logger: logger,
    # Path to the SSL root certificate file.
    ssl_ca_cert: File.read('./cert.pem'),
    sasl_plain_username: username,
    sasl_plain_password: password,
    )

producer = kafka.producer

# Read lines from stdin and send each line as a message.
begin
    $stdin.each_with_index do |line, index|

    producer.produce(line, topic: topic)

    producer.deliver_messages
end

ensure
    # Flush any remaining messages and release resources.
    producer.deliver_messages

    producer.shutdown
end

Endpoint par défaut

# frozen_string_literal: true

$LOAD_PATH.unshift(File.expand_path("../../lib", __FILE__))

require "kafka"

logger = Logger.new($stdout)
logger.level = Logger::INFO

# Replace the placeholders with your actual values.
brokers = "<your-broker-host1:port,your-broker-host2:port>"
topic = "<your-topic-name>"

kafka = Kafka.new(
    seed_brokers: brokers,
    client_id: "simple-producer",
    logger: logger,
    )

producer = kafka.producer

# Read lines from stdin and send each line as a message.
begin
    $stdin.each_with_index do |line, index|

    producer.produce(line, topic: topic)

    producer.deliver_messages
end

ensure
    # Flush any remaining messages and release resources.
    producer.deliver_messages

    producer.shutdown
end

Exécutez le producteur :

ruby producer.ruby

Saisissez chaque message sur une nouvelle ligne et appuyez sur Enter. Le producteur envoie chaque ligne au topic.

Consommer des messages

Créez un fichier nommé consumer.ruby en utilisant l’un des exemples ci-dessous, selon votre type d’endpoint.

Endpoint SSL

# frozen_string_literal: true

$LOAD_PATH.unshift(File.expand_path("../../lib", __FILE__))

require "kafka"

logger = Logger.new(STDOUT)
logger.level = Logger::INFO

# Replace the placeholders with your actual values.
brokers = "<your-broker-host1:port,your-broker-host2:port>"
topic = "<your-topic-name>"
username = "<your-sasl-username>"
password = "<your-sasl-password>"
consumerGroup = "<your-consumer-group-id>"

kafka = Kafka.new(
        seed_brokers: brokers,
        client_id: "sasl-consumer",
        socket_timeout: 20,
        logger: logger,
        # Path to the SSL root certificate file.
        ssl_ca_cert: File.read('./cert.pem'),
        sasl_plain_username: username,
        sasl_plain_password: password,
        )

consumer = kafka.consumer(group_id: consumerGroup)
consumer.subscribe(topic, start_from_beginning: false)

# Stop the consumer gracefully on TERM or INT signals.
# To stop the process, run: kill -s TERM <process-id>
trap("TERM") { consumer.stop }
trap("INT") { consumer.stop }

begin
    consumer.each_message(max_bytes: 64 * 1024) do |message|
    logger.info("Get message: #{message.value}")
    end
rescue Kafka::ProcessingError => e
    # If message processing fails, pause the affected partition
    # for 20 seconds before retrying. This prevents tight retry
    # loops that could overwhelm the broker.
    warn "Got error: #{e.cause}"
    consumer.pause(e.topic, e.partition, timeout: 20)

    retry
end

Endpoint par défaut

# frozen_string_literal: true

$LOAD_PATH.unshift(File.expand_path("../../lib", __FILE__))

require "kafka"

logger = Logger.new(STDOUT)
logger.level = Logger::INFO

# Replace the placeholders with your actual values.
brokers = "<your-broker-host1:port,your-broker-host2:port>"
topic = "<your-topic-name>"
consumerGroup = "<your-consumer-group-id>"

kafka = Kafka.new(
        seed_brokers: brokers,
        client_id: "test",
        socket_timeout: 20,
        logger: logger,
        )

consumer = kafka.consumer(group_id: consumerGroup)
consumer.subscribe(topic, start_from_beginning: false)

# Stop the consumer gracefully on TERM or INT signals.
# To stop the process, run: kill -s TERM <process-id>
trap("TERM") { consumer.stop }
trap("INT") { consumer.stop }

begin
    consumer.each_message(max_bytes: 64 * 1024) do |message|
    logger.info("Get message: #{message.value}")
    end
rescue Kafka::ProcessingError => e
    # If message processing fails, pause the affected partition
    # for 20 seconds before retrying. This prevents tight retry
    # loops that could overwhelm the broker.
    warn "Got error: #{e.cause}"
    consumer.pause(e.topic, e.partition, timeout: 20)

    retry
end

Exécutez le consommateur :

ruby consumer.ruby

Le consommateur rejoint le groupe de consommateurs spécifié et récupère les nouveaux messages. Chaque message est journalisé sur stdout. Appuyez sur Ctrl+C pour arrêter le consommateur correctement.

Utiliser le projet de démonstration

Des fichiers de démonstration sont disponibles dans le dépôt aliware-kafka-demos :

  1. Téléchargez et extrayez le dépôt.

  2. Accédez au dossier kafka-ruby-demo et ouvrez le sous-dossier correspondant à votre type d’endpoint (SSL ou par défaut).

  3. Modifiez les fichiers producer.ruby et consumer.ruby avec vos paramètres de connexion.

  4. Transférez les fichiers vers votre serveur. Si vous utilisez l’endpoint SSL, placez le fichier cert.pem dans le même répertoire.

Rubriques connexes