Este tópico descreve como usar o SDK para C# na conexão com o ApsaraMQ for Kafka para envio e recebimento de mensagens.
Requisitos de ambiente
Instale o .NET. Para mais informações, consulte Instale o .NET.
Instale a biblioteca de dependências do C#
Execute o comando a seguir para instalar a biblioteca de dependências do C#.
dotnet add package -v 1.5.2 Confluent.Kafka
Crie arquivos de configuração
(Opcional) Baixe o certificado raiz SSL (Secure Sockets Layer). A instalação do certificado é obrigatória caso você utilize o endpoint SSL para conectar-se à instância do ApsaraMQ for Kafka.
-
Configure os arquivos producer.cs e consumer.cs.
Tabela 1. Parâmetros Parâmetro
Descrição
BootstrapServers
Endpoint SSL da instância do ApsaraMQ for Kafka. Obtenha o endpoint na seção Endpoint Information da página Instance Details no ApsaraMQ for Kafka console.
SslCaLocation
Caminho do certificado raiz SSL baixado. Este parâmetro é necessário apenas se você utilizar o endpoint SSL para conectar-se à instância do ApsaraMQ for Kafka.
SaslMechanism
Mecanismo de segurança para envio e recebimento de mensagens.
Ao usar o endpoint SSL para conectar-se à instância do ApsaraMQ for Kafka, defina este parâmetro como SaslMechanism.Plain.
Para conexões via endpoint SASL (Simple Authentication and Security Layer), defina como SaslMechanism.Plain para o mecanismo PLAIN ou como SaslMechanism.ScramSha256 para o mecanismo SCRAM (Salted Challenge Response Authentication Mechanism).
SecurityProtocol
Protocolo de segurança para envio e recebimento de mensagens.
Se a conexão utilizar o endpoint SSL, configure este parâmetro como SecurityProtocol.SaslSsl.
Caso utilize o endpoint SASL, defina como SecurityProtocol.SaslPlaintext tanto para o mecanismo PLAIN quanto para o mecanismo SCRAM.
SaslUsername
Nome de usuário do SASL. Este parâmetro não está disponível ao utilizar o endpoint padrão para conectar-se à instância do ApsaraMQ for Kafka.
NotaQuando o recurso ACL não estiver ativado na instância do ApsaraMQ for Kafka, obtenha o nome de usuário e a senha do SASL nos parâmetros Username e Password, localizados na seção Configuration Information da página Instance Details no ApsaraMQ for Kafka console.
Se o recurso ACL estiver ativado na instância do ApsaraMQ for Kafka, garanta que o usuário SASL tenha autorização para enviar e receber mensagens por meio da instância. Para mais informações, consulte Conceder permissões a usuários SASL.
SaslPassword
Senha do usuário SASL. Indisponível quando a conexão com a instância do ApsaraMQ for Kafka utiliza o endpoint padrão.
topic
Nome do tópico. Consulte o nome do tópico na página Topics no ApsaraMQ for Kafka console.
GroupId
ID do grupo. Encontre o ID do grupo na página Groups no ApsaraMQ for Kafka console.
Envie mensagens
Execute o comando a seguir para executar o arquivo producer.cs e enviar mensagens:
dotnet run producer.cs
O exemplo a seguir mostra o conteúdo do arquivo producer.cs.
Para obter informações sobre os parâmetros do código de exemplo, consulte Parâmetros.
O código de exemplo utiliza o endpoint SSL. Exclua ou modifique o código relacionado aos parâmetros conforme o endpoint usado para conectar-se à instância do ApsaraMQ for Kafka.
using System;
using Confluent.Kafka;
class Producer
{
public static void Main(string[] args)
{
var conf = new ProducerConfig {
BootstrapServers = "XXX,XXX,XXX",
SslCaLocation = "XXX/only-4096-ca-cert.pem",
SaslMechanism = SaslMechanism.Plain,
SecurityProtocol = SecurityProtocol.SaslSsl,
SslEndpointIdentificationAlgorithm = SslEndpointIdentificationAlgorithm.None,
SaslUsername = "XXX",
SaslPassword = "XXX",
};
Action<DeliveryReport<Null, string>> handler = r =>
Console.WriteLine(!r.Error.IsError
? $"Delivered message to {r.TopicPartitionOffset}"
: $"Delivery Error: {r.Error.Reason}");
string topic ="XXX";
using (var p = new ProducerBuilder<Null, string>(conf).Build())
{
for (int i=0; i<100; ++i)
{
p.Produce(topic, new Message<Null, string> { Value = i.ToString() }, handler);
}
p.Flush(TimeSpan.FromSeconds(10));
}
}
}
Receba mensagens
Execute o comando a seguir para executar o arquivo consumer.cs e receber mensagens:
dotnet run consumer.cs
O exemplo a seguir mostra o conteúdo do arquivo consumer.cs.
Para obter informações sobre os parâmetros do código de exemplo, consulte Parâmetros.
O código de exemplo utiliza o endpoint SSL. Exclua ou modifique o código relacionado aos parâmetros conforme o endpoint usado para conectar-se à instância do ApsaraMQ for Kafka.
using System;
using System.Threading;
using Confluent.Kafka;
class Consumer
{
public static void Main(string[] args)
{
var conf = new ConsumerConfig {
GroupId = "XXX",
BootstrapServers = "XXX,XXX,XXX",
SslCaLocation = "XXX/only-4096-ca-cert.pem",
SaslMechanism = SaslMechanism.Plain,
SslEndpointIdentificationAlgorithm = SslEndpointIdentificationAlgorithm.None,
SecurityProtocol = SecurityProtocol.SaslSsl,
SaslUsername = "XXX",
SaslPassword = "XXX",
AutoOffsetReset = AutoOffsetReset.Earliest
};
string topic = "XXX";
using (var c = new ConsumerBuilder<Ignore, string>(conf).Build())
{
c.Subscribe(topic);
CancellationTokenSource cts = new CancellationTokenSource();
Console.CancelKeyPress += (_, e) => {
e.Cancel = true;
cts.Cancel();
};
try
{
while (true)
{
try
{
var cr = c.Consume(cts.Token);
Console.WriteLine($"Consumed message '{cr.Value}' at: '{cr.TopicPartitionOffset}'.");
}
catch (ConsumeException e)
{
Console.WriteLine($"Error occured: {e.Error.Reason}");
}
}
}
catch (OperationCanceledException)
{
c.Close();
}
}
}
}