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 |
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 :
Téléchargez et extrayez le dépôt.
Accédez au dossier
kafka-ruby-demoet ouvrez le sous-dossier correspondant à votre type d’endpoint (SSL ou par défaut).Modifiez les fichiers
producer.rubyetconsumer.rubyavec vos paramètres de connexion.Transférez les fichiers vers votre serveur. Si vous utilisez l’endpoint SSL, placez le fichier
cert.pemdans le même répertoire.