Для получения доступа к Kafka необходимо запросить через техподдержку РосДомофон следующую информацию:
Настройка осуществляется примерно следующим образом:
'metadata.broker.list' => 'kafka.rosdomofon.com:443',
'group.id' => 'имя группы',
'security.protocol' => 'SASL_SSL',
'sasl.mechanisms' => "SCRAM-SHA-512",
'ssl.certificate.location' => DIR . '/путь/до/сертификата/client.cer.pem',
'ssl.ca.location' => DIR . '/путь/до/сертификата/client.cer.pem',
'sasl.password' => 'имя пользователя',
'sasl.username' => 'пароль пользователя',
package main
import (
"crypto/tls"
"crypto/x509"
"flag"
"io/ioutil"
"log"
"os"
"strings"
"github.com/Shopify/sarama"
)
// keytool -importkeystore -srckeystore kafka.client.truststore.jks -destkeystore client.p12 -deststoretype PKCS12
// openssl pkcs12 -in client.p12 -nokeys -out client.cer.pem
func init() {
sarama.Logger = log.New(os.Stdout, "[Sarama] ", log.LstdFlags)
}
var (
brokers = flag.String("brokers", "kafka.rosdomofon.com:443", "The Kafka brokers to connect to, as a comma separated list")
userName = flag.String("username", "<username>", "The SASL username")
passwd = flag.String("passwd", "<password>", "The SASL password")
algorithm = flag.String("algorithm", "sha512", "The SASL SCRAM SHA algorithm sha256 or sha512 as mechanism")
certFile = flag.String("certificate", "ssl/kafka-client.cer.pem", "The optional certificate file for client authentication")
keyFile = flag.String("key", "", "The optional key file for client authentication")
caFile = flag.String("ca", "", "The optional certificate authority file for TLS client authentication")
verifySSL = flag.Bool("verify", true, "Optional verify ssl certificates chain")
useTLS = flag.Bool("tls", true, "Use TLS to communicate with the cluster")
logger = log.New(os.Stdout, "[Producer] ", log.LstdFlags)
)
func createTLSConfiguration() (t *tls.Config) {
t = &tls.Config{
InsecureSkipVerify: *verifySSL,
}
if *certFile != "" && *keyFile != "" && *caFile != "" {
cert, err := tls.LoadX509KeyPair(*certFile, *keyFile)
if err != nil {
log.Fatal(err)
}
caCert, err := ioutil.ReadFile(*caFile)
if err != nil {
log.Fatal(err)
}
caCertPool := x509.NewCertPool()
caCertPool.AppendCertsFromPEM(caCert)
t = &tls.Config{
Certificates: []tls.Certificate{cert},
RootCAs: caCertPool,
InsecureSkipVerify: *verifySSL,
}
}
return t
}
func main() {
flag.Parse()
if *brokers == "" {
log.Fatalln("at least one brocker is required")
}
if *userName == "" {
log.Fatalln("SASL username is required")
}
if *passwd == "" {
log.Fatalln("SASL password is required")
}
conf := sarama.NewConfig()
conf.Producer.Retry.Max = 1
conf.Producer.RequiredAcks = sarama.WaitForAll
conf.Producer.Return.Successes = true
conf.Metadata.Full = true
conf.Version = sarama.V0_10_0_0
conf.ClientID = "sasl_scram_client"
conf.Metadata.Full = true
conf.Net.SASL.Enable = true
conf.Net.SASL.User = *userName
conf.Net.SASL.Password = *passwd
conf.Net.SASL.Handshake = true
if *algorithm == "sha512" {
conf.Net.SASL.SCRAMClientGeneratorFunc = func() sarama.SCRAMClient { return &XDGSCRAMClient{HashGeneratorFcn: SHA512} }
conf.Net.SASL.Mechanism = sarama.SASLMechanism(sarama.SASLTypeSCRAMSHA512)
} else if *algorithm == "sha256" {
conf.Net.SASL.SCRAMClientGeneratorFunc = func() sarama.SCRAMClient { return &XDGSCRAMClient{HashGeneratorFcn: SHA256} }
conf.Net.SASL.Mechanism = sarama.SASLMechanism(sarama.SASLTypeSCRAMSHA256)
} else {
log.Fatalf("invalid SHA algorithm \"%s\": can be either \"sha256\" or \"sha512\"", *algorithm)
}
if *useTLS {
conf.Net.TLS.Enable = true
conf.Net.TLS.Config = createTLSConfiguration()
}
syncConsumer, err := sarama.NewConsumer(strings.Split(*brokers, ","), conf)
if err != nil {
logger.Fatalln("failed to create consumer: ", err)
}
topics, err := syncConsumer.Topics()
if err != nil {
logger.Fatalln("failed to get topics ", err)
}
logger.Printf("get topics list: %v", topics)
_ = syncConsumer.Close()
logger.Println("Bye now !")
}
public class KafkaConsumerConfig {
private static final String KAFKA_ADMIN_USER = System.getenv("KAFKA_USER_NAME");
private static final String KAFKA_ADMIN_PASSWORD = System.getenv("KAFKA_USER_PASSWORD");
private static final String KAFKA_TRUSTSTORE_LOCATION = System.getenv("KAFKA_TRUSTSTORE_LOCATION");
private static final String KAFKA_TRUSTSTORE_PASSWORD = System.getenv("KAFKA_TRUSTSTORE_PASSWORD");
@Value("${spring.kafka.bootstrap-servers}")
private String bootstrapServers;
@Bean
public Map<String, Object> consumerConfigs() {
Map<String, Object> props = new HashMap<>();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "<kafka_group>");
props.put(CommonClientConfigs.METADATA_MAX_AGE_CONFIG, 20000);
props.put("security.protocol", "SASL_SSL");
props.put("sasl.mechanism", "SCRAM-SHA-512");
props.put("ssl.truststore.type", "PKCS12");
props.put("sasl.jaas.config", "org.apache.kafka.common.security.scram.ScramLoginModule required username=\"" + KAFKA_USER_NAME + "\" password=\"" + KAFKA_USER_PASSWORD + "\";");
props.put("ssl.truststore.location", KAFKA_TRUSTSTORE_LOCATION);
props.put("ssl.truststore.password", KAFKA_TRUSTSTORE_PASSWORD);
props.put("ssl.endpoint.identification.algorithm", "");
return props;
}
@Bean
public ConsumerFactory<String, MessageRequest> consumerFactory() {
return new DefaultKafkaConsumerFactory<>(
consumerConfigs(),
new StringDeserializer(),
new JsonDeserializer<>(MessageRequest.class));
}
@Bean
public ConcurrentKafkaListenerContainerFactory<String, MessageRequest> kafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<String, MessageRequest> factory =
new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory());
return factory;
}
}