|
|
@ -51,7 +51,8 @@ func (s *KafkaWriter) createTopic() error { |
|
|
|
// return era
|
|
|
|
// return era
|
|
|
|
// }
|
|
|
|
// }
|
|
|
|
|
|
|
|
|
|
|
|
// //
|
|
|
|
//
|
|
|
|
|
|
|
|
controllerConn := conn |
|
|
|
// controllerConn, eru := kafka.Dial("tcp", net.JoinHostPort(controller.Host, strconv.Itoa(controller.Port)))
|
|
|
|
// controllerConn, eru := kafka.Dial("tcp", net.JoinHostPort(controller.Host, strconv.Itoa(controller.Port)))
|
|
|
|
// if eru != nil {
|
|
|
|
// if eru != nil {
|
|
|
|
// return eru
|
|
|
|
// return eru
|
|
|
@ -67,7 +68,7 @@ func (s *KafkaWriter) createTopic() error { |
|
|
|
}, |
|
|
|
}, |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
return conn.CreateTopics(topicConfigs...) |
|
|
|
return controllerConn.CreateTopics(topicConfigs...) |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
func (s *KafkaWriter) checkTopic() error { |
|
|
|
func (s *KafkaWriter) checkTopic() error { |
|
|
|