package main import ( "bytes" "database/sql" "fmt" _ "github.com/go-sql-driver/mysql" influx "github.com/influxdata/influxdb1-client/v2" druid "github.com/shunfei/godruid" "log" "net/http" "os" "strconv" "time" ) // All the required ENV Vars must be defined here. type Config struct { JOB_NAME, TZ, WINDOW_SIZE string MYSQL_HOST, MYSQL_USER, MYSQL_PASSWORD, MYSQL_DB, MYSQL_TABLE string DRUID_HOST string INFLUX_HOST, INFLUX_PORT, INFLUX_DB, INFLUX_RP, INFLUX_MEASUREMENT string DVO_API_CHECKIN_ENDPOINT string } type Stats struct { Views, Clicks, Leads sql.NullFloat64 } // Initialize a Druid Client given host. // Args: druidHost - the hostname of druid instance // Returns: pointer to the initialized Druid client. func initDruidClient(druidHost string) *druid.Client { client := druid.Client{ Url: fmt.Sprintf("http://%s:8082", druidHost), Debug: true, } return &client } // Initialize a MySql Db. Fail if can't connect. // Args: appConfig - pointer to the app Config struct // Returns: pointer to the initialized MySql Db func initMySqlDb(appConfig *Config) *sql.DB { connectionStr := fmt.Sprintf("%s:%s@tcp(%s:3306)/%s", appConfig.MYSQL_USER, appConfig.MYSQL_PASSWORD, appConfig.MYSQL_HOST, appConfig.MYSQL_DB) db, err := sql.Open("mysql", connectionStr) if err != nil { log.Fatal(err) } err = db.Ping() if err != nil { log.Fatal(err) } log.Print("mysql datasource connected!") return db } // Initialize an InfluxDB client. Fail if can't connect. // Args: appConfig - pointer to the app Config // Returns: pointer to the initialized InfluxDB Client func initInfluxClient(appConfig *Config) *influx.Client { // Make client c, err := influx.NewHTTPClient(influx.HTTPConfig{ Addr: fmt.Sprintf("http://%s:%s", appConfig.INFLUX_HOST, appConfig.INFLUX_PORT), }) if err != nil { log.Println("Error creating InfluxDB Client") log.Fatal(err) } _, _, err = c.Ping(0) if err != nil { log.Println("Error pinging InfluxDB") log.Fatal(err) } return &c } // Builds and runs the Druid query, and stores the results into a Stats struct. // If the Druid query fails or returns nothing, then use 0 values for stats. // Args: client - pointer to the initialized druid client // interval - the date string to get the values for // Returns: Stats struct containing Druid values (Views, Clicks, Leads) func getDruidStatsForInterval(client *druid.Client, interval string) Stats { var druidStats Stats query := &druid.QueryGroupBy{ DataSource: "RevenueReport", Intervals: []string{fmt.Sprintf("%s/PT1H", interval)}, Granularity: druid.GranPeriod{Type: "period", Period: "PT1H", TimeZone: "America/Toronto"}, Aggregations: []druid.Aggregation{ druid.AggLongSum("Views", "Views"), druid.AggLongSum("Clicks", "Clicks"), druid.AggDoubleSum("Leads", "Leads"), }, } err := client.Query(query) if err != nil { log.Println("Error while running the Druid Query", err.Error()) } //log.Println("requst", client.LastRequest) //log.Println("response", client.LastResponse) if (err != nil) || len(query.QueryResult) == 0 { log.Println("Druid query returned no results! Using 0 values.") druidStats.Views = sql.NullFloat64{Float64: float64(0), Valid: true} druidStats.Clicks = sql.NullFloat64{Float64: float64(0), Valid: true} druidStats.Leads = sql.NullFloat64{Float64: float64(0), Valid: true} return druidStats } views := query.QueryResult[0].Event["Views"].(float64) clicks := query.QueryResult[0].Event["Clicks"].(float64) leads := query.QueryResult[0].Event["Leads"].(float64) druidStats.Views = sql.NullFloat64{Float64: views, Valid: true} druidStats.Clicks = sql.NullFloat64{Float64: clicks, Valid: true} druidStats.Leads = sql.NullFloat64{Float64: leads, Valid: true} return druidStats } // Prepare a sql statement to be used repeatedly later. // Args: db - pointer to the initialized sql DB // Returns: a sql Prepared statement func getSqlPreparedStatement(db *sql.DB, table string) *sql.Stmt { query := fmt.Sprintf(` SELECT SUM(Views) as views, SUM(Clicks) as clicks, SUM(Leads) as leads FROM %s WHERE Stamp = ? AND CampaignID <> 278258 `, table) stmt, err := db.Prepare(query) if err != nil { log.Println("Error while preparing sql query") log.Fatal(err) } return stmt } // Runs the MySql query from a prepared statement and stores the results into a Stats struct // Args: stmt - pointer to the sql prepared statement // interval - the date string to get the values for // Returns: Stats struct containing MySql values (Views, Clicks, Leads) func getMySqlStatsForInterval(stmt *sql.Stmt, interval string) Stats { var mySqlStats Stats //log.Println("running sql...") err := stmt.QueryRow(interval).Scan(&mySqlStats.Views, &mySqlStats.Clicks, &mySqlStats.Leads) if err != nil { log.Println("Error while running sql query") log.Fatal(err) } //log.Println(mySqlStats) return mySqlStats } // Write the MySql and Druid Stats to InfluxDB. // Args: influxClient - pointer to the initialized InfluxDB client // appConfig - pointer to the app Config // mySqlStats - Stats struct contianing MySql values (Views, Leads, Clicks) // druidStats - Stats struct contianing Druid values (Views, Leads, Clicks) // t - timestamp object rounded down to the hour // Returns: the name of the measurement that was written to func writeToInflux(influxClient *influx.Client, appConfig *Config, mySqlStats *Stats, druidStats *Stats, t time.Time) string { // Create a new InfluxDB point batch bp, _ := influx.NewBatchPoints(influx.BatchPointsConfig{ Database: appConfig.INFLUX_DB, Precision: "s", RetentionPolicy: appConfig.INFLUX_RP, }) // Create a single point tags := map[string]string{} fields := map[string]interface{}{ "m.views": int64(mySqlStats.Views.Float64), "m.clicks": int64(mySqlStats.Clicks.Float64), "m.leads": int64(mySqlStats.Leads.Float64), "d.views": int64(druidStats.Views.Float64), "d.clicks": int64(druidStats.Clicks.Float64), "d.leads": int64(druidStats.Leads.Float64), } pt, err := influx.NewPoint(appConfig.INFLUX_MEASUREMENT, tags, fields, t) if err != nil { log.Println("Failed to crate valid Influx Point.", err.Error()) } log.Println(pt.String()) // Add to batch bp.AddPoint(pt) // Write the batch err = (*influxClient).Write(bp) if err != nil { log.Println("Failed to write point to Influx.", err.Error()) } return appConfig.INFLUX_MEASUREMENT } // A helper to perform checkin for execution of this job. // Args: jobName - the name of this job // endPoint - the checkin api endpoint url func checkInJob(jobName string, endPoint string) { log.Println("Cheking-in to", endPoint) var jsonStr = []byte(`{"source":"` + jobName + `"}`) req, err := http.NewRequest("POST", endPoint, bytes.NewBuffer(jsonStr)) req.Header.Set("Content-Type", "application/json") client := &http.Client{} resp, err := client.Do(req) if err != nil { log.Println(err.Error()) } defer resp.Body.Close() log.Println("Checkin response Status:", resp.Status) } // A helper to check if a env var is defined. If not defined, show error. // Args: key - the env var to check for // Returns if exists, the value of the var func getEnv(key string) string { value := os.Getenv(key) if len(value) == 0 { log.Fatalf("Env var %s is required and missing.", key) } log.Printf("%s=%s", key, value) return value } // A helper to check and assign all required ENV vars to the Config struct // Returns a pointer to the Config struct func loadConfig() *Config { conf := Config{ JOB_NAME: getEnv("JOB_NAME"), TZ: getEnv("TZ"), WINDOW_SIZE: getEnv("WINDOW_SIZE"), MYSQL_HOST: getEnv("MYSQL_HOST"), MYSQL_USER: getEnv("MYSQL_USER"), MYSQL_PASSWORD: getEnv("MYSQL_PASSWORD"), MYSQL_DB: getEnv("MYSQL_DB"), MYSQL_TABLE: getEnv("MYSQL_TABLE"), DRUID_HOST: getEnv("DRUID_HOST"), INFLUX_HOST: getEnv("INFLUX_HOST"), INFLUX_PORT: getEnv("INFLUX_PORT"), INFLUX_DB: getEnv("INFLUX_DB"), INFLUX_RP: getEnv("INFLUX_RP"), INFLUX_MEASUREMENT: getEnv("INFLUX_MEASUREMENT"), DVO_API_CHECKIN_ENDPOINT: getEnv("DVO_API_CHECKIN_ENDPOINT"), } return &conf } func main() { // setup appConfig := loadConfig() loc, _ := time.LoadLocation(appConfig.TZ) windowSize, err := strconv.Atoi(appConfig.WINDOW_SIZE) if err != nil { log.Fatal(err) } // druid druidClient := initDruidClient(appConfig.DRUID_HOST) // mysql mySqlDb := initMySqlDb(appConfig) preparedStmt := getSqlPreparedStatement(mySqlDb, appConfig.MYSQL_TABLE) // influx influxClient := initInfluxClient(appConfig) // main logic now := time.Now().UTC() startTime := now.Add(time.Duration(windowSize) * -time.Hour) endTime := now.Add(1 * -time.Hour) // skip the current hour for d := endTime; d.Before(startTime) == false; d = d.Add(-1 * time.Hour) { // get the formatted hour for mysql and druid estHrMysql := fmt.Sprintf("%s:00:00", d.In(loc).Format("2006-01-02 15")) // using GO constant format https://golang.org/src/pkg/time/format.go utcHrDruid := fmt.Sprintf("%s", d.Format("2006-01-02T15Z")) log.Printf("Mysql date: %s | Druid date: %s\n", estHrMysql, utcHrDruid) // get the mysql stats for this hour mySqlStats := getMySqlStatsForInterval(preparedStmt, estHrMysql) if mySqlStats.Views.Valid { log.Printf("MYSQL STATS: Views: %f | Clicks: %f | Leads: %f", mySqlStats.Views.Float64, mySqlStats.Clicks.Float64, mySqlStats.Leads.Float64) } else { log.Println("no sql row was returned for this hour") } // get the druid stats for this hour druidStats := getDruidStatsForInterval(druidClient, utcHrDruid) // store in influx if druidStats.Views.Valid { log.Printf("DRUID STATS: Views: %f | Clicks: %f | Leads: %f", druidStats.Views.Float64, druidStats.Clicks.Float64, druidStats.Leads.Float64) } else { log.Println("no druid result was returned for this hour") } // store results in influx if mySqlStats.Views.Valid && druidStats.Views.Valid { _ = writeToInflux(influxClient, appConfig, &mySqlStats, &druidStats, d.In(loc).Truncate(time.Hour)) } } // checkin job checkInJob(appConfig.JOB_NAME, appConfig.DVO_API_CHECKIN_ENDPOINT) }