|
| 1 | +package main |
| 2 | + |
| 3 | +import ( |
| 4 | +"flag" |
| 5 | +"fmt" |
| 6 | +"os" |
| 7 | +"os/signal" |
| 8 | +"strings" |
| 9 | +"syscall" |
| 10 | + |
| 11 | +"github.com/siddontang/go-mysql/canal" |
| 12 | +) |
| 13 | + |
| 14 | +var host = flag.String("host", "127.0.0.1", "MySQL host") |
| 15 | +var port = flag.Int("port", 3306, "MySQL port") |
| 16 | +var user = flag.String("user", "root", "MySQL user, must have replication privilege") |
| 17 | +var password = flag.String("password", "", "MySQL password") |
| 18 | + |
| 19 | +var flavor = flag.String("flavor", "mysql", "Flavor: mysql or mariadb") |
| 20 | + |
| 21 | +var dataDir = flag.String("data-dir", "./var", "Path to store data, like master.info") |
| 22 | + |
| 23 | +var serverID = flag.Int("server-id", 101, "Unique Server ID") |
| 24 | +var mysqldump = flag.String("mysqldump", "mysqldump", "mysqldump execution path") |
| 25 | + |
| 26 | +var dbs = flag.String("dbs", "test", "dump databases, seperated by comma") |
| 27 | +var tables = flag.String("tables", "", "dump tables, seperated by comma, will overwrite dbs") |
| 28 | +var tableDB = flag.String("table_db", "test", "database for dump tables") |
| 29 | +var ignoreTables = flag.String("ignore_tables", "", "ignore tables, must be database.table format, separated by comma") |
| 30 | + |
| 31 | +func main() { |
| 32 | +cfg := canal.NewDefaultConfig() |
| 33 | +cfg.Addr = fmt.Sprintf("%s:%d", *host, *port) |
| 34 | +cfg.User = *user |
| 35 | +cfg.Password = *password |
| 36 | +cfg.Flavor = *flavor |
| 37 | +cfg.DataDir = *dataDir |
| 38 | + |
| 39 | +cfg.ServerID = uint32(*serverID) |
| 40 | +cfg.Dump.ExecutionPath = *mysqldump |
| 41 | +cfg.Dump.DiscardErr = false |
| 42 | + |
| 43 | +c, err := canal.NewCanal(cfg) |
| 44 | +if err != nil { |
| 45 | +fmt.Printf("create canal err %v", err) |
| 46 | +os.Exit(1) |
| 47 | +} |
| 48 | + |
| 49 | +if len(*ignoreTables) == 0 { |
| 50 | +subs := strings.Split(*ignoreTables, ",") |
| 51 | +for _, sub := range subs { |
| 52 | +if seps := strings.Split(sub, "."); len(seps) == 2 { |
| 53 | +c.AddDumpIgnoreTables(seps[0], seps[1]) |
| 54 | +} |
| 55 | +} |
| 56 | +} |
| 57 | + |
| 58 | +if len(*tables) > 0 && len(*tableDB) > 0 { |
| 59 | +subs := strings.Split(*tables, ",") |
| 60 | +c.AddDumpTables(*tableDB, subs...) |
| 61 | +} else if len(*dbs) > 0 { |
| 62 | +subs := strings.Split(*dbs, ",") |
| 63 | +c.AddDumpDatabases(subs...) |
| 64 | +} |
| 65 | + |
| 66 | +c.RegRowsEventHandler(&handler{}) |
| 67 | + |
| 68 | +err = c.Start() |
| 69 | +if err != nil { |
| 70 | +fmt.Printf("start canal err %V", err) |
| 71 | +os.Exit(1) |
| 72 | +} |
| 73 | + |
| 74 | +sc := make(chan os.Signal, 1) |
| 75 | +signal.Notify(sc, |
| 76 | +os.Kill, |
| 77 | +os.Interrupt, |
| 78 | +syscall.SIGHUP, |
| 79 | +syscall.SIGINT, |
| 80 | +syscall.SIGTERM, |
| 81 | +syscall.SIGQUIT) |
| 82 | + |
| 83 | +<-sc |
| 84 | + |
| 85 | +c.Close() |
| 86 | +} |
| 87 | + |
| 88 | +type handler struct { |
| 89 | +} |
| 90 | + |
| 91 | +func (h *handler) Do(e *canal.RowsEvent) error { |
| 92 | +fmt.Printf("%v\n", e) |
| 93 | + |
| 94 | +return nil |
| 95 | +} |
| 96 | + |
| 97 | +func (h *handler) String() string { |
| 98 | +return "TestHandler" |
| 99 | +} |
0 commit comments