优化
This commit is contained in:
parent
783c289abe
commit
ba035974b5
@ -57,9 +57,6 @@ func MysqlConsumer() {
|
|||||||
|
|
||||||
func AddMetaEvent(kafkaData model.KafkaData) (err error) {
|
func AddMetaEvent(kafkaData model.KafkaData) (err error) {
|
||||||
if kafkaData.ReportType == model.EventReportType {
|
if kafkaData.ReportType == model.EventReportType {
|
||||||
redisConn := db.RedisPool.Get()
|
|
||||||
defer redisConn.Close()
|
|
||||||
|
|
||||||
b := bytes.Buffer{}
|
b := bytes.Buffer{}
|
||||||
b.WriteString(kafkaData.TableId)
|
b.WriteString(kafkaData.TableId)
|
||||||
b.WriteString("_")
|
b.WriteString("_")
|
||||||
@ -128,7 +125,6 @@ func AddTableColumn(kafkaData model.KafkaData, failFunc func(data consumer_data.
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
b := bytes.Buffer{}
|
b := bytes.Buffer{}
|
||||||
|
|
||||||
obj.Visit(func(key []byte, v *fastjson.Value) {
|
obj.Visit(func(key []byte, v *fastjson.Value) {
|
||||||
@ -207,7 +203,6 @@ func AddTableColumn(kafkaData model.KafkaData, failFunc func(data consumer_data.
|
|||||||
}()
|
}()
|
||||||
|
|
||||||
})
|
})
|
||||||
|
|
||||||
if foundNewKey {
|
if foundNewKey {
|
||||||
dims, err = sinker.ChangeSchema(newKeys, model.GlobConfig.Comm.ClickHouse.DbName, tableName, dims)
|
dims, err = sinker.ChangeSchema(newKeys, model.GlobConfig.Comm.ClickHouse.DbName, tableName, dims)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
@ -245,6 +245,7 @@ func main() {
|
|||||||
markFn()
|
markFn()
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
//生成表名
|
//生成表名
|
||||||
tableName := kafkaData.GetTableName()
|
tableName := kafkaData.GetTableName()
|
||||||
|
|
||||||
@ -260,6 +261,7 @@ func main() {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
//添加元数据
|
//添加元数据
|
||||||
if err := action.AddMetaEvent(kafkaData); err != nil {
|
if err := action.AddMetaEvent(kafkaData); err != nil {
|
||||||
logs.Logger.Error("addMetaEvent err", zap.Error(err))
|
logs.Logger.Error("addMetaEvent err", zap.Error(err))
|
||||||
|
@ -50,7 +50,6 @@ func GetDims(database, table string, excludedColumns []string, conn *sqlx.DB) (d
|
|||||||
|
|
||||||
redisConn := db.RedisPool.Get()
|
redisConn := db.RedisPool.Get()
|
||||||
defer redisConn.Close()
|
defer redisConn.Close()
|
||||||
|
|
||||||
dimsBytes, redisErr := redis.Bytes(redisConn.Do("get", dimsCachekey))
|
dimsBytes, redisErr := redis.Bytes(redisConn.Do("get", dimsCachekey))
|
||||||
|
|
||||||
if redisErr == nil && len(dimsBytes) != 0 {
|
if redisErr == nil && len(dimsBytes) != 0 {
|
||||||
|
Loading…
x
Reference in New Issue
Block a user