get schema versions
This commit is contained in:
12
main.go
12
main.go
@@ -3,11 +3,12 @@ package main
|
||||
|
||||
import (
|
||||
config "com.navi.medici.janus/config"
|
||||
// producer_module "com.navi.medici.janus/producer"
|
||||
producer_module "com.navi.medici.janus/producer"
|
||||
server "com.navi.medici.janus/server"
|
||||
|
||||
"fmt"
|
||||
"log"
|
||||
"strings"
|
||||
"github.com/spf13/viper"
|
||||
)
|
||||
|
||||
@@ -33,13 +34,20 @@ func init() {
|
||||
fmt.Printf("Unable to decode into struct, %v", err)
|
||||
}
|
||||
|
||||
|
||||
configuration.Kafka.Bootstrap_Servers = viper.GetString(configuration.Kafka.Bootstrap_Servers)
|
||||
configuration.Kafka.Sasl_User = viper.GetString(configuration.Kafka.Sasl_User)
|
||||
configuration.Kafka.Sasl_Password = viper.GetString(configuration.Kafka.Sasl_Password)
|
||||
configuration.SchemaRegistry.Endpoint = viper.GetString(configuration.SchemaRegistry.Endpoint)
|
||||
configuration.SchemaRegistry.Topics = viper.GetString(configuration.SchemaRegistry.Topics)
|
||||
|
||||
port = configuration.Server.Port
|
||||
log.Printf("PORT IS: ", port)
|
||||
log.Printf(configuration.Kafka.Bootstrap_Servers)
|
||||
log.Printf(configuration.Kafka.Sasl_User)
|
||||
log.Printf(configuration.Kafka.Sasl_Password)
|
||||
log.Printf(configuration.SchemaRegistry.Endpoint)
|
||||
|
||||
producer_module.GetSchemaVersions(configuration.SchemaRegistry.Endpoint, strings.Split(configuration.SchemaRegistry.Topics, ","))
|
||||
|
||||
// sync producer using sarama
|
||||
// sync_producer := producer_module.GetSyncProducer(configuration.Kafka)
|
||||
|
||||
Reference in New Issue
Block a user